File
Blob: archive/sfu-bringup/investigate-datachannels.py
| 1 | #!/usr/bin/env python3 |
| 2 | """Bounded hardware experiment for explicit local DataChannel creation.""" |
| 3 | import asyncio |
| 4 | import json |
| 5 | from pathlib import Path |
| 6 | import runpy |
| 7 | import time |
| 8 | |
| 9 | from aiortc import RTCPeerConnection, RTCConfiguration, RTCSessionDescription |
| 10 | |
| 11 | ROOT = Path(__file__).resolve().parents[2] |
| 12 | probe = runpy.run_path(str(Path(__file__).with_name("probe-sfu.py"))) |
| 13 | |
| 14 | |
| 15 | async def main(): |
| 16 | api = probe["SFU"]() |
| 17 | device = probe["Device"]() |
| 18 | host = RTCPeerConnection(RTCConfiguration(iceServers=[])) |
| 19 | resources = [] |
| 20 | result = {} |
| 21 | try: |
| 22 | device.send(cmd="ping") |
| 23 | await device.wait(lambda e: e.get("cmd") == "ping") |
| 24 | device.send(cmd="peer_init", initiator=True) |
| 25 | initialized = await device.wait(lambda e: e.get("cmd") == "peer_init") |
| 26 | if initialized["result"] != 0: |
| 27 | raise RuntimeError(f"Device peer initialization: {initialized['result']}") |
| 28 | offer = await device.wait(lambda e: e.get("event") == "sdp") |
| 29 | source = await api.call("/sessions/new", { |
| 30 | "sessionDescription": {"type": "offer", "sdp": offer["text"]}}) |
| 31 | source_id = source["sessionId"] |
| 32 | device.send(cmd="sdp", text=source["sessionDescription"]["sdp"]) |
| 33 | await device.wait(lambda e: e.get("event") == "peer_state" and e.get("state") == 9) |
| 34 | |
| 35 | sink_id, transport = await api.establish() |
| 36 | await host.setRemoteDescription(RTCSessionDescription(**transport["sessionDescription"])) |
| 37 | await host.setLocalDescription(await host.createAnswer()) |
| 38 | await api.call(f"/sessions/{sink_id}/renegotiate", { |
| 39 | "sessionDescription": {"type": "answer", "sdp": host.localDescription.sdp}}, method="PUT") |
| 40 | deadline = time.monotonic() + 20 |
| 41 | while host.connectionState != "connected" and time.monotonic() < deadline: |
| 42 | device.poll() |
| 43 | await asyncio.sleep(0.02) |
| 44 | if host.connectionState != "connected": |
| 45 | raise RuntimeError("Host did not connect") |
| 46 | |
| 47 | created = await api.call(f"/sessions/{source_id}/datachannels/new", { |
| 48 | "dataChannels": [{"location": "local", "dataChannelName": "robot", "ordered": True}]}) |
| 49 | local_id = created["dataChannels"][0]["id"] |
| 50 | resources.append((source_id, local_id)) |
| 51 | pulled = await api.call(f"/sessions/{sink_id}/datachannels/new", { |
| 52 | "dataChannels": [{"location": "remote", "sessionId": source_id, "dataChannelName": "robot", |
| 53 | "ordered": True, "waitForAck": True, "canReply": True}]}) |
| 54 | remote_id = pulled["dataChannels"][0]["id"] |
| 55 | resources.append((sink_id, remote_id)) |
| 56 | result.update(sfu_source_id=local_id, sfu_subscriber_id=remote_id) |
| 57 | print("SFU IDs:", local_id, remote_id, flush=True) |
| 58 | channel = host.createDataChannel("robot", negotiated=True, id=remote_id) |
| 59 | received = [] |
| 60 | channel.on("message", received.append) |
| 61 | deadline = time.monotonic() + 15 |
| 62 | while channel.readyState != "open": |
| 63 | if time.monotonic() >= deadline or channel.readyState == "closed" or host.connectionState in {"failed", "closed"}: |
| 64 | raise TimeoutError("Data channel did not open before the deadline") |
| 65 | device.poll() |
| 66 | await asyncio.sleep(0.02) |
| 67 | channel.send("ack") |
| 68 | |
| 69 | device.send(cmd="create_channel") |
| 70 | created = await device.wait(lambda e: e.get("cmd") == "create_channel") |
| 71 | result["create_result"] = created["result"] |
| 72 | deadline = time.monotonic() + 4 |
| 73 | while time.monotonic() < deadline: |
| 74 | device.poll() |
| 75 | await asyncio.sleep(0.02) |
| 76 | result["opened_channels"] = [e for e in device.events if e.get("event") == "channel_open"] |
| 77 | result["send_results"] = {} |
| 78 | for stream_id in sorted({0, local_id}): |
| 79 | device.send(cmd="send", stream_id=stream_id, |
| 80 | text=json.dumps({"probe": "explicit_creation", "stream_id": stream_id})) |
| 81 | sent = await device.wait(lambda e: e.get("cmd") == "send") |
| 82 | result["send_results"][stream_id] = sent["result"] |
| 83 | deadline = time.monotonic() + 3 |
| 84 | while time.monotonic() < deadline: |
| 85 | device.poll() |
| 86 | await asyncio.sleep(0.02) |
| 87 | result["received"] = [m.decode() if isinstance(m, bytes) else m for m in received] |
| 88 | except Exception as exc: |
| 89 | result["error"] = f"{type(exc).__name__}: {exc}" |
| 90 | finally: |
| 91 | for sid, cid in reversed(resources): |
| 92 | try: |
| 93 | await api.call(f"/sessions/{sid}/datachannels/close", {"dataChannels": [{"id": cid}]}, method="PUT") |
| 94 | except Exception: |
| 95 | result["cleanup_incomplete"] = True |
| 96 | await host.close() |
| 97 | device.close() |
| 98 | out = ROOT / "artifacts/archive/sfu-bringup/explicit-channel-result.json" |
| 99 | out.write_text(json.dumps(result, indent=2) + "\n") |
| 100 | print(json.dumps(result, indent=2), flush=True) |
| 101 | |
| 102 | |
| 103 | if __name__ == "__main__": |
| 104 | asyncio.run(main()) |