mirror of
https://github.com/saphid/frame-control.git
synced 2026-10-06 06:00:33 +02:00
Merge main into family-comfort; a session state without 'started' can't crash status
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
commit
1203ec0a9e
180 files changed
+62069
-234
No files matched your search
+589
-13
@@ -10,9 +10,11 @@ Env: FRAME_ALIAS (default frame)
|
||||
FRAME_LOCAL=1 run on the Frame itself (the iPhone app starts it there over SSH)
|
||||
FRAME_UI_KEY required X-Frame-UI value (the iPhone app passes a fresh one)
|
||||
FRAME_DEVICE what to call the device the page runs on (e.g. iPhone)
|
||||
FRAME_CLIENT a stable id for that device (keyboard-and-trackpad pairing is kept per id)
|
||||
"""
|
||||
import argparse
|
||||
import base64
|
||||
import hashlib
|
||||
import http.client
|
||||
import json
|
||||
import os
|
||||
@@ -23,6 +25,7 @@ import shlex
|
||||
import shutil
|
||||
import signal
|
||||
import socket
|
||||
import socketserver
|
||||
import subprocess
|
||||
import sys
|
||||
import tempfile
|
||||
@@ -36,12 +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_comfort # noqa: E402
|
||||
import frame_host # noqa: E402
|
||||
import frame_macview # noqa: E402
|
||||
import frame_media # 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
|
||||
|
||||
@@ -56,12 +65,14 @@ if LOCAL:
|
||||
os.environ["PATH"] = f"{HERE / 'local-bin'}{os.pathsep}{os.environ.get('PATH', '')}"
|
||||
UI_KEY = os.environ.get("FRAME_UI_KEY") or "1"
|
||||
DEVICE = os.environ.get("FRAME_DEVICE") or "phone"
|
||||
# What the Frame's KDE Connect calls this device (keyboard and trackpad).
|
||||
INPUT_NAME = DEVICE if LOCAL else socket.gethostname().split(".")[0]
|
||||
FRAME = os.environ.get("FRAME_ALIAS", "frame")
|
||||
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")
|
||||
MUX = ["ssh", "-o", "BatchMode=yes", *(["-o", f"ControlPath={CONTROL}"] if CONTROL else [])]
|
||||
# Commands use the master when it's up and connect directly when it isn't.
|
||||
SSH = [*MUX, *(["-o", "ControlMaster=no"] if CONTROL else []), "-o", "ConnectTimeout=5"]
|
||||
@@ -163,8 +174,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())
|
||||
@@ -250,7 +263,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 comfort(body):
|
||||
@@ -463,7 +481,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):
|
||||
@@ -531,16 +558,338 @@ def clipboard(body):
|
||||
return {"message": ssh(PASTE_CMD, stdin=text, timeout=30).strip()}
|
||||
|
||||
|
||||
# ---- keyboard and pointer (KDE Connect on the Frame, see frame_input_agent.py) ----
|
||||
|
||||
INPUT_FLAGS = ("singleclick", "doubleclick", "middleclick", "rightclick", "singlehold", "singlerelease",
|
||||
"scroll", "ctrl", "alt", "shift", "super")
|
||||
INPUT_MOVE_LIMIT = 2000 # pixels per event
|
||||
INPUT_TEXT_LIMIT = 500 # characters per event
|
||||
INPUT_BATCH_LIMIT = 200 # events per request
|
||||
|
||||
|
||||
def input_event(event):
|
||||
"""A KDE Connect remote-input body with only the fields it knows, in range."""
|
||||
if not isinstance(event, dict):
|
||||
raise Failure("each input event must be an object", 400)
|
||||
out = {}
|
||||
for name in ("dx", "dy"):
|
||||
value = event.get(name)
|
||||
if value is None:
|
||||
continue
|
||||
if isinstance(value, bool) or not isinstance(value, (int, float)) or value != value:
|
||||
raise Failure(f"{name} must be a number", 400)
|
||||
out[name] = max(-INPUT_MOVE_LIMIT, min(INPUT_MOVE_LIMIT, round(float(value), 2)))
|
||||
for name in INPUT_FLAGS:
|
||||
if event.get(name) is True:
|
||||
out[name] = True
|
||||
key = event.get("key")
|
||||
if key is not None:
|
||||
if not isinstance(key, str) or not 0 < len(key) <= INPUT_TEXT_LIMIT:
|
||||
raise Failure(f"key must be text of 1 to {INPUT_TEXT_LIMIT} characters", 400)
|
||||
out["key"] = key
|
||||
special = event.get("specialKey")
|
||||
if special is not None:
|
||||
# KDE Connect's numbering: 1 Backspace … 14 Escape, 21-32 F1-F12.
|
||||
if isinstance(special, bool) or not isinstance(special, int) or not 1 <= special <= 32:
|
||||
raise Failure("specialKey must be a whole number from 1 to 32", 400)
|
||||
out["specialKey"] = special
|
||||
if not set(out) - {"ctrl", "alt", "shift", "super"}:
|
||||
raise Failure("input event has nothing to do", 400)
|
||||
return out
|
||||
|
||||
|
||||
def input_client():
|
||||
"""A stable id for this device, so KDE Connect on the Frame keeps its pairing apart.
|
||||
|
||||
The iPhone app passes one (FRAME_CLIENT). A computer makes one the first time
|
||||
and keeps it: host names alone can clash (desk.home and desk.office).
|
||||
"""
|
||||
if os.environ.get("FRAME_CLIENT"):
|
||||
return os.environ["FRAME_CLIENT"]
|
||||
if LOCAL:
|
||||
return DEVICE
|
||||
path = frame_host.data_dir("input-client-id")
|
||||
try:
|
||||
saved = path.read_text().strip()
|
||||
if re.fullmatch(r"[A-Za-z0-9_-]{4,64}", saved):
|
||||
return saved
|
||||
except OSError:
|
||||
pass
|
||||
made = (re.sub(r"[^A-Za-z0-9-]", "", INPUT_NAME)[:24] or "computer") + "-" + secrets.token_hex(4)
|
||||
try:
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
path.write_text(made)
|
||||
except OSError:
|
||||
pass # still works this time; it pairs again next time
|
||||
return made
|
||||
|
||||
|
||||
# KDE Connect for the Frame, as Frame Control ships it (frame/kdeconnect/NOTICE.md).
|
||||
KDECONNECT = HERE.parent / "frame" / "kdeconnect"
|
||||
KDECONNECT_HOME = ".local/share/frame-control/kdeconnect"
|
||||
|
||||
|
||||
def kdeconnect_packages():
|
||||
"""[(file, sha256), ...] from frame/kdeconnect/packages.json; [] if it's missing."""
|
||||
try:
|
||||
manifest = json.loads((KDECONNECT / "packages.json").read_text())
|
||||
return [(p["file"], p["sha256"]) for p in manifest["packages"]]
|
||||
except (OSError, ValueError, KeyError, TypeError):
|
||||
return []
|
||||
|
||||
|
||||
def kdeconnect_stamp(packages):
|
||||
"""What the agent writes once these are unpacked (frame_input_agent.stamp)."""
|
||||
return "".join(f"{sha} {name}\n" for name, sha in packages)
|
||||
|
||||
|
||||
def file_sha256(path):
|
||||
digest = hashlib.sha256()
|
||||
with open(path, "rb") as f:
|
||||
for block in iter(lambda: f.read(1 << 20), b""):
|
||||
digest.update(block)
|
||||
return digest.hexdigest()
|
||||
|
||||
|
||||
class InputAgent:
|
||||
"""frame_input_agent.py running on the Frame, fed events over one long-lived ssh.
|
||||
|
||||
Before starting it, copies KDE Connect to the Frame if it isn't there yet. The
|
||||
agent unpacks it, pairs, and reports its state ({"state": "off" | "installing" |
|
||||
"starting" | "pairing" | "ready" | "error"}).
|
||||
"""
|
||||
|
||||
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
|
||||
# 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()
|
||||
client = input_client()
|
||||
return ("python3 -u -c " + shlex.quote(
|
||||
f"import base64;exec(compile(base64.b64decode('{code}'),'frame_input_agent','exec'))")
|
||||
+ f" {shlex.quote(client)} {shlex.quote(INPUT_NAME)} {shlex.quote(folder)}"
|
||||
+ f" {shlex.quote(json.dumps(self.packages, separators=(',', ':')))}")
|
||||
|
||||
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, unless the Frame already has them unpacked (`force`:
|
||||
the agent found it didn't after all).
|
||||
"""
|
||||
if LOCAL:
|
||||
return str(KDECONNECT / "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 ""
|
||||
# A folder of its own: a cancelled start's agent may still be cleaning up another.
|
||||
folder = f"{KDECONNECT_HOME}/incoming/{secrets.token_hex(8)}"
|
||||
ssh(f"mkdir -p {folder}", timeout=20)
|
||||
try:
|
||||
for name, sha in self.packages:
|
||||
path = KDECONNECT / "packages" / name
|
||||
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.rsplit('-', 3)[0]})")
|
||||
quoted = shlex.quote(name)
|
||||
ssh(f"cd {folder} && cat > {quoted}.part && mv {quoted}.part {quoted}",
|
||||
stdin=path.read_bytes(), text=False, timeout=600)
|
||||
except (Failure, OSError):
|
||||
self.discard(folder) # a partial copy is no use to anyone
|
||||
raise
|
||||
return f"~/{folder}"
|
||||
|
||||
def discard(self, folder):
|
||||
"""Remove a copy no agent will take over (best effort; agents tidy up old ones too)."""
|
||||
folder = folder.removeprefix("~/")
|
||||
if folder.startswith(f"{KDECONNECT_HOME}/incoming/") and not LOCAL:
|
||||
try:
|
||||
ssh(f"rm -rf {folder}", timeout=20)
|
||||
except (Failure, OSError):
|
||||
pass # never let tidying up get in the way of reporting and retrying
|
||||
|
||||
def start(self):
|
||||
with self.lock:
|
||||
if self.launching == self.generation or (self.proc and self.proc.poll() is None):
|
||||
return
|
||||
# (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, force=False):
|
||||
try:
|
||||
self._launch_once(generation, force)
|
||||
except Exception as e: # whatever went wrong, never leave it stuck "starting"
|
||||
with self.lock:
|
||||
if self.launching == generation:
|
||||
self.launching = None
|
||||
if self.generation == generation and self.status.get("state") in ("starting", "installing"):
|
||||
self.status = {"state": "error", "message": f"Couldn't start the keyboard and trackpad: {e}"}
|
||||
|
||||
def _launch_once(self, generation, force):
|
||||
def report(message):
|
||||
with self.lock:
|
||||
if self.generation == generation:
|
||||
self.status = {"state": "installing", "message": message}
|
||||
folder = ""
|
||||
try:
|
||||
errors = tempfile.TemporaryFile()
|
||||
ensure_master()
|
||||
folder = self.deliver(report, force)
|
||||
with self.lock:
|
||||
stopped = self.generation != generation
|
||||
if stopped: # turned off while copying
|
||||
raise Failure("stopped")
|
||||
proc = subprocess.Popen([*SSH, FRAME, self.command(folder)], stdin=subprocess.PIPE,
|
||||
stdout=subprocess.PIPE, stderr=errors)
|
||||
except (Failure, OSError) as e:
|
||||
self.discard(folder) # no agent will take the copy over
|
||||
message = str(e)
|
||||
friendly = unreachable(message)
|
||||
with self.lock:
|
||||
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:
|
||||
if self.launching == generation:
|
||||
self.launching = None
|
||||
stale = self.generation != generation
|
||||
if not stale:
|
||||
self.proc = proc
|
||||
if stale: # turned off meanwhile
|
||||
proc.terminate()
|
||||
wanted, heard = self._watch(proc, errors, retry=not force)
|
||||
if not heard:
|
||||
# The agent never started (or was stopped first), so it can't tidy the copy up.
|
||||
self.discard(folder)
|
||||
if wanted 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:
|
||||
# Unless a start() already took over (it launches, and copies if still needed).
|
||||
if self.proc is not proc or self.generation != generation or self.launching is not None:
|
||||
return
|
||||
self.proc, self.launching, self.status = None, generation, {"state": "starting"}
|
||||
self._launch(generation, force=True)
|
||||
|
||||
def _watch(self, proc, errors, retry=False):
|
||||
"""Follow the agent's status until it exits. Returns whether it asked for the
|
||||
packages (`retry`: the caller will send them, so that isn't an error yet), and
|
||||
whether it said anything at all (then it holds its copy and tidies it up)."""
|
||||
wanted = heard = False
|
||||
for line in proc.stdout:
|
||||
try:
|
||||
status = json.loads(line)
|
||||
except ValueError:
|
||||
continue
|
||||
if isinstance(status, dict) and isinstance(status.get("state"), str):
|
||||
heard = True
|
||||
if status["state"] == "need-packages":
|
||||
wanted = True
|
||||
continue
|
||||
with self.lock:
|
||||
if self.proc is proc:
|
||||
self.status = status
|
||||
proc.wait()
|
||||
_live_tunnels.discard(proc)
|
||||
try:
|
||||
errors.seek(0)
|
||||
detail = strip_ansi(errors.read().decode(errors="replace")).strip()
|
||||
except OSError:
|
||||
detail = ""
|
||||
with self.lock:
|
||||
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, heard
|
||||
|
||||
def send(self, events):
|
||||
"""Forward events if the agent is ready; start it if it isn't running.
|
||||
|
||||
Returns the state, with "sent" saying whether the events went; if not,
|
||||
the page keeps them and sends them again once the state is "ready".
|
||||
"""
|
||||
with self.lock:
|
||||
proc, ready = self.proc, self.status.get("state") == "ready"
|
||||
sent = False
|
||||
if not (proc and proc.poll() is None):
|
||||
self.start()
|
||||
elif ready and events:
|
||||
try:
|
||||
proc.stdin.write((json.dumps(events) + "\n").encode())
|
||||
proc.stdin.flush()
|
||||
sent = True
|
||||
except (BrokenPipeError, OSError, ValueError):
|
||||
pass # _watch reports how it ended
|
||||
with self.lock:
|
||||
return {**self.status, "sent": sent}
|
||||
|
||||
def stop(self):
|
||||
with self.lock:
|
||||
proc, self.proc, self.status = self.proc, None, {"state": "off"}
|
||||
self.generation += 1
|
||||
if proc and proc.poll() is None:
|
||||
proc.terminate()
|
||||
|
||||
|
||||
_input = InputAgent()
|
||||
|
||||
|
||||
def licenses():
|
||||
"""The notices and licence texts for what Frame Control ships (the About dialog)."""
|
||||
found = [("Third-party notices", HERE.parent / "THIRD_PARTY_NOTICES.md"), ("KDE Connect for the Frame", KDECONNECT / "NOTICE.md"),
|
||||
("Frame Control (MIT)", HERE.parent / "LICENSE")]
|
||||
found += [(f"{p.parent.name}: {p.stem}", p) for p in sorted((KDECONNECT / "LICENSES").glob("*/*.txt"))]
|
||||
notices = []
|
||||
for title, path in found:
|
||||
try:
|
||||
notices.append({"title": title, "text": path.read_text(errors="replace")})
|
||||
except OSError:
|
||||
pass
|
||||
return notices
|
||||
|
||||
|
||||
def remote_input(body):
|
||||
"""{"events": [...]} sends keyboard and pointer events; {} (or none yet) just starts the agent."""
|
||||
events = body.get("events", [])
|
||||
if not isinstance(events, list) or len(events) > INPUT_BATCH_LIMIT:
|
||||
raise Failure(f"events must be a list of at most {INPUT_BATCH_LIMIT}", 400)
|
||||
return _input.send([input_event(e) for e in events])
|
||||
|
||||
|
||||
def flatpak(body):
|
||||
app, action = str(body.get("id", "")), body.get("action")
|
||||
if not FLATPAK_ID.match(app):
|
||||
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":
|
||||
@@ -637,13 +986,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
|
||||
@@ -697,14 +1070,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())
|
||||
@@ -1157,10 +1534,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)
|
||||
@@ -1259,14 +1642,159 @@ def _sweep_one(prefix, d):
|
||||
pass
|
||||
|
||||
|
||||
POST = {"/api/comfort": comfort, "/api/android/display": android_display, "/api/android": android, "/api/titles": titles, "/api/launch": launch, "/api/steam": steam, "/api/volume": set_volume, "/api/clipboard": clipboard,
|
||||
# ---- Our Frame-side media player -----------------------------------------
|
||||
|
||||
_MEDIA_LOCK = threading.Lock()
|
||||
|
||||
|
||||
def media(body):
|
||||
action = body.get("action")
|
||||
if action not in ("list", "status", "play", "stop"):
|
||||
raise Failure("Media action must be list, status, play or stop", 400)
|
||||
if action == "play":
|
||||
identity = body.get("id")
|
||||
if not isinstance(identity, str) or not re.fullmatch(r"[0-9a-f]{32}/[^/\\\x00]+", identity):
|
||||
raise Failure("Invalid media id", 400)
|
||||
if body.get("layout", "auto") not in frame_media.LAYOUTS:
|
||||
raise Failure("Invalid media layout", 400)
|
||||
if type(body.get("theatre", False)) is not bool:
|
||||
raise Failure("theatre must be true or false", 400)
|
||||
with _MEDIA_LOCK:
|
||||
# Ship only our small stdlib modules, atomically, to the user account.
|
||||
sources = {name: (HERE / name).read_text() for name in (
|
||||
"frame_media.py", "frame_media_player.py", "frame_media_remote.py", "frame_splat.py")}
|
||||
installer = """import json, os, pathlib, sys, tempfile
|
||||
root = pathlib.Path.home()/'.local/share/frame-control/media'
|
||||
root.mkdir(parents=True, exist_ok=True)
|
||||
for name, source in json.load(sys.stdin).items():
|
||||
path = root/name
|
||||
fd, temp = tempfile.mkstemp(dir=root, prefix=name+'.')
|
||||
with os.fdopen(fd, 'w') as f:
|
||||
f.write(source)
|
||||
os.replace(temp, path)
|
||||
"""
|
||||
ssh("python3 -c " + shlex.quote(installer), stdin=json.dumps(sources))
|
||||
try:
|
||||
out = ssh("python3 ~/.local/share/frame-control/media/frame_media_remote.py",
|
||||
stdin=json.dumps(body), timeout=75) # remote worst case: ffprobe 30 + reset-failed 10 + systemd-run 15 s
|
||||
except Failure as e:
|
||||
for line in reversed(getattr(e, "stdout", "").splitlines()):
|
||||
try:
|
||||
detail = json.loads(line).get("error")
|
||||
except (ValueError, AttributeError):
|
||||
continue
|
||||
if detail:
|
||||
raise Failure(detail) from None
|
||||
raise
|
||||
return json.loads(out)
|
||||
|
||||
|
||||
def push_media(path):
|
||||
# Validate the format, but leave layout selection until playback (ffprobe
|
||||
# can then read metadata on the Frame, where it is installed).
|
||||
name = Path(path).name
|
||||
frame_media.plan(name, "mono")
|
||||
if name.startswith(".") or "\\" in name:
|
||||
raise Failure("Rename the file: media names can't start with a dot or contain a backslash", 400)
|
||||
token = secrets.token_hex(16)
|
||||
dest = "Videos/FrameControl/" + token + "/"
|
||||
ssh("mkdir -p ~/" + dest)
|
||||
try:
|
||||
push_file(path, dest)
|
||||
except Exception:
|
||||
try:
|
||||
ssh("rm -rf ~/" + dest)
|
||||
except Exception:
|
||||
pass # keep the copy error; an empty folder isn't listed as media
|
||||
raise
|
||||
return {"message": "Media sent. Choose its layout and press Play.",
|
||||
"id": token + "/" + name}
|
||||
|
||||
|
||||
# ---- 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/comfort": comfort, "/api/media": media, "/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/input": remote_input,
|
||||
"/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/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}
|
||||
|
||||
|
||||
# ---- 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
|
||||
@@ -1363,6 +1891,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,
|
||||
@@ -1377,6 +1911,10 @@ class Handler(BaseHTTPRequestHandler):
|
||||
self.send_json({"titles": frame_titles.list_titles()})
|
||||
elif path == "/api/titles/job":
|
||||
self.send_json(title_job(url.query))
|
||||
elif path == "/api/licenses":
|
||||
self.send_json({"notices": licenses()})
|
||||
elif path == "/api/input":
|
||||
self.send_json(_input.send([]) if parse_qs(url.query).get("start") == ["1"] else dict(_input.status))
|
||||
elif path == "/api/job":
|
||||
self.send_json(job_status(url.query))
|
||||
elif path == "/api/android/displays":
|
||||
@@ -1386,6 +1924,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":
|
||||
@@ -1409,15 +1953,19 @@ 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):
|
||||
if not self.local_request():
|
||||
return
|
||||
path = urlparse(self.path).path
|
||||
body = None
|
||||
try:
|
||||
if path == "/api/upload":
|
||||
self.send_json(self.upload())
|
||||
@@ -1434,12 +1982,16 @@ class Handler(BaseHTTPRequestHandler):
|
||||
raise Failure("request body must be a JSON object", 400)
|
||||
self.send_json(handler(body))
|
||||
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 stream_video(self, query):
|
||||
@@ -1522,6 +2074,8 @@ class Handler(BaseHTTPRequestHandler):
|
||||
raise Failure("upload interrupted", 400)
|
||||
f.write(chunk)
|
||||
remaining -= len(chunk)
|
||||
if mode == "media":
|
||||
return push_media(dest)
|
||||
if mode == "apkinfo":
|
||||
# Read an APK for a report without installing it.
|
||||
try:
|
||||
@@ -1539,26 +2093,46 @@ 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:
|
||||
m = frame_android.install(str(dest), source=name)
|
||||
display = self.headers.get("X-APK-Display", "auto")
|
||||
if display not in ("auto", "flat", "vr"):
|
||||
raise frame_android.FrameError("invalid APK display mode")
|
||||
m = frame_android.install(str(dest), source=name,
|
||||
flatscreen=None if display == "auto" else display == "flat")
|
||||
except frame_android.FrameError as e:
|
||||
raise Failure(str(e), 400)
|
||||
return {"message": f"Installed {m['label']} as its own app in the Steam library", "app": m}
|
||||
kind = "VR app" if not m['flatscreen'] else "app"
|
||||
notes = " ".join(m.get("vr_issues", []))
|
||||
return {"message": f"Installed {m['label']} as its own {kind} in the Steam library. {notes}".strip(), "app": m}
|
||||
return {"message": push_file(dest)}
|
||||
finally:
|
||||
if not keep:
|
||||
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]
|
||||
|
||||
|
||||
def main():
|
||||
ap = argparse.ArgumentParser(description=__doc__.splitlines()[0])
|
||||
ap.add_argument("--port", type=int, default=int(os.environ.get("PORT", 47810)))
|
||||
@@ -1566,8 +2140,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()
|
||||
if not frame_host.WINDOWS:
|
||||
signal.signal(signal.SIGTERM, lambda *_: (_ for _ in ()).throw(KeyboardInterrupt))
|
||||
if args.exit_on_eof:
|
||||
@@ -1587,6 +2162,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 CONTROL:
|
||||
subprocess.run([*MUX, "-O", "exit", FRAME], capture_output=True, stdin=subprocess.DEVNULL)
|
||||
|
||||
Reference in new issue
Block a user