mirror of
https://github.com/saphid/frame-control.git
synced 2026-10-06 06:00:33 +02:00
Bundled KDE Connect: close the races the second review found
- The need-packages retry only runs if no start() has taken over, so two agents can't end up running for one session. - Each copy goes to its own incoming folder; the agent removes its own once connected (and any a cancelled start left over an hour ago), so a stopped start can't delete files another is copying or still needs. - A missing package folder means need-packages, not an error. - Tests keep fake agents running, so they check ready and which agent owns the session, plus a bounded retry and the tidy rule. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
1 parent
75b2478a55
commit
6765fbcb20
3 files changed
+84
-31
No files matched your search
+63
-19
@@ -173,9 +173,11 @@ class Bundled(unittest.TestCase):
|
|||||||
self.deliver(frame_has=False, damaged=True)
|
self.deliver(frame_has=False, damaged=True)
|
||||||
|
|
||||||
def lifecycle(self, agent_lines, deliver=None):
|
def lifecycle(self, agent_lines, deliver=None):
|
||||||
"""An InputAgent whose ssh and agent are fakes; returns it and the folders each launch used."""
|
"""An InputAgent whose ssh and agent are fakes. Returns it, the folders each launch
|
||||||
launches = []
|
used, and the fake agents. A fake agent reports its lines, then keeps running
|
||||||
test = self
|
(unless its lines end with need-packages) until the test ends."""
|
||||||
|
launches, procs, done = [], [], threading.Event()
|
||||||
|
self.addCleanup(done.set)
|
||||||
|
|
||||||
class Agent(self.server.InputAgent):
|
class Agent(self.server.InputAgent):
|
||||||
def deliver(self, report, force=False):
|
def deliver(self, report, force=False):
|
||||||
@@ -191,34 +193,60 @@ class Bundled(unittest.TestCase):
|
|||||||
stdin = None
|
stdin = None
|
||||||
|
|
||||||
def __init__(self, lines):
|
def __init__(self, lines):
|
||||||
self.stdout = iter(lines)
|
self.lines, self.ended = lines, threading.Event()
|
||||||
|
procs.append(self)
|
||||||
|
|
||||||
|
@property
|
||||||
|
def stdout(self):
|
||||||
|
yield from self.lines
|
||||||
|
if 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()
|
||||||
|
|
||||||
def wait(self):
|
def wait(self):
|
||||||
|
self.ended.wait(5)
|
||||||
return 0
|
return 0
|
||||||
|
|
||||||
def poll(self):
|
def poll(self):
|
||||||
return None
|
return 0 if self.ended.is_set() else None
|
||||||
|
|
||||||
def terminate(self):
|
def terminate(self):
|
||||||
pass
|
self.ended.set()
|
||||||
|
|
||||||
def popen(*a, **k):
|
def popen(*a, **k):
|
||||||
return Proc(agent_lines.pop(0) if agent_lines else [b'{"state": "ready"}\n'])
|
return Proc(agent_lines.pop(0) if agent_lines else [b'{"state": "ready"}\n'])
|
||||||
old = self.server.ensure_master, self.server.subprocess.Popen
|
old = self.server.ensure_master, self.server.subprocess.Popen
|
||||||
self.server.ensure_master, self.server.subprocess.Popen = lambda: None, popen
|
self.server.ensure_master, self.server.subprocess.Popen = lambda: None, popen
|
||||||
test.addCleanup(lambda: (setattr(self.server, "ensure_master", old[0]),
|
self.addCleanup(lambda: (setattr(self.server, "ensure_master", old[0]),
|
||||||
setattr(self.server.subprocess, "Popen", old[1])))
|
setattr(self.server.subprocess, "Popen", old[1])))
|
||||||
return Agent(packages=[("a.pkg.tar.zst", "0" * 64)]), launches
|
return Agent(packages=[("a.pkg.tar.zst", "0" * 64)]), launches, procs
|
||||||
|
|
||||||
|
def wait_for(self, check):
|
||||||
|
for _ in range(250):
|
||||||
|
if check():
|
||||||
|
return
|
||||||
|
time.sleep(0.02)
|
||||||
|
self.fail("timed out")
|
||||||
|
|
||||||
def test_agent_asking_for_packages_gets_them_and_starts_again(self):
|
def test_agent_asking_for_packages_gets_them_and_starts_again(self):
|
||||||
agent, launches = self.lifecycle([[b'{"state": "need-packages"}\n'], [b'{"state": "ready"}\n']])
|
agent, launches, procs = self.lifecycle([[b'{"state": "need-packages"}\n'], [b'{"state": "ready"}\n']])
|
||||||
agent._launch(agent.generation)
|
agent.start()
|
||||||
|
self.wait_for(lambda: agent.status == {"state": "ready"})
|
||||||
self.assertEqual(launches, ["", "~/copied"])
|
self.assertEqual(launches, ["", "~/copied"])
|
||||||
self.assertNotIn("reach", agent.status.get("message", ""))
|
self.assertIs(agent.proc, procs[1])
|
||||||
|
|
||||||
|
def test_asking_twice_is_an_error_not_a_loop(self):
|
||||||
|
need = [b'{"state": "need-packages"}\n']
|
||||||
|
agent, launches, _ = self.lifecycle([need, list(need)])
|
||||||
|
agent.start()
|
||||||
|
self.wait_for(lambda: agent.status.get("state") == "error")
|
||||||
|
self.assertEqual(launches, ["", "~/copied"])
|
||||||
|
self.assertIn("didn't reach", agent.status["message"])
|
||||||
|
|
||||||
def test_start_after_stop_during_the_copy_still_starts(self):
|
def test_start_after_stop_during_the_copy_still_starts(self):
|
||||||
gate, entered = threading.Event(), threading.Event()
|
gate, entered = threading.Event(), threading.Event()
|
||||||
agent, launches = self.lifecycle([], deliver=lambda: (entered.set(), gate.wait(5)))
|
agent, launches, procs = self.lifecycle([], deliver=lambda: (entered.set(), gate.wait(5)))
|
||||||
agent.start()
|
agent.start()
|
||||||
self.assertTrue(entered.wait(5))
|
self.assertTrue(entered.wait(5))
|
||||||
agent.stop()
|
agent.stop()
|
||||||
@@ -226,11 +254,17 @@ class Bundled(unittest.TestCase):
|
|||||||
agent.start() # while the first launch is still copying
|
agent.start() # while the first launch is still copying
|
||||||
self.assertTrue(entered.wait(5), "the second start didn't launch")
|
self.assertTrue(entered.wait(5), "the second start didn't launch")
|
||||||
gate.set()
|
gate.set()
|
||||||
for _ in range(100):
|
self.wait_for(lambda: len(procs) == 2 and agent.status == {"state": "ready"})
|
||||||
if len(launches) == 2:
|
# The stopped launch's agent was ended; the new one is the one in use.
|
||||||
break
|
self.wait_for(lambda: sum(p.ended.is_set() for p in procs) == 1)
|
||||||
time.sleep(0.02)
|
self.assertFalse(agent.proc.ended.is_set())
|
||||||
self.assertEqual(len(launches), 2) # the stopped launch ran its agent too, then ended it
|
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)
|
||||||
|
self.assertNotEqual(first, second)
|
||||||
|
|
||||||
|
|
||||||
@unittest.skipIf(sys.platform == "win32", "the agent runs on the Frame (Linux)")
|
@unittest.skipIf(sys.platform == "win32", "the agent runs on the Frame (Linux)")
|
||||||
@@ -296,11 +330,21 @@ class AgentInstall(unittest.TestCase):
|
|||||||
def test_asks_for_packages_it_was_not_sent(self):
|
def test_asks_for_packages_it_was_not_sent(self):
|
||||||
self.agent.listening, old = (lambda: False), self.agent.listening
|
self.agent.listening, old = (lambda: False), self.agent.listening
|
||||||
try:
|
try:
|
||||||
with self.assertRaises(self.agent.NeedPackages):
|
for folder in ("", str(self.dir / "gone")): # none sent, or already tidied away
|
||||||
self.agent.ensure_daemon("", [("new.pkg.tar.zst", "1" * 64)])
|
with self.assertRaises(self.agent.NeedPackages):
|
||||||
|
self.agent.ensure_daemon(folder, [("new.pkg.tar.zst", "1" * 64)])
|
||||||
finally:
|
finally:
|
||||||
self.agent.listening = old
|
self.agent.listening = old
|
||||||
|
|
||||||
|
def test_tidy_keeps_other_starts_copies(self):
|
||||||
|
incoming = self.agent.BASE / "incoming"
|
||||||
|
mine, theirs, stale = incoming / "mine", incoming / "theirs", incoming / "stale"
|
||||||
|
for d in (mine, theirs, stale):
|
||||||
|
d.mkdir(parents=True)
|
||||||
|
os.utime(stale, (time.time() - 7200,) * 2)
|
||||||
|
self.agent.tidy_incoming(str(mine))
|
||||||
|
self.assertEqual(sorted(p.name for p in incoming.iterdir()), ["theirs"])
|
||||||
|
|
||||||
def test_server_and_agent_agree_on_the_stamp(self):
|
def test_server_and_agent_agree_on_the_stamp(self):
|
||||||
import server
|
import server
|
||||||
packages = [("a.pkg.tar.zst", "1" * 64), ("b.pkg.tar.zst", "2" * 64)]
|
packages = [("a.pkg.tar.zst", "1" * 64), ("b.pkg.tar.zst", "2" * 64)]
|
||||||
|
|||||||
+16
-7
@@ -166,7 +166,7 @@ def ensure_daemon(folder, packages):
|
|||||||
"""
|
"""
|
||||||
system = SYSTEM_DAEMON.exists()
|
system = SYSTEM_DAEMON.exists()
|
||||||
if not system and not installed(packages) and not (listening() and our_daemons()):
|
if not system and not installed(packages) and not (listening() and our_daemons()):
|
||||||
if not folder:
|
if not folder or not Path(folder).is_dir():
|
||||||
raise NeedPackages()
|
raise NeedPackages()
|
||||||
install(folder, packages)
|
install(folder, packages)
|
||||||
if listening():
|
if listening():
|
||||||
@@ -346,6 +346,20 @@ def main():
|
|||||||
stop_daemon()
|
stop_daemon()
|
||||||
|
|
||||||
|
|
||||||
|
def tidy_incoming(folder):
|
||||||
|
"""Remove this start's copy of the packages, and any a cancelled start left over an hour ago."""
|
||||||
|
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)
|
||||||
|
incoming.rmdir()
|
||||||
|
except OSError:
|
||||||
|
pass # none, or another start's copy is still there
|
||||||
|
|
||||||
|
|
||||||
class daemon_lock:
|
class daemon_lock:
|
||||||
"""Installing, starting and restarting KDE Connect happen one agent at a time."""
|
"""Installing, starting and restarting KDE Connect happen one agent at a time."""
|
||||||
|
|
||||||
@@ -361,12 +375,6 @@ def run(client, name, folder, packages):
|
|||||||
try:
|
try:
|
||||||
with daemon_lock():
|
with daemon_lock():
|
||||||
ensure_daemon(folder, packages)
|
ensure_daemon(folder, packages)
|
||||||
if folder.startswith(str(BASE / "incoming") + "/"):
|
|
||||||
shutil.rmtree(folder, ignore_errors=True) # unpacked; the copy isn't needed again
|
|
||||||
try:
|
|
||||||
(BASE / "incoming").rmdir()
|
|
||||||
except OSError:
|
|
||||||
pass # another device's copy is still there
|
|
||||||
device, cert, key = identity(client)
|
device, cert, key = identity(client)
|
||||||
say("pairing")
|
say("pairing")
|
||||||
seen = our_daemons()
|
seen = our_daemons()
|
||||||
@@ -390,6 +398,7 @@ def run(client, name, folder, packages):
|
|||||||
except (OSError, RuntimeError, subprocess.SubprocessError) as e:
|
except (OSError, RuntimeError, subprocess.SubprocessError) as e:
|
||||||
say("error", message=str(e))
|
say("error", message=str(e))
|
||||||
return 1
|
return 1
|
||||||
|
tidy_incoming(folder) # unpacked or not needed: the copy has done its job
|
||||||
say("ready", keyboard=link.keyboard is not False)
|
say("ready", keyboard=link.keyboard is not False)
|
||||||
stdin, pending = sys.stdin.fileno(), b""
|
stdin, pending = sys.stdin.fileno(), b""
|
||||||
sel = selectors.DefaultSelector()
|
sel = selectors.DefaultSelector()
|
||||||
|
|||||||
+5
-5
@@ -635,8 +635,8 @@ class InputAgent:
|
|||||||
" && echo yes || true", stdin=kdeconnect_stamp(self.packages).encode(), text=False, timeout=20)
|
" && echo yes || true", stdin=kdeconnect_stamp(self.packages).encode(), text=False, timeout=20)
|
||||||
if have.strip() == b"yes":
|
if have.strip() == b"yes":
|
||||||
return ""
|
return ""
|
||||||
client = "".join(c for c in input_client() if c.isalnum() or c in "-_")[:64] or "default"
|
# A folder of its own: a cancelled start's agent may still be cleaning up another.
|
||||||
folder = f"{KDECONNECT_HOME}/incoming/{client}"
|
folder = f"{KDECONNECT_HOME}/incoming/{secrets.token_hex(8)}"
|
||||||
ssh(f"mkdir -p {folder}", timeout=20)
|
ssh(f"mkdir -p {folder}", timeout=20)
|
||||||
for name, sha in self.packages:
|
for name, sha in self.packages:
|
||||||
path = KDECONNECT / "packages" / name
|
path = KDECONNECT / "packages" / name
|
||||||
@@ -645,8 +645,7 @@ class InputAgent:
|
|||||||
" (a build runs app/build/fetch-deps.js to add it)", 500)
|
" (a build runs app/build/fetch-deps.js to add it)", 500)
|
||||||
report(f"Copying KDE Connect to the Frame ({name.rsplit('-', 3)[0]})")
|
report(f"Copying KDE Connect to the Frame ({name.rsplit('-', 3)[0]})")
|
||||||
quoted = shlex.quote(name)
|
quoted = shlex.quote(name)
|
||||||
# Its own temporary name: two starts can be copying at once.
|
ssh(f"cd {folder} && cat > {quoted}.part && mv {quoted}.part {quoted}",
|
||||||
ssh(f"cd {folder} && cat > {quoted}.part$$ && mv {quoted}.part$$ {quoted}",
|
|
||||||
stdin=path.read_bytes(), text=False, timeout=600)
|
stdin=path.read_bytes(), text=False, timeout=600)
|
||||||
return f"~/{folder}"
|
return f"~/{folder}"
|
||||||
|
|
||||||
@@ -693,7 +692,8 @@ class InputAgent:
|
|||||||
# It needed the packages after all (another device changed what's
|
# It needed the packages after all (another device changed what's
|
||||||
# installed after we looked): copy them and start once more.
|
# installed after we looked): copy them and start once more.
|
||||||
with self.lock:
|
with self.lock:
|
||||||
if self.proc is not proc or self.generation != generation:
|
# Unless a start() already took over (it launches, and copies if still needed).
|
||||||
|
if self.proc is not proc or self.generation != generation or self.launching is not None:
|
||||||
return
|
return
|
||||||
self.proc, self.launching, self.status = None, generation, {"state": "starting"}
|
self.proc, self.launching, self.status = None, generation, {"state": "starting"}
|
||||||
self._launch(generation, force=True)
|
self._launch(generation, force=True)
|
||||||
|
|||||||
Reference in new issue
Block a user