mirror of
https://github.com/saphid/frame-control.git
synced 2026-10-06 01:00:18 +02:00
Bundled KDE Connect: fixes from review
- iPhone build fails if the bundle can't be made, instead of shipping a stale archive with an empty version. - A different build that another device is using right now is left running and replaced once nobody is, rather than stopped under them. - If what's installed changed after the server looked, the agent asks for the packages (need-packages) and the server copies them and starts again. - The installed-build check sends bytes, so Windows' CRLF can't break it. - Starting again while a stopped start is still copying launches anew; concurrent copies use their own temporary names. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
1 parent
5bf1268c82
commit
75b2478a55
5 files changed
+148
-25
No files matched your search
@@ -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 */
|
||||
|
||||
|
||||
+5
-1
@@ -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
|
||||
|
||||
+82
-2
@@ -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)]
|
||||
|
||||
+17
-3
@@ -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
|
||||
|
||||
+43
-18
@@ -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.
|
||||
|
||||
Reference in new issue
Block a user