mirror of
https://github.com/saphid/frame-control.git
synced 2026-10-06 02:00:19 +02:00
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) <noreply@anthropic.com>
This commit is contained in:
1 parent
efffd72c52
commit
e4421d966a
2 files changed
+46
-14
No files matched your search
+30
-5
@@ -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)
|
||||
|
||||
+16
-9
@@ -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.
|
||||
|
||||
Reference in new issue
Block a user