Skip to content
File

Blob: archive/legacy-c/run-device.py

python110 lines
1#!/usr/bin/env python3
2"""Supply USB signaling for the archived C radio through the Worker API."""
3import argparse
4import asyncio
5import importlib.util
6import json
7from pathlib import Path
8import sys
9from urllib.error import HTTPError
10from urllib.request import Request, urlopen
11 
12ROOT = Path(__file__).resolve().parents[2]
13sys.path.insert(0, str(ROOT / "scripts"))
14from config import read_env
15spec = importlib.util.spec_from_file_location("probe", ROOT / "archive/sfu-bringup/probe-sfu.py")
16probe = importlib.util.module_from_spec(spec)
17spec.loader.exec_module(probe)
18OUT = ROOT / "artifacts/archive/legacy-c"
19probe.OUT = OUT
20env = read_env(ROOT / "worker/.dev.vars")
21 
22 
23class SignalingError(RuntimeError):
24 def __init__(self, status, message):
25 super().__init__(f"Signaling HTTP {status}: {message}")
26 self.status = status
27 
28 
29def request(base, path, data):
30 req = Request(base + path, data=json.dumps(data).encode(), headers={
31 "Content-Type": "application/json", "Authorization": "Bearer " + env["DEVICE_TOKEN"]})
32 try:
33 with urlopen(req, timeout=45) as response:
34 return json.load(response)
35 except HTTPError as error:
36 try: message = json.loads(error.read(2048)).get("error", "request failed")
37 except ValueError: message = "request failed"
38 raise SignalingError(error.code, message) from None
39 
40 
41async def main():
42 parser = argparse.ArgumentParser(description=__doc__)
43 parser.add_argument("--url", default="http://127.0.0.1:11880")
44 parser.add_argument("--no-reset", action="store_true", help="Use immediately after flashing a fresh board")
45 args = parser.parse_args()
46 device = probe.Device()
47 try:
48 await device.handshake()
49 if not args.no_reset:
50 device.send(cmd="restart_device")
51 await asyncio.sleep(2)
52 device.close()
53 device = probe.Device()
54 await device.handshake()
55 device.send(cmd="peer_init")
56 opened = await device.wait(lambda e: e.get("cmd") == "peer_init")
57 if opened["result"] != 0:
58 raise RuntimeError(f"S3 radio initialization failed ({opened['result']}). Check its music partition and restart the board.")
59 offer = await device.wait(lambda e: e.get("event") == "sdp")
60 started = await asyncio.to_thread(request, args.url, "/api/device/start", {
61 "sessionDescription": {"type": "offer", "sdp": offer["text"]}})
62 for name, sdp in (("s3-offer", offer["text"]), ("sfu-answer", started["sessionDescription"]["sdp"])):
63 path = OUT / (name + ".sdp")
64 path.parent.mkdir(parents=True, exist_ok=True, mode=0o700)
65 path.touch(mode=0o600, exist_ok=True)
66 path.chmod(0o600)
67 path.write_text(sdp)
68 generation = started["generation"]
69 device.send(cmd="sdp", text=started["sessionDescription"]["sdp"])
70 await device.wait(lambda e: e.get("event") == "peer_state" and e.get("state") == 9)
71 created = await asyncio.to_thread(request, args.url, "/api/device/channels", {"generation": generation})
72 ids = {}
73 for name in ("robot", "spectrum"):
74 channel = next(c for c in created["channels"] if c["dataChannelName"] == name)
75 ids[name] = channel["id"]
76 device.send(cmd="create_channel")
77 reply = await device.wait(lambda e: e.get("cmd") == "create_channel")
78 if reply["result"] != 0: raise RuntimeError(f"Local {name} creation failed")
79 device.send(cmd="start", robot_id=ids["robot"], spectrum_id=ids["spectrum"])
80 reply = await device.wait(lambda e: e.get("cmd") == "start")
81 if reply["result"] != 0: raise RuntimeError("The device could not start playback")
82 await asyncio.to_thread(request, args.url, "/api/device/ready", {"generation": generation})
83 print("S3 radio is streaming Opus, telemetry and spectrum through Cloudflare. Open " + args.url, flush=True)
84 print("USB carries setup and a service heartbeat only. Leave this helper running during local development.", flush=True)
85 failures = 0
86 while True:
87 for _ in range(50):
88 device.poll()
89 device.events.clear()
90 await asyncio.sleep(0.1)
91 try:
92 await asyncio.to_thread(request, args.url, "/api/device/heartbeat", {"generation": generation})
93 if failures: print("Signaling heartbeat recovered.", flush=True)
94 failures = 0
95 except (SignalingError, OSError) as error:
96 if isinstance(error, SignalingError) and error.status < 500:
97 raise
98 failures += 1
99 if failures >= 10: raise
100 print(f"Signaling heartbeat temporarily unavailable; retry {failures}/10. The board keeps streaming.", flush=True)
101 finally:
102 # esp_peer_close currently hangs. A subsequent run restarts this demo
103 # firmware, and the server closes the old SFU resources before setup.
104 device.close()
105 
106 
107if __name__ == "__main__":
108 try: asyncio.run(main())
109 except KeyboardInterrupt: print("Device helper stopped.")