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.