File
Blob: archive/legacy-c/run-device.py
| 1 | #!/usr/bin/env python3 |
| 2 | """Supply USB signaling for the archived C radio through the Worker API.""" |
| 3 | import argparse |
| 4 | import asyncio |
| 5 | import importlib.util |
| 6 | import json |
| 7 | from pathlib import Path |
| 8 | import sys |
| 9 | from urllib.error import HTTPError |
| 10 | from urllib.request import Request, urlopen |
| 11 | |
| 12 | ROOT = Path(__file__).resolve().parents[2] |
| 13 | sys.path.insert(0, str(ROOT / "scripts")) |
| 14 | from config import read_env |
| 15 | spec = importlib.util.spec_from_file_location("probe", ROOT / "archive/sfu-bringup/probe-sfu.py") |
| 16 | probe = importlib.util.module_from_spec(spec) |
| 17 | spec.loader.exec_module(probe) |
| 18 | OUT = ROOT / "artifacts/archive/legacy-c" |
| 19 | probe.OUT = OUT |
| 20 | env = read_env(ROOT / "worker/.dev.vars") |
| 21 | |
| 22 | |
| 23 | class SignalingError(RuntimeError): |
| 24 | def __init__(self, status, message): |
| 25 | super().__init__(f"Signaling HTTP {status}: {message}") |
| 26 | self.status = status |
| 27 | |
| 28 | |
| 29 | def 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 | |
| 41 | async 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 | |
| 107 | if __name__ == "__main__": |
| 108 | try: asyncio.run(main()) |
| 109 | except KeyboardInterrupt: print("Device helper stopped.") |