Merge main into devices: MCP, analytics, Mac view alongside several headsets

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
saphidandClaude Opus 5.5 committed 2026-09-29 10:14:50 +10:00
commit 7d6ff7f919
84 files changed
+33966 -59

No files matched your search

+186 -11
View File
@@ -25,6 +25,7 @@ import shlex
import shutil
import signal
import socket
import socketserver
import subprocess
import sys
import tempfile
@@ -38,13 +39,18 @@ from urllib.parse import parse_qs, unquote, urlparse
# sys.path, so add it for the sibling modules below.
sys.path.insert(0, str(Path(__file__).resolve().parent))
import frame_agent # noqa: E402
import frame_assistant # noqa: E402
import frame_android # noqa: E402
import frame_apk_versions # noqa: E402
import frame_catalog # noqa: E402
import frame_devices # noqa: E402
import frame_host # noqa: E402
import frame_link # noqa: E402
import frame_macview # noqa: E402
import frame_report # noqa: E402
import frame_store # noqa: E402
import frame_telemetry # noqa: E402
import frame_titles # noqa: E402
import frame_webinstall # noqa: E402
@@ -65,7 +71,7 @@ if not re.fullmatch(r"[A-Za-z0-9][A-Za-z0-9._-]*", FRAME):
sys.exit(f"FRAME_ALIAS must be a plain host alias, not {FRAME!r}")
# Reuse one SSH connection for the frequent status/screenshot calls, where ssh
# supports it (not on Windows: there every command connects on its own).
CONTROL = None if LOCAL else frame_host.control_path()
CONTROL = None if LOCAL else frame_host.control_path(private=os.environ.get("FRAME_PRIVATE_SSH") == "1")
# The ControlPath itself is per headset: the connector puts it in HOST_OPTS.
MUX = ["ssh", "-o", "BatchMode=yes"]
MUX_BASE = list(MUX)
@@ -229,8 +235,10 @@ def start_job(label, 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)}
frame_telemetry.diagnostic(f"job {label.split()[0]}", e)
except Exception as e:
fields = {"error": f"{type(e).__name__}: {e}"}
frame_telemetry.diagnostic(f"job {label.split()[0]}", e)
finally:
with _jobs_lock:
_jobs[job].update(fields, done=True, time=time.time())
@@ -290,7 +298,12 @@ def terminal(argv):
# ---- actions ---------------------------------------------------------------
def status(_body):
return json.loads(ssh("python3 -", stdin=(HERE / "frame_status.py").read_text(), timeout=20))
s = json.loads(ssh("python3 -", stdin=(HERE / "frame_status.py").read_text(), timeout=20))
osr = s.get("os") if isinstance(s, dict) else None
if isinstance(osr, dict):
frame_telemetry.frame_seen(osr.get("build"), osr.get("version"))
frame_report.frame.update(build=osr.get("build"), version=osr.get("version"))
return s
def headset_view():
@@ -478,7 +491,16 @@ def steam(body):
raise Failure("bad appid", 400)
if action not in ("install", "store"):
raise Failure("action must be install or store", 400)
return steam_frame(action, appid)
if action == "store":
return steam_frame(action, appid)
# Starts Steam's download; Steam reports the rest in the headset.
try:
res = steam_frame(action, appid)
except Failure as e:
frame_telemetry.install_finished("steam", False, error=e, steam_appid=appid)
raise
frame_telemetry.install_finished("steam", True, steam_appid=appid)
return res
def steam_search(query):
@@ -552,10 +574,16 @@ def flatpak(body):
raise Failure("bad Flatpak app ID", 400)
if action == "install":
def work():
# Per-user, so it survives SteamOS updates and needs no sudo (as install-apps.sh).
ssh("flatpak remote-add --user --if-not-exists flathub "
"https://dl.flathub.org/repo/flathub.flatpakrepo && "
f"flatpak install --user -y --noninteractive flathub {shlex.quote(app)}", timeout=1800)
start = time.time()
try:
# Per-user, so it survives SteamOS updates and needs no sudo (as install-apps.sh).
ssh("flatpak remote-add --user --if-not-exists flathub "
"https://dl.flathub.org/repo/flathub.flatpakrepo && "
f"flatpak install --user -y --noninteractive flathub {shlex.quote(app)}", timeout=1800)
except Failure as e:
frame_telemetry.install_finished("flatpak", False, time.time() - start, e, flatpak_id=app)
raise
frame_telemetry.install_finished("flatpak", True, time.time() - start, flatpak_id=app)
return {"message": f"Installed {app}"}
return start_job(f"Install {app}", work)
if action == "uninstall":
@@ -658,13 +686,37 @@ def android(body):
runtime=body.get("runtime") or "instance",
label=body.get("label"), source=body.get("source"))
name = r.get("label") or pkg
where = "" if frame_catalog.compat_db.shared() else " on this computer"
where = ("" if frame_catalog.compat_db.shared() else
" and shared it" if frame_telemetry.enabled("compat") else " on this computer")
return {"message": f"Saved your report for {name}{where}", "report": r}
except frame_android.FrameError as e:
raise Failure(str(e))
raise Failure("unknown action", 400)
# Errors that are the APK's own fault, so they belong in the compatibility
# database as install_failed. Connection trouble and the like don't.
APK_FAULTS = {"android_installer", "apk_needs_newer_android", "apk_wrong_abi"}
def apk_installed(info, meta, error, seconds):
"""Every APK install (catalogue, dropped file, web link): usage analytics, and an
install_failed report when the APK itself wouldn't install."""
pkg = (info or {}).get("package")
by_pkg = frame_catalog._cache.get("by_pkg") or {}
in_catalog = bool(pkg) and pkg in by_pkg
# Package names only for catalogue apps, which are public; a private APK's name stays here.
# No version: a local rebuild can share a catalogue app's package name but carry anything in its version.
frame_telemetry.install_finished("apk", error is None, seconds, error, catalog=in_catalog,
package=pkg if in_catalog else None)
if error is not None and pkg and frame_telemetry.categorize(error)[0] in APK_FAULTS:
frame_catalog.add_report(pkg, info.get("version"), result="install_failed", notes=str(error)[:300],
via="install", label=info.get("label"))
frame_android.install_hooks.append(apk_installed)
# ---- Sideloaded titles (Linux/Windows builds as Steam Devkit Games) --------
#
# Installing is two steps: inspect (a dropped file is uploaded and a zip
@@ -719,14 +771,18 @@ def _run_title_install(token, entry, name, exe, runtime):
with _titles_lock:
_title_jobs[token].update(fields)
start = time.time()
try:
m = frame_titles.install_plan(entry["plan"], name=name, exe=exe, runtime=runtime,
progress=lambda stage, fraction: update(stage=stage, fraction=fraction))
update(title=m, message=f"Installed {m['id']} in the Steam library ({m['runtime_label']})")
frame_telemetry.install_finished("title", True, time.time() - start, runtime=m.get("runtime"))
except frame_android.FrameError as e:
update(error=str(e))
frame_telemetry.install_finished("title", False, time.time() - start, e)
except Exception as e:
update(error=f"{type(e).__name__}: {e}")
frame_telemetry.install_finished("title", False, time.time() - start, e)
finally:
_drop_staged(entry)
update(done=True, time=time.time())
@@ -1181,10 +1237,16 @@ def _webinstall_run(plan, job):
ensure_master()
res = frame_webinstall.dispatch(path, name=plan["name"], exe=plan["exe"], progress=detail, source=plan["url"])
job["message"], job["phase"] = res["message"], "done"
if res.get("kind") != "apk": # APKs are counted by apk_installed
frame_telemetry.install_finished("web", True, kind_detail=res.get("kind"))
except Exception as e:
stage = job.get("phase") # download or install, before it becomes "error"
known = (frame_webinstall.WebInstallError, Failure, frame_android.FrameError)
job["error"] = str(e) if isinstance(e, known) else f"{type(e).__name__}: {e}"
job["phase"] = "error"
# An APK that failed to install was counted by apk_installed.
if not isinstance(e, frame_webinstall.Cancelled) and not (stage == "install" and plan.get("kind") == "apk"):
frame_telemetry.install_finished("web", False, error=e, stage=stage, kind_detail=plan.get("kind"))
finally:
with _web_lock:
job.pop("_conn", None)
@@ -1283,10 +1345,81 @@ def _sweep_one(prefix, d):
pass
POST = {"/api/android/display": android_display, "/api/android": android, "/api/titles": titles, "/api/launch": launch, "/api/steam": steam, "/api/volume": set_volume, "/api/clipboard": clipboard,
# ---- Mac in the headset (frame_macview.py) ----------------------------------
# The tunnel gets its own connection: the shared master's options would win
# over anything added after them.
macview = frame_macview.MacView(["ssh", "-o", "BatchMode=yes", "-o", "ConnectTimeout=8"],
lambda remote, stdin=None, timeout=30: ssh(remote, stdin=stdin, timeout=timeout),
FRAME, track=_live_tunnels.add)
def macview_state(query=None):
if LOCAL:
return {"available": False, "reason": "Runs on your computer, not the headset."}
try:
# Refresh restarts the helper, which picks up a permission just granted.
if (query or {}).get("restart") == ["1"]:
macview.restart_agent()
return macview.state()
except frame_macview.MacViewError as e:
raise Failure(str(e), 500)
def macview_action(body):
"""{action: show|stop|permissions, src, quality, w, h}."""
if LOCAL:
raise Failure("Runs on your computer, not the headset.", 400)
action = body.get("action")
try:
if action == "show":
w, h = body.get("w"), body.get("h")
return macview.show(str(body.get("src") or ""), str(body.get("quality") or "balanced"),
int(w) if w else None, int(h) if h else None)
if action == "stop":
return macview.stop(body.get("src") or None)
if action == "permissions":
return macview.request_permissions()
except frame_macview.MacViewError as e:
raise Failure(str(e), 502)
raise Failure("unknown action", 400)
# ---- Report a problem (frame_report.py) --------------------------------------
def report_preview(body):
"""Exactly the diagnostics a report would include, for the dialog to show first."""
return {"text": frame_report.diagnostics(body.get("activity") or (), include_logs=bool(body.get("includeLogs")))}
def report_send(body):
try:
return frame_report.send(body)
except frame_report.ReportError as e:
raise Failure(str(e))
def agent_call(body):
return frame_agent.call(sys.modules[__name__], body)
def assistant_chat(body):
return frame_assistant.chat(body, headset_view)
def agent_approval(body):
return frame_agent.approvals.decide(body.get("confirmation"), body.get("accept"))
POST = {"/api/agent/call": agent_call, "/api/agent/approval": agent_approval,
"/api/assistant/chat": assistant_chat, "/api/android/display": android_display, "/api/android": android, "/api/titles": titles, "/api/launch": launch, "/api/steam": steam, "/api/volume": set_volume, "/api/clipboard": clipboard,
"/api/flatpak": flatpak, "/api/open": open_thing, "/api/shots/save": save_shots,
"/api/webinstall/check": webinstall_check, "/api/webinstall/start": webinstall_start,
"/api/webinstall/cancel": webinstall_cancel, "/api/devices": lambda body: devices_post(body)}
"/api/webinstall/cancel": webinstall_cancel,
"/api/telemetry": frame_telemetry.update_settings, "/api/telemetry/event": frame_telemetry.page_event,
"/api/report/preview": report_preview, "/api/report": report_send, "/api/macview": macview_action,
"/api/devices": lambda body: devices_post(body)}
# ---- headsets and the connection (frame_devices.py, frame_link.py) ----------
@@ -1334,6 +1467,12 @@ def connection_state():
# ---- HTTP ------------------------------------------------------------------
def action_of(body):
"""The action a request asked for, for diagnostics: a short word, never user data."""
a = body.get("action") if isinstance(body, dict) else None
return a if isinstance(a, str) and re.fullmatch(r"[a-z]{1,20}", a) else ""
def _pipe_reader(pipe):
"""Chunks from a pipe via a thread; select() can't wait on pipes on Windows."""
chunks = queue.Queue() # unbounded: the pump never blocks, so it ends at EOF
@@ -1430,6 +1569,12 @@ class Handler(BaseHTTPRequestHandler):
try:
if path in ("/", "/index.html"):
self.send_bytes((HERE / "index.html").read_bytes(), "text/html; charset=utf-8")
elif path == "/assistant":
page = (HERE / "assistant.html").read_text().replace("__FRAME_KEY__", json.dumps(UI_KEY).replace("<", "\\u003c"))
self.send_bytes(page.encode(), "text/html; charset=utf-8")
elif path == "/api/agent/approval":
token = (parse_qs(url.query).get("confirmation") or [""])[0]
self.send_json(frame_agent.approvals.inspect(token))
elif path == "/api/host":
self.send_json({"os": "SteamOS", "fileManager": None, "computer": DEVICE, "mobile": True} if LOCAL else
{"os": frame_host.NAME, "fileManager": frame_host.FILE_MANAGER,
@@ -1459,6 +1604,12 @@ class Handler(BaseHTTPRequestHandler):
"shared": frame_catalog.compat_db.shared()})
elif path == "/api/android/catalog":
self.send_json({"apps": frame_catalog.catalog()})
elif path == "/api/macview":
self.send_json(macview_state(parse_qs(url.query)))
elif path == "/api/telemetry":
self.send_json(frame_telemetry.state())
elif path == "/api/computer/state":
self.send_json(json.loads(ssh("python3 -", stdin=(HERE / "frame_computer.py").read_text(), timeout=20)))
elif path == "/api/status":
self.send_json(status({}))
elif path == "/api/steam/owned":
@@ -1482,9 +1633,12 @@ class Handler(BaseHTTPRequestHandler):
self.send_json({"error": "not found"}, 404)
except Failure as e:
self.send_error_json(str(e), e.status, e.apk)
except ValueError as e:
self.send_json({"error": str(e)}, 400)
except frame_android.FrameError as e:
self.send_error_json(str(e), 502)
except Exception as e:
frame_telemetry.diagnostic(f"GET {path}", e)
self.send_json({"error": f"{type(e).__name__}: {e}"}, 500)
def do_POST(self):
@@ -1492,6 +1646,7 @@ class Handler(BaseHTTPRequestHandler):
return
path = urlparse(self.path).path
meant = self.headers.get("X-Frame-Device")
body = None
try:
if path == "/api/upload":
with working(meant):
@@ -1511,12 +1666,16 @@ class Handler(BaseHTTPRequestHandler):
result = handler(body)
self.send_json(result)
except Failure as e:
if e.status >= 500:
frame_telemetry.diagnostic(f"POST {path} {action_of(body)}", e)
self.send_error_json(str(e), e.status, e.apk)
except (ValueError, TypeError) as e:
self.send_json({"error": f"bad request: {e}"}, 400)
except frame_android.FrameError as e:
frame_telemetry.diagnostic(f"POST {path} {action_of(body)}", e)
self.send_error_json(str(e), 502)
except Exception as e:
frame_telemetry.diagnostic(f"POST {path} {action_of(body)}", e)
self.send_json({"error": f"{type(e).__name__}: {e}"}, 500)
def connection_events(self):
@@ -1640,13 +1799,18 @@ class Handler(BaseHTTPRequestHandler):
keep = True # stage_title owns tmp now, and removes it on failure
return stage_title(str(dest), temp_dir=str(tmp))
if mode == "apk":
# Checked here, before install(), to hand the page a blocker it can offer
# alternatives for; report these failures the way install() would have.
start = time.time()
try:
info = frame_android.apk_info(str(dest))
except frame_android.FrameError as e:
frame_android._after_install(None, None, e, start)
raise Failure(str(e), 400)
try:
frame_android.check_installable(info)
except frame_android.FrameError as e:
frame_android._after_install(info, None, e, start)
raise Failure(str(e), 400, {"package": info["package"], "version_code": info.get("version_code"), "blocker": str(e)})
ensure_master()
try:
@@ -1660,6 +1824,15 @@ class Handler(BaseHTTPRequestHandler):
shutil.rmtree(tmp, ignore_errors=True)
class LoopbackServer(ThreadingHTTPServer):
def server_bind(self):
# HTTPServer.server_bind resolves socket.getfqdn(host), a reverse-DNS
# lookup that can stall for seconds (verified on GitHub's macOS runners).
# Loopback needs no hostname.
socketserver.TCPServer.server_bind(self)
self.server_name, self.server_port = "127.0.0.1", self.server_address[1]
_ONE_SERVER = None
@@ -1684,8 +1857,9 @@ def main():
help="stop cleanly when stdin closes (the app closes it on quit; "
"Windows has no SIGTERM to catch)")
args = ap.parse_args()
httpd = ThreadingHTTPServer(("127.0.0.1", args.port), Handler)
httpd = LoopbackServer(("127.0.0.1", args.port), Handler)
sweep_tmp()
frame_telemetry.start()
global LINK, _ONE_SERVER
if not LOCAL:
_ONE_SERVER = one_server()
@@ -1712,6 +1886,7 @@ def main():
if not frame_host.WINDOWS:
signal.signal(signal.SIGTERM, signal.SIG_IGN)
webinstall_shutdown()
macview.shutdown() # close the headset's viewers before the agent goes
# The master was started with -N, so it stays up until told to exit.
if LINK:
LINK.stop()