From 18334b5989846138232478298daf3c16f8943233 Mon Sep 17 00:00:00 2001 From: saphid <4596216+saphid@users.noreply.github.com> Date: Mon, 28 Sep 2026 21:38:41 +1000 Subject: [PATCH] Devices: fixes from review round 2 - Switching headsets is serialized with the start of any install; background work counts as running from before its thread starts. use() reroutes every command at once and makes ensure() wait for the new headset. - A new user or port reroutes commands even if the attempt then fails. - ~/.ssh/config edits take a lock file shared with Set Up Connection (frame_connect.py and connect.sh, which now also writes atomically). - A finished attempt no longer writes its older settings over a change Set Up Connection made meanwhile. - Pin edits are locked and swapped atomically. - A bare alias behind ProxyJump/ProxyCommand is left to ssh to reach. - The page drops answers about the previous headset after a switch; the header switcher takes clicks in the macOS title bar. Co-Authored-By: Claude Opus 5.5 (1M context) --- app/main.js | 2 +- scripts/connect.sh | 19 +++++++-- tests/test_devices.py | 2 +- tests/test_link.py | 23 ++++++++++- ui/frame_connect.py | 34 +++++++++++++++- ui/frame_devices.py | 95 ++++++++++++++++++++++++++++++++++--------- ui/frame_link.py | 80 ++++++++++++++++++++++++++++-------- ui/index.html | 14 ++++++- ui/server.py | 36 ++++++++-------- 9 files changed, 241 insertions(+), 64 deletions(-) diff --git a/app/main.js b/app/main.js index aac8d35..b159ac2 100644 --- a/app/main.js +++ b/app/main.js @@ -189,7 +189,7 @@ async function restartServer() { // On macOS the page's sticky header becomes the title bar, clear of the traffic lights. const CHROME_CSS = IS_MAC && ` header { padding-left: 92px !important; -webkit-app-region: drag; user-select: none; } - header a, header button, header input, header .chip { -webkit-app-region: no-drag; } + header a, header button, header input, header select, header .chip { -webkit-app-region: no-drag; } `; // Restart Server can start a new load while an older one is still waiting for diff --git a/scripts/connect.sh b/scripts/connect.sh index 00bcb41..477020b 100755 --- a/scripts/connect.sh +++ b/scripts/connect.sh @@ -101,10 +101,22 @@ make_key() { # path type comment [extra ssh-keygen args] fi } -# Checks each step itself: pair_with_devkit calls this from an `elif`, where set -e is off. +# Takes the lock Frame Control uses to edit ~/.ssh/config (ui/frame_devices.py), so a +# running app and this script never write over each other's change. write_config() { + local lockfd="" rc + zmodload zsh/system 2>/dev/null + touch "$CONFIG.frame-control.lock" 2>/dev/null + zsystem flock -t 30 -f lockfd "$CONFIG.frame-control.lock" 2>/dev/null || lockfd="" + write_config_locked; rc=$? + [[ -n "$lockfd" ]] && zsystem flock -u "$lockfd" + return $rc +} + +# Checks each step itself: pair_with_devkit calls this from an `elif`, where set -e is off. +write_config_locked() { touch "$CONFIG" && chmod 600 "$CONFIG" || return 1 - local tmp + local tmp new="$CONFIG.frame-control.$$" tmp=$(mktemp) || return 1 # Drop any previous managed block, then PREPEND a fresh one: ssh uses the first # value it sees per option, so this block must precede any other "Host frame" @@ -126,7 +138,8 @@ write_config() { print -r -- "Host *" print -r -- "$END_MARK" cat "$tmp" - } > "$CONFIG" || { print -u2 "!! Writing $CONFIG failed; its previous contents are in $tmp"; return 1; } + } > "$new" && chmod 600 "$new" && mv -f "$new" "$CONFIG" \ + || { rm -f "$new"; print -u2 "!! Writing $CONFIG failed; its previous contents are in $tmp"; return 1; } rm -f "$tmp" } diff --git a/tests/test_devices.py b/tests/test_devices.py index a115071..abad757 100644 --- a/tests/test_devices.py +++ b/tests/test_devices.py @@ -179,7 +179,7 @@ class ConfigRewrite(Base): t.join() blocks = fd.parse_blocks((self.ssh / "config").read_text()) self.assertEqual([b["hostname"] for b in blocks], ["10.0.0.14", "10.0.1.14"]) - self.assertEqual([p.name for p in self.ssh.iterdir() if "frame-control." in p.name], []) # no temp files left + self.assertEqual([p.name for p in self.ssh.iterdir() if "frame-control." in p.name and not p.name.endswith(".lock")], []) # no temp files left def test_zone_is_escaped_and_read_back(self): fd.rewrite_block("frame", hostname="fe80::1%en0") diff --git a/tests/test_link.py b/tests/test_link.py index c2ea013..1f46be2 100644 --- a/tests/test_link.py +++ b/tests/test_link.py @@ -208,13 +208,31 @@ class Connecting(unittest.TestCase): def test_a_bare_alias_lets_ssh_config_decide(self): self.link.override = "frame-bare" self.hosts({"frame-bare": "ok"}) # the stand-in ssh has no config: the alias is the host - with mock.patch.object(fl, "ssh_g", return_value=("localhost", self.port, "tester")): + with mock.patch.object(fl, "ssh_g", return_value=("localhost", self.port, "tester", False)): self.link.connect(["start"]) s = self.link.snapshot() self.assertEqual(s["phase"], "connected", s["error"]) self.assertTrue(s["device"]["transient"]) self.assertEqual(self.routes[-1], ("frame-bare", [])) # no HostName override, ssh's own known_hosts + def test_a_bare_alias_behind_a_jump_host_is_left_to_ssh(self): + self.link.override = "frame-jump" + self.hosts({"frame-jump": "ok"}) + with mock.patch.object(fl, "ssh_g", return_value=("10.99.99.99", 22, "tester", True)): + self.link.connect(["start"]) + s = self.link.snapshot() + self.assertEqual(s["phase"], "connected", s["error"]) + self.assertEqual(s["via"]["why"], "through a jump host") + + def test_changing_the_port_reroutes_even_if_the_attempt_fails(self): + d = self.device("localhost") + self.hosts({"localhost": "ok"}) + self.link.connect(["start"]) + self.reg.update_device(d["id"], port=1) # nothing listens there + self.link.connect(["switch"]) + self.assertEqual(self.link.snapshot()["phase"], "failed") + self.assertIn("Port=1", self.routes[-1][1]) + def test_test_now_checks_every_address_without_touching_the_connection(self): d = self.device("127.0.0.1", "localhost", "nothing.invalid") (self.dir / "ssh" / "frame-control_known_hosts").write_text(f"frame-control-{d['id']} ssh-ed25519 AAAA\n") @@ -234,6 +252,9 @@ class Connecting(unittest.TestCase): other = self.reg.add_device("frame-other", port=self.port) self.reg.add_address(other["id"], "nothing.invalid") self.link.use(other["id"]) + self.assertEqual(self.routes[-1][0], "frame-other") # at once, before any attempt + self.assertEqual(self.link.snapshot()["phase"], "connecting") + self.assertFalse(self.link.alive()) # so ensure() waits instead of using the old master self.link.connect(["switch"]) self.assertEqual(self.link.snapshot()["phase"], "failed") alias, opts = self.routes[-1] diff --git a/ui/frame_connect.py b/ui/frame_connect.py index 6aefc4b..a459a32 100644 --- a/ui/frame_connect.py +++ b/ui/frame_connect.py @@ -284,10 +284,40 @@ def config_block(host, port=22, user=FRAME_USER): " IdentitiesOnly yes", " ServerAliveInterval 30", "Host *", END] +class config_lock: + """The lock Frame Control takes to edit ~/.ssh/config (frame_devices.file_lock), so a + running app and this setup never write over each other's change.""" + + def __enter__(self): + self.fh = open(SSH_DIR / "config.frame-control.lock", "a+") + for _ in range(300): + try: + if os.name == "nt": + import msvcrt + self.fh.seek(0) + msvcrt.locking(self.fh.fileno(), msvcrt.LK_NBLCK, 1) + else: + import fcntl + fcntl.lockf(self.fh, fcntl.LOCK_EX | fcntl.LOCK_NB) + return self + except OSError: + time.sleep(0.1) + return self # 30 s: go ahead rather than fail the setup + + def __exit__(self, *exc): + self.fh.close() # closing releases the lock + return False + + def write_config(host, port=22, user=FRAME_USER): + make_ssh_dir() + with config_lock(): + _write_config(host, port, user) + + +def _write_config(host, port, user): """Replace our managed block and put it first: ssh uses the first value it sees per option. The trailing "Host *" returns the rest of the file to global scope.""" - make_ssh_dir() old = CONFIG.read_text(encoding="utf-8") if CONFIG.exists() else "" kept, skip = [], False for line in old.splitlines(): @@ -298,7 +328,7 @@ def write_config(host, port=22, user=FRAME_USER): elif not skip: kept.append(line) block = config_block(host, port, user) - tmp = CONFIG.with_name("config.frame-control.tmp") + tmp = CONFIG.with_name(f"config.frame-control.{os.getpid()}.tmp") tmp.write_text("\n".join(block + kept) + "\n", encoding="utf-8") if os.name != "nt": tmp.chmod(0o600) diff --git a/ui/frame_devices.py b/ui/frame_devices.py index b207ad9..2fd49f1 100644 --- a/ui/frame_devices.py +++ b/ui/frame_devices.py @@ -27,6 +27,7 @@ different device answering at a remembered IP is caught. Python stdlib only. """ +import contextlib import copy import json import os @@ -179,9 +180,48 @@ def read_config(path=None): return "" -# One edit of ~/.ssh/config at a time from this app (the connector and the page can -# both want one); _edit_config also notices another program writing in between. +# One edit of ~/.ssh/config at a time: between this app's threads (_config_lock) and +# with Set Up Connection (frame_connect.py and scripts/connect.sh take the same lock +# file). _edit_config also notices any other program writing in between. _config_lock = threading.Lock() +LOCK_NAME = "config.frame-control.lock" + + +@contextlib.contextmanager +def file_lock(path, timeout=30): + """An exclusive lock on `path` (created if need be) shared with other processes: + POSIX record locks (what zsh's `zsystem flock` takes), or msvcrt on Windows.""" + path.parent.mkdir(parents=True, exist_ok=True) + fh = open(path, "a+") + try: + deadline = time.monotonic() + timeout + while True: + try: + if frame_host.WINDOWS: + import msvcrt + fh.seek(0) + msvcrt.locking(fh.fileno(), msvcrt.LK_NBLCK, 1) + else: + import fcntl + fcntl.lockf(fh, fcntl.LOCK_EX | fcntl.LOCK_NB) + break + except OSError: + if time.monotonic() > deadline: + raise OSError(f"{path} stayed locked (is Set Up Connection running?)") + time.sleep(0.1) + yield + finally: + try: + if frame_host.WINDOWS: + import msvcrt + fh.seek(0) + msvcrt.locking(fh.fileno(), msvcrt.LK_UNLCK, 1) + else: + import fcntl + fcntl.lockf(fh, fcntl.LOCK_UN) + except OSError: + pass + fh.close() def _write_config(path, text, expected): @@ -211,7 +251,7 @@ def _write_config(path, text, expected): def _edit_config(path, change): """Apply change(lines) -> new lines or None to the file, retrying if another program wrote it meanwhile. -> True if the file changed.""" - with _config_lock: + with _config_lock, file_lock(path.with_name(LOCK_NAME)): for _ in range(5): text = read_config(path) new = change(text.splitlines()) @@ -283,6 +323,10 @@ def _keygen(*args): return None +def _pin_lock(target): + return file_lock(target.with_name(target.name + ".lock")) + + def pinned(device_id, path=None): """Whether a key is saved for the device. The app's own entries are plain text (it passes HashKnownHosts=no), but ask ssh-keygen too in case one was hashed.""" @@ -322,11 +366,11 @@ def seed_pin(device_id, hosts, port=22, sources=None, path=None): keys.append(f"{name} {f[1]} {f[2]}") if keys: target = Path(path or known_hosts()) - target.parent.mkdir(parents=True, exist_ok=True) - with open(target, "a", encoding="utf-8") as fh: - fh.write("\n".join(dict.fromkeys(keys)) + "\n") - if not frame_host.WINDOWS: - target.chmod(0o600) + with _pin_lock(target): + with open(target, "a", encoding="utf-8") as fh: + fh.write("\n".join(dict.fromkeys(keys)) + "\n") + if not frame_host.WINDOWS: + target.chmod(0o600) return True return False @@ -336,17 +380,30 @@ def forget_pin(device_id, path=None): trusts whatever key the headset shows, as a first connection does.""" target = Path(path or known_hosts()) name = host_key_alias(device_id) - lines = _pin_lines(target) - kept = [line for line in lines if not (line.strip() and name in line.split(None, 1)[0].split(","))] - removed = kept != lines - if removed: - target.write_text("".join(line + "\n" for line in kept), encoding="utf-8") - if target.is_file() and pinned(device_id, target): # a hashed entry: ssh-keygen finds it - r = _keygen("-R", name, "-f", str(target)) - removed = removed or bool(r and r.returncode == 0) - old = target.with_name(target.name + ".old") # ssh-keygen -R leaves a backup - if old.exists(): - old.unlink() + with _pin_lock(target): + # ssh itself may append a first-seen key meanwhile (accept-new): swap the file + # only if it still holds what was read, else read it again. + removed = False + for _ in range(5): + text = "\n".join(_pin_lines(target)) + lines = text.splitlines() + kept = [line for line in lines if not (line.strip() and name in line.split(None, 1)[0].split(","))] + if kept == lines: + break + fd_, tmp = tempfile.mkstemp(prefix=target.name + ".", dir=str(target.parent)) + with os.fdopen(fd_, "w", encoding="utf-8") as fh: + fh.write("".join(line + "\n" for line in kept)) + if "\n".join(_pin_lines(target)) == text: + os.replace(tmp, target) + removed = True + break + os.unlink(tmp) + if target.is_file() and pinned(device_id, target): # a hashed entry: ssh-keygen finds it + r = _keygen("-R", name, "-f", str(target)) + removed = removed or bool(r and r.returncode == 0) + old = target.with_name(target.name + ".old") # ssh-keygen -R leaves a backup + if old.exists(): + old.unlink() return removed diff --git a/ui/frame_link.py b/ui/frame_link.py index f23fa02..ed88e7a 100644 --- a/ui/frame_link.py +++ b/ui/frame_link.py @@ -68,7 +68,8 @@ def now(): def ssh_g(alias): - """(hostname, port, user) from `ssh -G ALIAS`, for a headset that's only an ssh alias.""" + """(hostname, port, user, proxied) from `ssh -G ALIAS`, for a headset that's only an + ssh alias. proxied: it goes through ProxyJump or ProxyCommand, so only ssh can reach it.""" try: out = subprocess.run(["ssh", "-G", alias], capture_output=True, stdin=subprocess.DEVNULL, text=True, timeout=10).stdout @@ -77,10 +78,11 @@ def ssh_g(alias): got = {} for line in out.splitlines(): k, _, v = line.partition(" ") - if k in ("hostname", "port", "user") and k not in got: + if k in ("hostname", "port", "user", "proxyjump", "proxycommand") and k not in got: got[k] = v.strip() port = int(got["port"]) if got.get("port", "").isdigit() else 22 - return got.get("hostname") or alias, port, got.get("user") + proxied = any(got.get(k) not in (None, "", "none") for k in ("proxyjump", "proxycommand")) + return got.get("hostname") or alias, port, got.get("user"), proxied def probe(host, port, timeout=PROBE_TIMEOUT, update=None): @@ -246,11 +248,27 @@ class Link: self.state["phase"] != "connecting"), wait) def use(self, device_id): - """Switch to another headset.""" + """Switch to another headset. Commands go to it from now on (never the last one), + and wait in ensure() for the connector to reach it.""" self.reg.set_active(device_id) self.override = None + device = self.active_device() + self.apply(device["alias"], self.first_route(device)) + with self.cond: + self.state["phase"] = "connecting" + self.kicks.append("switch") + self.version += 1 + self.cond.notify_all() self.devices_changed() - self.kick("switch") + + def first_route(self, device): + """Where commands go before any address has answered: the first one, with the + headset's own pinned identity, so nothing reaches another device meanwhile.""" + return self.host_opts(device, device["addresses"][0]["host"] if device["addresses"] else None) + + @staticmethod + def route_key(device): + return device["id"], device.get("user"), device.get("port") def lost(self, message): """A command couldn't reach the headset (Windows has no master to watch).""" @@ -357,13 +375,12 @@ class Link: self.last_attempt = now() self.close_master() device = self.active_device() - if device["id"] != self.routed: - # Another headset: nothing may go on reaching the last one, even if this one - # never answers. Its first address (and its own pinned identity) until one does. - first = device["addresses"][0]["host"] if device["addresses"] else None - self.alias, self.opts = device["alias"], self.host_opts(device, first) + if self.route_key(device) != self.routed: + # Another headset, or a new user or port: nothing may go on using the old + # route, even if this attempt fails. + self.alias, self.opts = device["alias"], self.first_route(device) self.apply(self.alias, self.opts) - self.routed = device["id"] + self.routed = self.route_key(device) with self.cond: self.state.update(phase="connecting", reason=why, device=self.public_device(device), via=None, error=None, retry_at=None, attempt=self.state["attempt"] + 1, started=now(), @@ -432,12 +449,21 @@ class Link: # 2. find the headset self.stage("find", "active") port = device.get("port") or 22 - if device.get("transient") or not device["addresses"]: - host, port, user = ssh_g(device["alias"]) + bare = device.get("transient") or not device["addresses"] + if bare: + host, port, user, proxied = ssh_g(device["alias"]) if user and not device.get("user"): device["user"] = user - ranked = [({"host": host, "kind": frame_network.guess_kind(host), "label": "from ~/.ssh/config"}, - "from ~/.ssh/config")] + a = {"host": host, "kind": frame_network.guess_kind(host), "label": "from ~/.ssh/config"} + if proxied: + # Reached through a jump host: only ssh itself can find it. + with self.cond: + self.state["probes"] = [dict(a, why="through a jump host", state="answered", ip=None, rtt_ms=None, + detail="ssh's ProxyJump or ProxyCommand connects")] + self.stage("find", "done", f"{device['alias']} goes through a jump host; ssh finds it") + return self.handshake(device, a, {"ip": None, "rtt_ms": None}, device.get("user") or user) == "ok" \ + and self.finish_bare(a, net) + ranked = [(a, "from ~/.ssh/config")] else: ranked = frame_devices.order_addresses(device["addresses"], net.get("id"), bool(ts.get("up"))) with self.cond: @@ -502,6 +528,11 @@ class Link: self.fail("find", message, raw) return False + def finish_bare(self, a, net): + self.publish(via={"host": a["host"], "kind": a["kind"], "ip": None, "rtt_ms": None, + "why": "through a jump host", "network": net.get("id"), "network_name": net["name"]}) + return True + @staticmethod def pick(results, tried, done, deadline=None): """The next address to use: the best-ranked answer once every better-ranked @@ -534,12 +565,25 @@ class Link: if device.get("transient"): return self.reg.record_success(device["id"], host, net.get("id"), rtt) + # Set Up Connection may have changed the block while this attempt ran: take that + # in (and reconnect) rather than writing this attempt's older settings over it. + self.watch_config() + try: + now_dev = self.reg.get(device["id"]) + except frame_devices.DeviceError: + return + if self.route_key(now_dev) != self.route_key(device): + return # Terminal's `ssh ALIAS` and the helper scripts use ~/.ssh/config: point it here too. try: - if frame_devices.rewrite_block(device["alias"], hostname=host, user=device["user"], - port=device["port"]): + block = next((b for b in frame_devices.parse_blocks(frame_devices.read_config()) + if b["alias"] == device["alias"]), None) + moved = block and block["hostname"] != device.get("config_host") and device.get("config_host") + if not moved and frame_devices.rewrite_block(device["alias"], hostname=host, user=device["user"], + port=device["port"]): self.config_mtime = frame_devices.ssh_config().stat().st_mtime - self.reg.set_config_host(device["id"], host) + if block and not moved: + self.reg.set_config_host(device["id"], host) except OSError: pass # not fatal: the app itself doesn't need the file self.devices_changed() diff --git a/ui/index.html b/ui/index.html index e5b977c..3e7c6ee 100644 --- a/ui/index.html +++ b/ui/index.html @@ -920,7 +920,11 @@ async function saveToDevice(items) { } const savesToDevice = () => !!(window.frameApp && window.frameApp.saveImages); +// Bumped when the app moves to another headset: answers to requests made before +// then are about the old one, so they're dropped rather than shown. +let devGen = 0; async function api(path, body) { + const gen = devGen; const opts = body === undefined ? { headers: {"X-Frame-UI": UI_KEY} } : { method: "POST", headers: {"Content-Type": "application/json", "X-Frame-UI": UI_KEY}, body: JSON.stringify(body) }; let r; @@ -930,6 +934,11 @@ async function api(path, body) { : "Frame Control's local server isn't running. Restart the app (Frame → Restart Server)."); } const data = await r.json().catch(() => ({ error: `HTTP ${r.status}` })); + if (gen !== devGen && !path.startsWith("/api/devices") && !path.startsWith("/api/connection")) { + const err = new Error("Switched headsets"); + err.offline = true; // panels show "Waiting for the Frame" and reload from the new one + throw err; + } if (!r.ok) { const err = new Error(data.error || `HTTP ${r.status}`); err.offline = !!data.offline; @@ -2522,7 +2531,9 @@ function onConnection(s) { log(`Now using ${s.device.name}`, "ok"); state = null; online = null; - link.reload = true; // everything on the page was the other headset's + devGen++; // answers still on their way are about the other headset + refreshing = null; // and the next status check must be a new one + link.reload = true; // everything on the page was the other headset's $("battChip").hidden = true; } if (s.phase !== link.phase || devId !== link.device) { @@ -2616,6 +2627,7 @@ $("devSel").onchange = e => useDevice(e.target.value); async function useDevice(id) { const d = dv.list.find(x => x.id === id); if (!d || d.active) return; + devGen++; // from here on, answers about the last headset are dropped await act(`Switch to ${d.name}`, () => api("/api/devices", { action: "use", id })); } async function devAction(label, body, btn) { diff --git a/ui/server.py b/ui/server.py index 949d300..bc6c104 100755 --- a/ui/server.py +++ b/ui/server.py @@ -114,16 +114,18 @@ def working(): _work[0] -= 1 -def busy_while(fn): - def run(*args, **kwargs): - with working(): - return fn(*args, **kwargs) - return run - - -def busy(): +def busy_thread(fn, *args): + """A thread that counts as work from before it starts until it ends.""" with _work_lock: - return _work[0] + _work[0] += 1 + + def run(): + try: + fn(*args) + finally: + with _work_lock: + _work[0] -= 1 + return threading.Thread(target=run, daemon=True) APPID = re.compile(r"^\d{1,10}$") FLATPAK_ID = re.compile(r"^[A-Za-z0-9][A-Za-z0-9_-]*(\.[A-Za-z0-9_-]+){2,}$") @@ -215,8 +217,7 @@ def start_job(label, work): def run(): fields = {} try: - with working(): - result = work() + result = work() fields = {"message": result.get("message") or f"{label}: done", "result": result} except (Failure, frame_android.FrameError) as e: fields = {"error": unreachable(str(e)) or str(e)} @@ -226,7 +227,7 @@ def start_job(label, work): with _jobs_lock: _jobs[job].update(fields, done=True, time=time.time()) - threading.Thread(target=run, daemon=True).start() + busy_thread(run).start() return {"message": f"{label}…", "job": job} @@ -691,7 +692,6 @@ def stage_title(path, temp_dir=None, name=None): "token": token, "plan": frame_titles.public(plan)} -@busy_while def _run_title_install(token, entry, name, exe, runtime): def update(**fields): # the page reads jobs from other threads; change them under the lock with _titles_lock: @@ -733,8 +733,7 @@ def titles(body): "message": None, "title": None, "time": time.time()} ensure_master() opt = lambda k: str(body.get(k) or "") or None # noqa: E731 - threading.Thread(target=_run_title_install, daemon=True, - args=(token, entry, opt("name"), opt("exe"), opt("runtime"))).start() + busy_thread(_run_title_install, token, entry, opt("name"), opt("exe"), opt("runtime")).start() return {"message": f"Installing {entry['plan']['source']}", "job": token} if action not in ("launch", "remove"): raise Failure("unknown action", 400) @@ -1120,13 +1119,12 @@ def webinstall_start(body): _web_jobs.clear() _web_jobs[pid] = job # Started under the lock, so shutdown never sees a thread it can't join. - worker = threading.Thread(target=_webinstall_run, args=(plan, job), daemon=True) + worker = busy_thread(_webinstall_run, plan, job) _web_workers.add(worker) worker.start() return {"job": pid} -@busy_while def _webinstall_run(plan, job): tmp = None try: @@ -1281,7 +1279,9 @@ def devices_post(body): if not LINK: raise Failure("Headsets are managed from the computer app", 400) try: - return frame_link.devices_action(LINK, body, open_setup, busy) + # Under the work lock: nothing can start on the old headset while it switches. + with _work_lock: + return frame_link.devices_action(LINK, body, open_setup, lambda: _work[0]) except frame_devices.DeviceError as e: raise Failure(str(e), 400)