diff --git a/tests/test_input.py b/tests/test_input.py index 4f7e66c..5bd8707 100644 --- a/tests/test_input.py +++ b/tests/test_input.py @@ -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" diff --git a/ui/frame_input_agent.py b/ui/frame_input_agent.py index 2c45a04..fe437b9 100644 --- a/ui/frame_input_agent.py +++ b/ui/frame_input_agent.py @@ -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() diff --git a/ui/server.py b/ui/server.py index 4ab2425..794e269 100755 --- a/ui/server.py +++ b/ui/server.py @@ -638,15 +638,22 @@ 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) - for name, sha in self.packages: - path = KDECONNECT / "packages" / name - if not path.is_file() or file_sha256(path) != sha: - raise Failure(f"{name} is missing or damaged in this copy of Frame Control" - " (a build runs app/build/fetch-deps.js to add it)", 500) - report(f"Copying KDE Connect to the Frame ({name.rsplit('-', 3)[0]})") - quoted = shlex.quote(name) - ssh(f"cd {folder} && cat > {quoted}.part && mv {quoted}.part {quoted}", - stdin=path.read_bytes(), text=False, timeout=600) + try: + for name, sha in self.packages: + path = KDECONNECT / "packages" / name + if not path.is_file() or file_sha256(path) != sha: + raise Failure(f"{name} is missing or damaged in this copy of Frame Control" + " (a build runs app/build/fetch-deps.js to add it)", 500) + report(f"Copying KDE Connect to the Frame ({name.rsplit('-', 3)[0]})") + 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):