Bundled KDE Connect: tidy the copied packages however a start ends

- Each start keeps its copy (holding a lock on it) until it ends, however
  it ends, then removes it; a restart can still unpack it.
- The server removes a partial copy when copying fails.
- Other copies are removed only when unlocked and over an hour old.
- Tests: a start racing the need-packages retry (one agent, not two), a
  restart that must unpack its copy, a held copy surviving the sweep.
  Both regression tests fail without their fix.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
saphidandClaude Opus 5.5 committed 2026-09-28 21:36:22 +10:00
1 parent 6765fbcb20
commit d5e8193ce5
3 files changed
+105 -7

No files matched your search

+55 -1
View File
@@ -220,7 +220,16 @@ class Bundled(unittest.TestCase):
self.server.ensure_master, self.server.subprocess.Popen = lambda: None, popen
self.addCleanup(lambda: (setattr(self.server, "ensure_master", old[0]),
setattr(self.server.subprocess, "Popen", old[1])))
return Agent(packages=[("a.pkg.tar.zst", "0" * 64)]), launches, procs
agent = Agent(packages=[("a.pkg.tar.zst", "0" * 64)])
def settle(): # nothing of this test may still be launching when the next one starts
agent.stop()
for _ in range(250):
if agent.launching is None and all(p.ended.is_set() for p in procs):
return
time.sleep(0.02)
self.addCleanup(settle)
return agent, launches, procs
def wait_for(self, check):
for _ in range(250):
@@ -261,6 +270,26 @@ class Bundled(unittest.TestCase):
agent.stop()
self.wait_for(lambda: all(p.ended.is_set() for p in procs))
def test_retry_and_a_new_start_never_run_two_agents(self):
# A start() landing just as an agent that asked for the packages exits: exactly one
# of it and the retry launches, and stop() ends everything.
agent, launches, procs = self.lifecycle([[b'{"state": "need-packages"}\n'], [b'{"state": "ready"}\n']])
watch, raced = agent._watch, []
def racing_watch(proc, errors, retry=False):
wanted = watch(proc, errors, retry)
if wanted and not raced:
raced.append(True)
agent.start() # lands between the agent exiting and the retry
return wanted
agent._watch = racing_watch
agent.start()
self.wait_for(lambda: agent.status == {"state": "ready"})
time.sleep(0.2)
self.assertEqual(sum(not p.ended.is_set() for p in procs), 1, launches)
agent.stop()
self.wait_for(lambda: all(p.ended.is_set() for p in procs))
def test_copies_go_to_a_folder_of_their_own(self):
first, _, _ = self.deliver(frame_has=False)
second, _, _ = self.deliver(frame_has=False)
@@ -336,6 +365,31 @@ class AgentInstall(unittest.TestCase):
finally:
self.agent.listening = old
def test_restart_can_still_unpack_its_copy_then_tidies_it(self):
folder = self.agent.BASE / "incoming/abc"
folder.mkdir(parents=True)
seen = []
stubs = {"ensure_daemon": lambda f, p: seen.append(Path(f).is_dir()),
"identity": lambda c: ("id", "cert", "key"), "our_daemons": lambda: [123],
"connect": lambda *a: (_ for _ in ()).throw(OSError("no answer"))}
saved = {k: getattr(self.agent, k) for k in stubs}
self.addCleanup(lambda: [setattr(self.agent, k, v) for k, v in saved.items()])
for k, v in stubs.items():
setattr(self.agent, k, v)
self.assertEqual(self.agent.run("c", "n", str(folder), []), 1)
self.assertEqual(seen, [True, True]) # there for the first start and the restart
self.assertFalse(folder.exists())
def test_tidy_keeps_copies_in_use(self):
incoming = self.agent.BASE / "incoming"
held_dir = incoming / "held"
held_dir.mkdir(parents=True)
held = self.agent.hold_incoming(str(held_dir))
self.addCleanup(held.close)
os.utime(held_dir, (time.time() - 7200,) * 2)
self.agent.tidy_incoming("")
self.assertTrue(held_dir.is_dir())
def test_tidy_keeps_other_starts_copies(self):
incoming = self.agent.BASE / "incoming"
mine, theirs, stale = incoming / "mine", incoming / "theirs", incoming / "stale"
+43 -6
View File
@@ -347,17 +347,43 @@ def main():
def tidy_incoming(folder):
"""Remove this start's copy of the packages, and any a cancelled start left over an hour ago."""
"""Remove this start's copy of the packages, and others nobody is using.
Another copy goes only if no agent holds its .in-use lock and it's over an
hour old (so not one a server is still copying, before its agent starts).
"""
incoming = BASE / "incoming"
if folder.startswith(str(incoming) + "/"):
shutil.rmtree(folder, ignore_errors=True)
try:
for old in incoming.iterdir():
if time.time() - old.stat().st_mtime > 3600:
shutil.rmtree(old, ignore_errors=True)
others = list(incoming.iterdir())
except OSError:
return
for other in others:
try:
if time.time() - other.stat().st_mtime < 3600:
continue
with open(other / ".in-use", "a") as lock:
fcntl.flock(lock, fcntl.LOCK_EX | fcntl.LOCK_NB)
shutil.rmtree(other, ignore_errors=True)
except OSError:
pass # in use, or already gone
try:
incoming.rmdir()
except OSError:
pass # none, or another start's copy is still there
pass # another start's copy is still there
def hold_incoming(folder):
"""Mark this start's copy as in use (tidy_incoming leaves it alone); the lock lasts as long as the file."""
if not folder.startswith(str(BASE / "incoming") + "/"):
return None
try:
lock = open(Path(folder) / ".in-use", "a")
fcntl.flock(lock, fcntl.LOCK_SH)
return lock
except OSError:
return None
class daemon_lock:
@@ -372,6 +398,18 @@ class daemon_lock:
def run(client, name, folder, packages):
"""Set up, pair and forward events. This start's copy of the packages stays
until it ends, however it ends: restarting KDE Connect may need to unpack it."""
held = hold_incoming(folder)
try:
return serve(client, name, folder, packages)
finally:
if held:
held.close()
tidy_incoming(folder)
def serve(client, name, folder, packages):
try:
with daemon_lock():
ensure_daemon(folder, packages)
@@ -398,7 +436,6 @@ def run(client, name, folder, packages):
except (OSError, RuntimeError, subprocess.SubprocessError) as e:
say("error", message=str(e))
return 1
tidy_incoming(folder) # unpacked or not needed: the copy has done its job
say("ready", keyboard=link.keyboard is not False)
stdin, pending = sys.stdin.fileno(), b""
sel = selectors.DefaultSelector()
+7
View File
@@ -638,6 +638,7 @@ class InputAgent:
# A folder of its own: a cancelled start's agent may still be cleaning up another.
folder = f"{KDECONNECT_HOME}/incoming/{secrets.token_hex(8)}"
ssh(f"mkdir -p {folder}", timeout=20)
try:
for name, sha in self.packages:
path = KDECONNECT / "packages" / name
if not path.is_file() or file_sha256(path) != sha:
@@ -647,6 +648,12 @@ class InputAgent:
quoted = shlex.quote(name)
ssh(f"cd {folder} && cat > {quoted}.part && mv {quoted}.part {quoted}",
stdin=path.read_bytes(), text=False, timeout=600)
except Failure:
try:
ssh(f"rm -rf {folder}", timeout=20) # a partial copy is no use to anyone
except Failure:
pass # the agent tidies it later
raise
return f"~/{folder}"
def start(self):