From e4421d966ad734c4591d94eb66069bb0ed5971f6 Mon Sep 17 00:00:00 2001 From: saphid <4596216+saphid@users.noreply.github.com> Date: Mon, 28 Sep 2026 21:51:02 +1000 Subject: [PATCH] Bundled KDE Connect: the server owns a copy until its agent speaks The agent holds its copy from its first status line on and tidies it up however it ends. Before that (a launch error, or stopped before the agent ran) the server removes the copy itself. discard() only ever removes incoming copies, never the iPhone bundle's own. The retry race test waits for both contenders' decisions instead of sleeping. Co-Authored-By: Claude Opus 5.5 (1M context) --- tests/test_input.py | 35 ++++++++++++++++++++++++++++++----- ui/server.py | 25 ++++++++++++++++--------- 2 files changed, 46 insertions(+), 14 deletions(-) diff --git a/tests/test_input.py b/tests/test_input.py index f5b4175..a214a8b 100644 --- a/tests/test_input.py +++ b/tests/test_input.py @@ -217,7 +217,7 @@ class Bundled(unittest.TestCase): @property def stdout(self): yield from self.lines - if b"need-packages" not in self.lines[-1]: + if self.lines and b"need-packages" not in self.lines[-1]: while not (done.is_set() or self.ended.is_set()): self.ended.wait(0.02) self.ended.set() @@ -298,11 +298,11 @@ class Bundled(unittest.TestCase): watch, launch, deliver = agent._watch, agent._launch, agent.deliver def racing_watch(proc, errors, retry=False): - wanted = watch(proc, errors, retry) + wanted, heard = watch(proc, errors, retry) if wanted and not first: first.append(threading.current_thread()) agent.start() # lands between the agent exiting and the retry - return wanted + return wanted, heard def held_deliver(report, force=False): if force: @@ -319,8 +319,9 @@ class Bundled(unittest.TestCase): decided.set() # the retry declined agent._watch, agent._launch, agent.deliver = racing_watch, first_launch, held_deliver agent.start() + # Both have decided once the new start has launched its agent (it always does). + self.wait_for(lambda: decided.is_set() and launches.count("") == 2 and len(procs) == len(launches)) self.wait_for(lambda: agent.status == {"state": "ready"}) - time.sleep(0.1) 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)) @@ -331,13 +332,37 @@ class Bundled(unittest.TestCase): removed = [] agent.discard = removed.append agent.deliver = lambda report, force=False: (entered.set(), gate.wait(5), "~/incoming/x")[2] + self.addCleanup(lambda: self.assertEqual(removed, ["~/incoming/x"])) agent.start() self.assertTrue(entered.wait(5)) agent.stop() gate.set() - self.wait_for(lambda: removed == ["incoming/x"]) + self.wait_for(lambda: removed == ["~/incoming/x"]) self.assertEqual(procs, []) # no agent started for it + def test_a_copy_whose_agent_never_starts_is_removed(self): + agent, launches, procs = self.lifecycle([[]]) # the agent dies before saying anything + removed = [] + agent.discard = removed.append + agent.deliver = lambda report, force=False: "~/incoming/y" + agent.start() + self.wait_for(lambda: removed == ["~/incoming/y"]) + + def test_discard_only_touches_copies(self): + calls = [] + old = self.server.ssh, self.server.LOCAL + self.server.ssh, self.server.LOCAL = (lambda remote, **k: calls.append(remote)), False + try: + agent = self.server.InputAgent(packages=[]) + agent.discard("") + agent.discard("~/.local/share/frame-control/kdeconnect") + agent.discard("~/.local/share/frame-control/kdeconnect/incoming/ab12") + self.server.LOCAL = True + agent.discard("~/.local/share/frame-control/kdeconnect/incoming/ab12") + finally: + self.server.ssh, self.server.LOCAL = old + self.assertEqual(calls, ["rm -rf .local/share/frame-control/kdeconnect/incoming/ab12"]) + def test_copies_go_to_a_folder_of_their_own(self): first, _, _ = self.deliver(frame_has=False) second, _, _ = self.deliver(frame_has=False) diff --git a/ui/server.py b/ui/server.py index 137cc19..587a355 100755 --- a/ui/server.py +++ b/ui/server.py @@ -655,7 +655,8 @@ class InputAgent: def discard(self, folder): """Remove a copy no agent will take over (best effort; agents tidy up old ones too).""" - if folder and not LOCAL: + folder = folder.removeprefix("~/") + if folder.startswith(f"{KDECONNECT_HOME}/incoming/") and not LOCAL: try: ssh(f"rm -rf {folder}", timeout=20) except Failure: @@ -675,18 +676,18 @@ class InputAgent: with self.lock: if self.generation == generation: self.status = {"state": "installing", "message": message} - errors = tempfile.TemporaryFile() + errors, folder = tempfile.TemporaryFile(), "" try: ensure_master() folder = self.deliver(report, force) with self.lock: stopped = self.generation != generation - if stopped: # turned off while copying: no agent will take the copy over - self.discard(folder.removeprefix("~/")) + if stopped: # turned off while copying raise Failure("stopped") proc = subprocess.Popen([*SSH, FRAME, self.command(folder)], stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=errors) except (Failure, OSError) as e: + self.discard(folder) # no agent will take the copy over message = str(e) friendly = unreachable(message) with self.lock: @@ -705,7 +706,11 @@ class InputAgent: self.proc = proc if stale: # turned off meanwhile proc.terminate() - if self._watch(proc, errors, retry=not force) and not force: + wanted, heard = self._watch(proc, errors, retry=not force) + if not heard: + # The agent never started (or was stopped first), so it can't tidy the copy up. + self.discard(folder) + if wanted and not force: # It needed the packages after all (another device changed what's # installed after we looked): copy them and start once more. with self.lock: @@ -716,15 +721,17 @@ class InputAgent: self._launch(generation, force=True) def _watch(self, proc, errors, retry=False): - """Follow the agent's status until it exits. True if it asked for the packages - (`retry`: the caller will send them, so that isn't an error yet).""" - wanted = False + """Follow the agent's status until it exits. Returns whether it asked for the + packages (`retry`: the caller will send them, so that isn't an error yet), and + whether it said anything at all (then it holds its copy and tidies it up).""" + wanted = heard = False for line in proc.stdout: try: status = json.loads(line) except ValueError: continue if isinstance(status, dict) and isinstance(status.get("state"), str): + heard = True if status["state"] == "need-packages": wanted = True continue @@ -742,7 +749,7 @@ class InputAgent: message = "KDE Connect didn't reach the Frame" friendly = unreachable(message) self.status = {"state": "error", "message": friendly or message, **({"offline": True} if friendly else {})} - return wanted + return wanted, heard def send(self, events): """Forward events if the agent is ready; start it if it isn't running.