mirror of
https://github.com/saphid/frame-control.git
synced 2026-10-06 01:00:18 +02:00
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) <noreply@anthropic.com>
This commit is contained in:
1 parent
13eb65603b
commit
18334b5989
9 files changed
+241
-64
No files matched your search
+1
-1
@@ -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
|
||||
|
||||
+16
-3
@@ -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"
|
||||
}
|
||||
|
||||
|
||||
@@ -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")
|
||||
|
||||
+22
-1
@@ -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]
|
||||
|
||||
+32
-2
@@ -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)
|
||||
|
||||
+76
-19
@@ -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
|
||||
|
||||
|
||||
|
||||
+62
-18
@@ -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()
|
||||
|
||||
+13
-1
@@ -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) {
|
||||
|
||||
+18
-18
@@ -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)
|
||||
|
||||
|
||||
Reference in new issue
Block a user