diff --git a/ios/FrameControl.xcodeproj/project.pbxproj b/ios/FrameControl.xcodeproj/project.pbxproj index d23f1b8..c1fa8f9 100644 --- a/ios/FrameControl.xcodeproj/project.pbxproj +++ b/ios/FrameControl.xcodeproj/project.pbxproj @@ -247,7 +247,7 @@ ); runOnlyForDeploymentPostprocessing = 0; shellPath = /bin/sh; - shellScript = "mkdir -p \"${DERIVED_FILE_DIR}\"\npython3 \"${SRCROOT}/scripts/make_frame_bundle.py\" \"${DERIVED_FILE_DIR}/frame-bundle.tar.gz\" > \"${DERIVED_FILE_DIR}/frame-bundle.version\"\nmkdir -p \"${TARGET_BUILD_DIR}/${UNLOCALIZED_RESOURCES_FOLDER_PATH}\"\ncp \"${DERIVED_FILE_DIR}/frame-bundle.tar.gz\" \"${DERIVED_FILE_DIR}/frame-bundle.version\" \"${TARGET_BUILD_DIR}/${UNLOCALIZED_RESOURCES_FOLDER_PATH}/\"\n"; + shellScript = "set -e\nmkdir -p \"${DERIVED_FILE_DIR}\"\n# Written aside and moved into place together, so a failed download can't leave a stale pair.\npython3 \"${SRCROOT}/scripts/make_frame_bundle.py\" \"${DERIVED_FILE_DIR}/frame-bundle.new.tar.gz\" > \"${DERIVED_FILE_DIR}/frame-bundle.new.version\"\nmv \"${DERIVED_FILE_DIR}/frame-bundle.new.tar.gz\" \"${DERIVED_FILE_DIR}/frame-bundle.tar.gz\"\nmv \"${DERIVED_FILE_DIR}/frame-bundle.new.version\" \"${DERIVED_FILE_DIR}/frame-bundle.version\"\nmkdir -p \"${TARGET_BUILD_DIR}/${UNLOCALIZED_RESOURCES_FOLDER_PATH}\"\ncp \"${DERIVED_FILE_DIR}/frame-bundle.tar.gz\" \"${DERIVED_FILE_DIR}/frame-bundle.version\" \"${TARGET_BUILD_DIR}/${UNLOCALIZED_RESOURCES_FOLDER_PATH}/\"\n"; }; /* End PBXShellScriptBuildPhase section */ diff --git a/ios/project.yml b/ios/project.yml index 2b9743b..30a5aa8 100644 --- a/ios/project.yml +++ b/ios/project.yml @@ -52,8 +52,12 @@ targets: - name: Pack the Frame bundle # The server, headset helpers and catalogue, as the app copies them to the Frame. script: | + set -e mkdir -p "${DERIVED_FILE_DIR}" - python3 "${SRCROOT}/scripts/make_frame_bundle.py" "${DERIVED_FILE_DIR}/frame-bundle.tar.gz" > "${DERIVED_FILE_DIR}/frame-bundle.version" + # Written aside and moved into place together, so a failed download can't leave a stale pair. + python3 "${SRCROOT}/scripts/make_frame_bundle.py" "${DERIVED_FILE_DIR}/frame-bundle.new.tar.gz" > "${DERIVED_FILE_DIR}/frame-bundle.new.version" + mv "${DERIVED_FILE_DIR}/frame-bundle.new.tar.gz" "${DERIVED_FILE_DIR}/frame-bundle.tar.gz" + mv "${DERIVED_FILE_DIR}/frame-bundle.new.version" "${DERIVED_FILE_DIR}/frame-bundle.version" mkdir -p "${TARGET_BUILD_DIR}/${UNLOCALIZED_RESOURCES_FOLDER_PATH}" cp "${DERIVED_FILE_DIR}/frame-bundle.tar.gz" "${DERIVED_FILE_DIR}/frame-bundle.version" "${TARGET_BUILD_DIR}/${UNLOCALIZED_RESOURCES_FOLDER_PATH}/" basedOnDependencyAnalysis: false diff --git a/tests/test_input.py b/tests/test_input.py index bf38bba..f071ee9 100644 --- a/tests/test_input.py +++ b/tests/test_input.py @@ -144,7 +144,8 @@ class Bundled(unittest.TestCase): def ssh(remote, stdin=None, timeout=30, text=True): calls.append((remote, stdin)) - return "yes\n" if remote.startswith("{ test -x") and frame_has else "" + answer = "yes\n" if remote.startswith("{ test -x") and frame_has else "" + return answer if text else answer.encode() old = self.server.ssh, self.server.KDECONNECT, self.server.LOCAL self.server.ssh, self.server.KDECONNECT, self.server.LOCAL = ssh, folder, False try: @@ -157,7 +158,8 @@ class Bundled(unittest.TestCase): folder, calls, packages = self.deliver(frame_has=True) self.assertEqual(folder, "") self.assertEqual(len(calls), 1) - self.assertEqual(calls[0][1], "".join(f"{sha} {name}\n" for name, sha in packages)) + # Bytes: on Windows a text pipe would send CRLF and never match. + self.assertEqual(calls[0][1], "".join(f"{sha} {name}\n" for name, sha in packages).encode()) def test_copies_each_package_over_ssh(self): folder, calls, packages = self.deliver(frame_has=False) @@ -170,6 +172,66 @@ class Bundled(unittest.TestCase): with self.assertRaises(self.server.Failure): self.deliver(frame_has=False, damaged=True) + def lifecycle(self, agent_lines, deliver=None): + """An InputAgent whose ssh and agent are fakes; returns it and the folders each launch used.""" + launches = [] + test = self + + class Agent(self.server.InputAgent): + def deliver(self, report, force=False): + if deliver: + deliver() + return "~/copied" if force else "" + + def command(self, folder=""): + launches.append(folder) + return "agent" + + class Proc: + stdin = None + + def __init__(self, lines): + self.stdout = iter(lines) + + def wait(self): + return 0 + + def poll(self): + return None + + def terminate(self): + pass + + def popen(*a, **k): + return Proc(agent_lines.pop(0) if agent_lines else [b'{"state": "ready"}\n']) + old = self.server.ensure_master, self.server.subprocess.Popen + self.server.ensure_master, self.server.subprocess.Popen = lambda: None, popen + test.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 + + 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._launch(agent.generation) + self.assertEqual(launches, ["", "~/copied"]) + self.assertNotIn("reach", agent.status.get("message", "")) + + def test_start_after_stop_during_the_copy_still_starts(self): + gate, entered = threading.Event(), threading.Event() + agent, launches = self.lifecycle([], deliver=lambda: (entered.set(), gate.wait(5))) + agent.start() + self.assertTrue(entered.wait(5)) + agent.stop() + entered.clear() + agent.start() # while the first launch is still copying + self.assertTrue(entered.wait(5), "the second start didn't launch") + gate.set() + for _ in range(100): + if len(launches) == 2: + break + time.sleep(0.02) + self.assertEqual(len(launches), 2) # the stopped launch ran its agent too, then ended it + @unittest.skipIf(sys.platform == "win32", "the agent runs on the Frame (Linux)") class AgentInstall(unittest.TestCase): @@ -221,6 +283,24 @@ class AgentInstall(unittest.TestCase): with self.assertRaisesRegex(RuntimeError, "doesn't include"): self.agent.install(self.dir, []) + def test_leaves_another_devices_running_copy_alone(self): + # A different build is running for another device: use it, don't stop it to reinstall. + self.agent.listening, old = (lambda: True), self.agent.listening + self.agent.our_daemons, old_ours = (lambda: [123]), self.agent.our_daemons + try: + self.agent.ensure_daemon("", [("new.pkg.tar.zst", "1" * 64)]) + finally: + self.agent.listening, self.agent.our_daemons = old, old_ours + self.assertFalse(self.agent.ROOT.exists()) + + def test_asks_for_packages_it_was_not_sent(self): + self.agent.listening, old = (lambda: False), self.agent.listening + try: + with self.assertRaises(self.agent.NeedPackages): + self.agent.ensure_daemon("", [("new.pkg.tar.zst", "1" * 64)]) + finally: + self.agent.listening = old + def test_server_and_agent_agree_on_the_stamp(self): import server packages = [("a.pkg.tar.zst", "1" * 64), ("b.pkg.tar.zst", "2" * 64)] diff --git a/ui/frame_input_agent.py b/ui/frame_input_agent.py index 776bc0f..e7c1785 100644 --- a/ui/frame_input_agent.py +++ b/ui/frame_input_agent.py @@ -16,7 +16,7 @@ argv: client id, client name, the folder holding the packages, and a JSON list of [file, sha256] naming them (see frame/kdeconnect/packages.json). Status goes to stdout, one JSON object per line: -{"state": "installing" | "starting" | "pairing" | "ready" | "error", ...}. +{"state": "installing" | "starting" | "pairing" | "ready" | "error" | "need-packages", ...}. Standard library only: this runs on the Frame's own Python. """ @@ -153,10 +153,21 @@ def stop_daemon(): time.sleep(0.1) +class NeedPackages(Exception): + """This build of KDE Connect isn't unpacked and the server didn't send it (it thought it was there).""" + + def ensure_daemon(folder, packages): - """Start KDE Connect: the Frame's own if it ever has one, else ours, unpacked first if needed.""" + """Start KDE Connect: the Frame's own if it ever has one, else ours, unpacked first if needed. + + A copy from another Frame Control version that another device is using right + now is left running and used as it is (they speak the same protocol); it's + replaced the next time nobody is using it. + """ system = SYSTEM_DAEMON.exists() - if not system and not installed(packages): + if not system and not installed(packages) and not (listening() and our_daemons()): + if not folder: + raise NeedPackages() install(folder, packages) if listening(): return @@ -373,6 +384,9 @@ def run(client, name, folder, packages): stop_daemon() ensure_daemon(folder, packages) link = connect(device, cert, key, name) + except NeedPackages: + say("need-packages") # the server copies them and starts again + return 1 except (OSError, RuntimeError, subprocess.SubprocessError) as e: say("error", message=str(e)) return 1 diff --git a/ui/server.py b/ui/server.py index dd6b7db..f48d49a 100755 --- a/ui/server.py +++ b/ui/server.py @@ -607,7 +607,8 @@ class InputAgent: def __init__(self, source=HERE / "frame_input_agent.py", packages=None): self.source, self.proc, self.lock = source, None, threading.Lock() self.packages = kdeconnect_packages() if packages is None else packages - self.status, self.launching, self.generation = {"state": "off"}, False, 0 + # generation counts stop()s; launching is the generation a launch is under way for. + self.status, self.launching, self.generation = {"state": "off"}, None, 0 def command(self, folder=""): code = base64.b64encode(self.source.read_bytes()).decode() @@ -617,19 +618,23 @@ class InputAgent: + f" {shlex.quote(client)} {shlex.quote(INPUT_NAME)} {shlex.quote(folder)}" + f" {shlex.quote(json.dumps(self.packages, separators=(',', ':')))}") - def deliver(self, report): + def deliver(self, report, force=False): """Where the agent finds the packages on the Frame, copying them there first if needed. On the Frame itself (the iPhone app) they came with the bundle. Otherwise they - go over the SSH connection, only if the Frame doesn't already have them unpacked. + go over the SSH connection, unless the Frame already has them unpacked (`force`: + the agent found it didn't after all). """ if LOCAL: return str(KDECONNECT / "packages") - stamp = kdeconnect_stamp(self.packages) - have = ssh(f"{{ test -x /usr/lib/kdeconnectd || cmp -s - {KDECONNECT_HOME}/root/.frame-control-packages; }}" - " && echo yes || true", stdin=stamp, timeout=20).strip() - if have == "yes" or not self.packages: + if not self.packages: return "" # nothing to copy (the agent says so if it needed them) + if not force: + # Bytes, so Windows doesn't turn the stamp's line ends into CRLF. + have = ssh(f"{{ test -x /usr/lib/kdeconnectd || cmp -s - {KDECONNECT_HOME}/root/.frame-control-packages; }}" + " && echo yes || true", stdin=kdeconnect_stamp(self.packages).encode(), text=False, timeout=20) + if have.strip() == b"yes": + return "" client = "".join(c for c in input_client() if c.isalnum() or c in "-_")[:64] or "default" folder = f"{KDECONNECT_HOME}/incoming/{client}" ssh(f"mkdir -p {folder}", timeout=20) @@ -638,21 +643,23 @@ class InputAgent: 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.split('-')[0]})") + 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}", + # Its own temporary name: two starts can be copying at once. + ssh(f"cd {folder} && cat > {quoted}.part$$ && mv {quoted}.part$$ {quoted}", stdin=path.read_bytes(), text=False, timeout=600) return f"~/{folder}" def start(self): with self.lock: - if self.launching or (self.proc and self.proc.poll() is None): + if self.launching == self.generation or (self.proc and self.proc.poll() is None): return - self.launching, self.status = True, {"state": "starting"} + # (A launch from before a stop() may still be finishing; it ends itself.) + self.launching, self.status = self.generation, {"state": "starting"} generation = self.generation threading.Thread(target=self._launch, args=(generation,), daemon=True).start() - def _launch(self, generation): + def _launch(self, generation, force=False): def report(message): with self.lock: if self.generation == generation: @@ -660,35 +667,50 @@ class InputAgent: errors = tempfile.TemporaryFile() try: ensure_master() - folder = self.deliver(report) + folder = self.deliver(report, force) proc = subprocess.Popen([*SSH, FRAME, self.command(folder)], stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=errors) except (Failure, OSError) as e: message = str(e) friendly = unreachable(message) with self.lock: - self.launching = False + if self.launching == generation: + self.launching = None if self.generation == generation: self.status = {"state": "error", "message": friendly or message, **({"offline": True} if friendly else {})} return _live_tunnels.add(proc) with self.lock: - self.launching = False + if self.launching == generation: + self.launching = None stale = self.generation != generation if not stale: self.proc = proc if stale: # turned off meanwhile proc.terminate() - self._watch(proc, errors) + if self._watch(proc, errors, retry=not force) 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: + if self.proc is not proc or self.generation != generation: + return + self.proc, self.launching, self.status = None, generation, {"state": "starting"} + self._launch(generation, force=True) - def _watch(self, proc, errors): + 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 for line in proc.stdout: try: status = json.loads(line) except ValueError: continue if isinstance(status, dict) and isinstance(status.get("state"), str): + if status["state"] == "need-packages": + wanted = True + continue with self.lock: if self.proc is proc: self.status = status @@ -697,10 +719,13 @@ class InputAgent: errors.seek(0) detail = strip_ansi(errors.read().decode(errors="replace")).strip() with self.lock: - if self.proc is proc and self.status.get("state") != "error": + if self.proc is proc and self.status.get("state") != "error" and not (wanted and retry): message = detail.splitlines()[-1] if detail else "The connection to the Frame ended" + if wanted: + 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 def send(self, events): """Forward events if the agent is ready; start it if it isn't running.