From 4b1a0aecda584a06a918b46d48c55df7216c967b Mon Sep 17 00:00:00 2001 From: Paul Nechifor Date: Sat, 18 Jul 2026 15:34:14 +0300 Subject: [PATCH] feat(web): relay process management, e2e test, wheel packaging --- .github/workflows/release-build-check.yml | 2 + MANIFEST.in | 8 + dimos/web/relay_bridge/demo_smoke.py | 180 +++++++++++++ dimos/web/relay_bridge/relay_process.py | 154 +++++++++++ dimos/web/relay_bridge/test_relay_e2e.py | 311 ++++++++++++++++++++++ dimos/web/relay_bridge/test_wt_client.py | 2 +- dimos/web/relay_bridge/wt_client.py | 10 +- setup.py | 25 ++ 8 files changed, 686 insertions(+), 6 deletions(-) create mode 100644 dimos/web/relay_bridge/demo_smoke.py create mode 100644 dimos/web/relay_bridge/relay_process.py create mode 100644 dimos/web/relay_bridge/test_relay_e2e.py diff --git a/.github/workflows/release-build-check.yml b/.github/workflows/release-build-check.yml index 4ecf3c5ee0..4ab1c7be51 100644 --- a/.github/workflows/release-build-check.yml +++ b/.github/workflows/release-build-check.yml @@ -9,6 +9,8 @@ on: - .github/workflows/release.yml - .github/workflows/release-build-check.yml - pyproject.toml + - setup.py + - MANIFEST.in permissions: {} diff --git a/MANIFEST.in b/MANIFEST.in index 1536332725..069098507e 100644 --- a/MANIFEST.in +++ b/MANIFEST.in @@ -24,6 +24,14 @@ recursive-exclude dimos/web/websocket_vis/node_modules * recursive-exclude dimos/web/dimos_interface * recursive-include dimos/web/dimos_interface/api *.py *.html +# Deno relay sources (repo-root web/): shipped inside the wheel by the +# build_py hook in setup.py (copied to dimos/web/relay_bridge/_relay_dist). +# The graft keeps them in the sdist so sdist->wheel builds (uv build, pip) +# can reproduce that copy. +graft web +prune web/cockpit/node_modules +prune web/cockpit/dist + # Exclude development files exclude .gitignore exclude .gitattributes diff --git a/dimos/web/relay_bridge/demo_smoke.py b/dimos/web/relay_bridge/demo_smoke.py new file mode 100644 index 0000000000..acc8298a95 --- /dev/null +++ b/dimos/web/relay_bridge/demo_smoke.py @@ -0,0 +1,180 @@ +# Copyright 2026 Dimensional Inc. +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""Manual smoke demo for the relay chain. + +Spawns the Deno relay (unless --url points at a running one), then drives a +robot client pushing synthetic color_image JPEGs as fast as they encode +(latest-wins) plus odom at 20 Hz (reliable), and a viewer client receiving +both. Open the printed debug URL in Chrome/Firefox to watch the same stream. + +Run: uv run python -m dimos.web.relay_bridge.demo_smoke [--secs 20] [--url https://...] +""" + +from __future__ import annotations + +import argparse +import asyncio +import contextlib +import json +import math +import time +from typing import Any +from urllib.parse import urlparse + +import cv2 +import numpy as np + +from dimos.web.relay_bridge.protocol import DataFrame +from dimos.web.relay_bridge.relay_process import RelayProcess +from dimos.web.relay_bridge.wt_client import RelayClient + +WIDTH, HEIGHT = 640, 480 + + +def make_jpeg(seq: int) -> bytes: + """Synthetic camera frame: moving gradient + seq/timestamp overlay.""" + ramp = np.linspace(0, 255, WIDTH, dtype=np.uint8) + gray = np.roll(np.tile(ramp, (HEIGHT, 1)), (seq * 7) % WIDTH, axis=1) + image = cv2.cvtColor(gray, cv2.COLOR_GRAY2BGR) + cv2.putText( + image, + f"seq {seq} {time.strftime('%H:%M:%S')}", + (20, 60), + cv2.FONT_HERSHEY_SIMPLEX, + 1.2, + (0, 255, 0), + 2, + ) + ok, encoded = cv2.imencode(".jpg", image, [int(cv2.IMWRITE_JPEG_QUALITY), 75]) + assert ok + return encoded.tobytes() + + +class ViewerStats: + def __init__(self) -> None: + self.channels: dict[str, dict[str, Any]] = {} + + def on_frame(self, frame: DataFrame) -> None: + ch = self.channels.setdefault( + frame.header.ch, + {"frames": 0, "bytes": 0, "seqs": set(), "last": -1, "ooo": 0, "lat_ms": 0.0}, + ) + ch["frames"] += 1 + ch["bytes"] += len(frame.payload) + ch["seqs"].add(frame.header.seq) + if frame.header.seq < ch["last"]: + ch["ooo"] += 1 + ch["last"] = max(ch["last"], frame.header.seq) + ch["lat_ms"] = (time.time() - frame.header.ts) * 1000 + + def line(self) -> str: + parts = [] + for name, ch in sorted(self.channels.items()): + seqs = ch["seqs"] + span_loss = (max(seqs) - min(seqs) + 1 - len(seqs)) if seqs else 0 + parts.append( + f"{name}: {ch['frames']}f span_loss={span_loss} ooo={ch['ooo']} " + f"lat={ch['lat_ms']:.1f}ms" + ) + return " | ".join(parts) or "(nothing received yet)" + + +async def run(url: str, secs: float) -> None: + stats = ViewerStats() + async with ( + await RelayClient.connect(url, "robot") as robot, + await RelayClient.connect(url, "viewer") as viewer, + ): + await robot.hello() + await viewer.hello() + rtt = await viewer.ping() + print(f"connected; datagram RTT {rtt * 1000:.1f} ms") + + deadline = time.monotonic() + secs if secs > 0 else math.inf + image_writer = robot.latest_writer("color_image") + odom_sent = 0 + + async def image_pump() -> None: + seq = 0 + while time.monotonic() < deadline: + image_writer.offer(make_jpeg(seq), meta={"w": WIDTH, "h": HEIGHT}) + seq += 1 + await asyncio.sleep(0) # flat out: paced by encode + delivery + + async def odom_pump() -> None: + nonlocal odom_sent + while time.monotonic() < deadline: + t = time.time() + payload = json.dumps( + {"x": 3 * math.sin(t / 3), "y": 2 * math.sin(t / 2), "yaw": t % 6.28, "ts": t} + ).encode() + robot.send_frame("odom", payload, delivery="reliable") + odom_sent += 1 + await asyncio.sleep(1 / 20) + + async def viewer_pump() -> None: + async for frame in viewer.frames(): + stats.on_frame(frame) + + async def report() -> None: + while time.monotonic() < deadline: + await asyncio.sleep(2) + print( + f"tx img {image_writer.sent} (dropped {image_writer.dropped}, " + f"resets {image_writer.resets}) odom {odom_sent} | rx {stats.line()}" + ) + + viewer_task = asyncio.create_task(viewer_pump()) + try: + await asyncio.gather(image_pump(), odom_pump(), report()) + finally: + await asyncio.sleep(0.3) # let the tail of the stream arrive + viewer_task.cancel() + + print("\nsummary:") + print( + f" sent: color_image {image_writer.sent} (+{image_writer.dropped} shed " + f"at source), odom {odom_sent}" + ) + print(f" received: {stats.line()}") + odom = stats.channels.get("odom", {}) + odom_ok = odom and len(odom["seqs"]) == odom_sent and odom["last"] == odom_sent - 1 + img_frames = stats.channels.get("color_image", {}).get("frames", 0) + print( + f" verdict: odom {'complete' if odom_ok else 'INCOMPLETE'}, " + f"color_image {img_frames} frames delivered" + ) + + +def main() -> None: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--url", default=None, help="attach to a running relay (wtUrl)") + parser.add_argument("--secs", type=float, default=0, help="run time; 0 = until Ctrl-C") + args = parser.parse_args() + + if args.url is not None: + # /api/info's wtUrl ends in /viewer; strip the path so each role picks its own. + parsed = urlparse(args.url) + url = f"{parsed.scheme}://{parsed.netloc}" if parsed.netloc else args.url + asyncio.run(run(url, args.secs)) + return + with RelayProcess() as info: + print(f"relay up; open {info.debug_url} in Chrome/Firefox to watch") + with contextlib.suppress(KeyboardInterrupt): + asyncio.run(run(info.wt_url, args.secs)) + + +if __name__ == "__main__": + main() diff --git a/dimos/web/relay_bridge/relay_process.py b/dimos/web/relay_bridge/relay_process.py new file mode 100644 index 0000000000..4d5827ad29 --- /dev/null +++ b/dimos/web/relay_bridge/relay_process.py @@ -0,0 +1,154 @@ +# Copyright 2026 Dimensional Inc. +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""Spawn the Deno relay as a child process (used by tests and the smoke demo; +T2 grows this into the RelayBridge module's managed relay process).""" + +from __future__ import annotations + +import collections +from dataclasses import dataclass +import json +import os +from pathlib import Path +import queue +import subprocess +import threading +from typing import IO + +from dimos.utils.deno import ensure_deno +from dimos.utils.logging_config import setup_logger +from dimos.web.relay_bridge.locate import find_web_dir, relay_run_cmd + +logger = setup_logger() + +_STDERR_TAIL_LINES = 60 + + +@dataclass +class RelayReadyInfo: + http_port: int + wt_url: str + cert_hash: str + v: int + + @property + def debug_url(self) -> str: + return f"http://127.0.0.1:{self.http_port}/debug.html" + + +class RelayProcess: + """Relay child process with a parsed ready line and clean teardown.""" + + def __init__( + self, + *, + port: int = 0, + host: str = "127.0.0.1", + web_dir: Path | None = None, + timeout: float = 20.0, + ) -> None: + self._port = port + self._host = host + self._web_dir = web_dir + self._timeout = timeout + self._process: subprocess.Popen[str] | None = None + self._threads: list[threading.Thread] = [] + self._ready_queue: queue.Queue[RelayReadyInfo] = queue.Queue(maxsize=1) + self._stderr_tail: collections.deque[str] = collections.deque(maxlen=_STDERR_TAIL_LINES) + self.info: RelayReadyInfo | None = None + + def start(self) -> RelayReadyInfo: + deno = ensure_deno() + web_dir = self._web_dir or find_web_dir() + cmd = relay_run_cmd(deno, web_dir, "--port", str(self._port), "--host", self._host) + logger.info(f"starting relay: {' '.join(cmd)}") + env = os.environ | {"NO_COLOR": "1"} + self._process = subprocess.Popen( + cmd, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, env=env + ) + assert self._process.stdout is not None and self._process.stderr is not None + self._threads = [ + threading.Thread(target=self._read_stdout, args=(self._process.stdout,), daemon=True), + threading.Thread(target=self._read_stderr, args=(self._process.stderr,), daemon=True), + ] + for thread in self._threads: + thread.start() + try: + self.info = self._ready_queue.get(timeout=self._timeout) + except queue.Empty: + code = self._process.poll() + self.stop() + stderr = "\n".join(self._stderr_tail) + state = f"exited with {code}" if code is not None else "still running" + raise RuntimeError( + f"relay produced no ready line within {self._timeout} s ({state}); " + f"stderr tail:\n{stderr}" + ) from None + logger.info(f"relay ready: {self.info}") + return self.info + + def stop(self) -> None: + if self._process is None: + return + process, self._process = self._process, None + process.terminate() + try: + process.wait(timeout=5) + except subprocess.TimeoutExpired: + process.kill() + process.wait(timeout=2) + # The child is dead: its pipes are at EOF, so the reader threads have + # finished. Join them and close the pipes so no file object leaks. + for thread in self._threads: + thread.join(timeout=1) + self._threads.clear() + for stream in (process.stdout, process.stderr): + if stream is not None: + stream.close() + + def __enter__(self) -> RelayReadyInfo: + return self.start() + + def __exit__(self, *exc_info: object) -> None: + self.stop() + + def _read_stdout(self, stream: IO[str]) -> None: + for line in stream: + line = line.rstrip() + if not line: + continue + if self.info is None and line.startswith("{"): + try: + data = json.loads(line) + except ValueError: + data = None + if isinstance(data, dict) and data.get("event") == "ready": + self._ready_queue.put( + RelayReadyInfo( + http_port=int(data["httpPort"]), + wt_url=str(data["wtUrl"]), + cert_hash=str(data["certHash"]), + v=int(data["v"]), + ) + ) + continue + logger.debug(f"[relay stdout] {line}") + + def _read_stderr(self, stream: IO[str]) -> None: + for line in stream: + line = line.rstrip() + if line: + self._stderr_tail.append(line) + logger.debug(f"[relay stderr] {line}") diff --git a/dimos/web/relay_bridge/test_relay_e2e.py b/dimos/web/relay_bridge/test_relay_e2e.py new file mode 100644 index 0000000000..59e223e819 --- /dev/null +++ b/dimos/web/relay_bridge/test_relay_e2e.py @@ -0,0 +1,311 @@ +# Copyright 2026 Dimensional Inc. +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""End-to-end tests against a real relay child process (aioquic both legs). + +One file on purpose: --dist=loadfile keeps the module-scoped relay on a +single xdist worker. +""" + +import asyncio +from collections.abc import AsyncIterator, Callable, Iterator +import hashlib +import json +import statistics +import time +import urllib.request + +import pytest + +from dimos.web.relay_bridge.protocol import DataFrame, FrameHeader +from dimos.web.relay_bridge.relay_process import RelayProcess, RelayReadyInfo +from dimos.web.relay_bridge.wt_client import RelayClient + + +@pytest.fixture(scope="module") +def relay() -> Iterator[RelayReadyInfo]: + process = RelayProcess() + try: + yield process.start() + finally: + process.stop() + + +@pytest.fixture +def own_relay() -> Iterator[RelayProcess]: + """A started relay process the test may stop itself (stop() is idempotent).""" + process = RelayProcess() + try: + process.start() + yield process + finally: + process.stop() + + +@pytest.fixture +async def robot(relay: RelayReadyInfo) -> AsyncIterator[RelayClient]: + """A connected robot client with the hello handshake done.""" + async with await RelayClient.connect(relay.wt_url, "robot") as client: + await client.hello() + yield client + + +@pytest.fixture +async def viewer(relay: RelayReadyInfo) -> AsyncIterator[RelayClient]: + """A connected viewer client with the hello handshake done.""" + async with await RelayClient.connect(relay.wt_url, "viewer") as client: + await client.hello() + yield client + + +async def collect_until( + viewer: RelayClient, + done: Callable[[list[DataFrame]], bool], + timeout: float = 10.0, +) -> list[DataFrame]: + """Consume viewer frames until `done(frames)` or `timeout` (returns what arrived).""" + frames: list[DataFrame] = [] + + async def _consume() -> None: + async for frame in viewer.frames(): + frames.append(frame) + if done(frames): + return + + try: + await asyncio.wait_for(_consume(), timeout) + except asyncio.TimeoutError: + pass + return frames + + +async def fetch_stats(relay: RelayReadyInfo) -> dict: + def _get() -> dict: + with urllib.request.urlopen(f"http://127.0.0.1:{relay.http_port}/api/stats") as response: + return json.load(response) + + return await asyncio.to_thread(_get) + + +def test_info_matches_ready_line(relay: RelayReadyInfo) -> None: + with urllib.request.urlopen(f"http://127.0.0.1:{relay.http_port}/api/info") as response: + info = json.load(response) + assert info == {"wtUrl": f"{relay.wt_url}/viewer", "certHash": relay.cert_hash, "v": relay.v} + assert relay.wt_url.startswith("https://127.0.0.1:") + + +async def test_robot_handshake_and_datagram_rtt(relay: RelayReadyInfo) -> None: + # Connects manually: the hello handshake itself is under test here. + async with await RelayClient.connect(relay.wt_url, "robot") as robot: + await robot.hello() + rtts = [await robot.ping() for _ in range(20)] + assert statistics.median(rtts) < 0.1 + + +async def test_reliable_channel_is_complete_and_intact( + robot: RelayClient, viewer: RelayClient +) -> None: + count = 100 + payloads = [seq.to_bytes(4, "little") * 256 for seq in range(count)] + for seq, payload in enumerate(payloads): + robot.send_frame("odom", payload, delivery="reliable", meta={"i": seq}) + + frames = await collect_until( + viewer, + lambda fs: len({f.header.seq for f in fs if f.header.ch == "odom"}) >= count, + ) + odom = {f.header.seq: f for f in frames if f.header.ch == "odom"} + # Reliable = complete, no drops. One-stream-per-message may reorder; + # completeness is the contract, headers carry the sequence. + assert sorted(odom) == list(range(count)) + assert all(bytes(odom[seq].payload) == payloads[seq] for seq in range(count)) + assert odom[0].header.delivery == "reliable" + assert odom[0].header.meta == {"i": 0} + + +async def test_latest_channel_newest_wins(robot: RelayClient, viewer: RelayClient) -> None: + writer = robot.latest_writer("cam") + offered = 200 + for i in range(offered): + writer.offer(i.to_bytes(4, "little") + b"\xab" * 2000) + # Yield so the pump interleaves with the offers; without this all + # 200 land in one loop turn and the mailbox collapses to sent=1, + # never exercising the concurrent send-while-in-flight path. + await asyncio.sleep(0) + + def newest_arrived(frames: list[DataFrame]) -> bool: + return any( + f.header.ch == "cam" and f.payload[:4] == (offered - 1).to_bytes(4, "little") + for f in frames + ) + + frames = await collect_until(viewer, newest_arrived) + cam = [f for f in frames if f.header.ch == "cam"] + markers = [int.from_bytes(bytes(f.payload[:4]), "little") for f in cam] + # The newest offered frame always lands; the mailbox shed the rest. + assert newest_arrived(frames), f"newest frame missing; got markers {markers}" + assert writer.dropped + writer.sent == offered + assert 0 < len(cam) <= offered + # The interleaving must actually exercise multiple sends (the old + # single-turn version guaranteed sent==1). + assert writer.sent >= 2, f"pump never interleaved; sent={writer.sent}" + # Everything the writer actually sent arrived (loopback: no transport loss). + assert len(cam) == writer.sent + + +async def test_large_frame_1mib(robot: RelayClient, viewer: RelayClient) -> None: + payload = bytes(range(256)) * 4096 # 1 MiB + robot.send_frame("blob", payload, delivery="reliable") + frames = await collect_until(viewer, lambda fs: any(f.header.ch == "blob" for f in fs)) + blob = next(f for f in frames if f.header.ch == "blob") + assert len(blob.payload) == len(payload) + assert hashlib.sha256(blob.payload).hexdigest() == hashlib.sha256(payload).hexdigest() + + +async def test_reset_stale_discards_partial_frame(robot: RelayClient, viewer: RelayClient) -> None: + """A reset mid-frame must drop the partial on the relay and nothing else.""" + # 8 MiB cannot be flushed + ACKed within the same event-loop turn, so + # the reset below reliably lands mid-transfer. + big = robot.send_frame("cam", b"\xcd" * (8 * 1024 * 1024), delivery="latest") + assert robot._session.reset_if_in_flight(big) + small = b"\x01\x02\x03\x04" * 8 + robot.send_frame("cam", small, delivery="latest") + + frames = await collect_until(viewer, lambda fs: any(f.header.ch == "cam" for f in fs)) + cam = [f for f in frames if f.header.ch == "cam"] + assert [bytes(f.payload) for f in cam] == [small] + + # The relay survived the reset: control still answers. + assert await robot.ping() < 5.0 + + +async def test_reset_burst_does_not_wedge_robot_leg( + robot: RelayClient, viewer: RelayClient +) -> None: + """Resets racing stream acceptance must not kill the relay's robot data path. + + A stream reset before the relay has read its WebTransport preamble errors + Deno's wt.incomingBidirectionalStreams permanently (rejected pull), which + used to silently end the robot stream loop. Bursting resets in the same + event-loop turn as the sends makes that race near-certain. + """ + for rnd in range(5): + # The accept glue cannot have read all 50 preambles before the + # resets land, so some streams are reset pre-acceptance. + ids = [robot.send_frame("cam", b"\xcd" * (16 * 1024), delivery="latest") for _ in range(50)] + for stream_id in ids: + robot._session.reset_if_in_flight(stream_id) + marker = f"alive-{rnd}".encode() + robot.send_frame("cam", marker, delivery="latest") + + frames = await collect_until( + viewer, + lambda fs, marker=marker: any(bytes(f.payload) == marker for f in fs), + timeout=5.0, + ) + assert any(bytes(f.payload) == marker for f in frames), ( + f"robot data path wedged in round {rnd}" + ) + + +async def test_stats_reflect_traffic( + relay: RelayReadyInfo, robot: RelayClient, viewer: RelayClient +) -> None: + robot.send_frame("odom", b"{}", delivery="reliable") + await collect_until(viewer, lambda fs: len(fs) >= 1, timeout=5.0) + + stats = await fetch_stats(relay) + assert stats["robot"] is True + assert stats["viewers"] >= 1 + assert stats["channels"]["odom"]["framesIn"] >= 1 + assert stats["channels"]["odom"]["delivery"] == "reliable" + + +async def test_send_frame_paces_with_wait_delivered(robot: RelayClient) -> None: + start = time.monotonic() + stream_id = robot.send_frame("odom", b"x" * 1000, delivery="reliable") + assert await robot.wait_delivered(stream_id, timeout=5.0) + assert time.monotonic() - start < 5.0 + + +async def test_malformed_robot_frame_is_dropped( + relay: RelayReadyInfo, robot: RelayClient, viewer: RelayClient +) -> None: + """A well-framed frame with an invalid header is dropped, not fatal.""" + before = (await fetch_stats(relay)).get("framesDropped", 0) + # The client validates delivery, so model_construct skips it to put a + # bogus value on the wire; the relay's validator must reject it. + bad_header = FrameHeader.model_construct(ch="cam", seq=0, ts=time.time(), delivery="bogus") + bad_id = robot._session.send_frame(bad_header, b"junk") + assert await robot.wait_delivered(bad_id, timeout=5.0) + # A following valid frame proves the channel still forwards. + robot.send_frame("cam", b"good", delivery="reliable") + frames = await collect_until(viewer, lambda fs: any(f.header.ch == "cam" for f in fs)) + cam = [bytes(f.payload) for f in frames if f.header.ch == "cam"] + assert cam == [b"good"], f"only the valid frame should forward, got {cam}" + + # The drop was counted (poll: onRobotFrame runs just after the ACK). + after = before + for _ in range(100): + after = (await fetch_stats(relay)).get("framesDropped", 0) + if after - before >= 1: + break + await asyncio.sleep(0.05) + assert after - before == 1 + # The session survived the bad frame: control still answers. + assert await robot.ping() < 5.0 + + +async def test_latest_writer_resets_stale_stream(robot: RelayClient, viewer: RelayClient) -> None: + """The writer auto-resets an in-flight stream when a newer frame is waiting.""" + writer = robot.latest_writer("cam", stale_after=0.02) + # 8 MiB can't flush + ACK within stale_after, so it stays in flight. + writer.offer(b"\xcd" * (8 * 1024 * 1024)) + # Wait until the pump has begun sending the big frame. + for _ in range(1000): + if writer.sent >= 1: + break + await asyncio.sleep(0.005) + assert writer.sent >= 1, "pump never sent the first frame" + # A newer small frame makes the stalled big stream stale -> reset. + writer.offer(b"\x01\x02\x03\x04") + + frames = await collect_until( + viewer, lambda fs: any(f.header.ch == "cam" and len(f.payload) < 100 for f in fs) + ) + cam = [f for f in frames if f.header.ch == "cam"] + assert cam, "no cam frame reached the viewer" + # Only the small frame arrives; the 8 MiB frame was reset mid-flight. + assert all(len(f.payload) < 100 for f in cam) + assert writer.resets >= 1, "the stale stream was never reset" + # The relay survived the reset. + assert await robot.ping() < 5.0 + + +async def test_close_signal_stops_writer_and_wakes_waiter(own_relay: RelayProcess) -> None: + """Relay death terminates the connection, wakes wait_closed, stops the pump.""" + async with await RelayClient.connect(own_relay.info.wt_url, "robot") as robot: + await robot.hello() + writer = robot.latest_writer("cam") + writer.offer(b"x" * 1000) + await asyncio.sleep(0.1) # let the pump start + own_relay.stop() # graceful shutdown sends CONNECTION_CLOSE + + await asyncio.wait_for(robot.wait_closed(), timeout=10.0) + assert robot.is_closed + await asyncio.sleep(0.1) # let the pump observe the close + assert writer._task.done() + # A dead channel is visible at the producer. + with pytest.raises(RuntimeError): + writer.offer(b"y") diff --git a/dimos/web/relay_bridge/test_wt_client.py b/dimos/web/relay_bridge/test_wt_client.py index 2d62966a84..b84b31aa28 100644 --- a/dimos/web/relay_bridge/test_wt_client.py +++ b/dimos/web/relay_bridge/test_wt_client.py @@ -111,7 +111,7 @@ async def consume_forever() -> None: async for _ in client.frames(): pass - consumer = asyncio.ensure_future(consume_forever()) + consumer = asyncio.create_task(consume_forever()) await asyncio.sleep(0.05) # consumer now suspended waiting on an empty queue consumer.cancel() with pytest.raises(asyncio.CancelledError): diff --git a/dimos/web/relay_bridge/wt_client.py b/dimos/web/relay_bridge/wt_client.py index acc45b1041..25d182b125 100644 --- a/dimos/web/relay_bridge/wt_client.py +++ b/dimos/web/relay_bridge/wt_client.py @@ -201,10 +201,10 @@ async def frames(self) -> AsyncIterator[DataFrame]: consumer mid-iteration never orphans the queue getter (which would steal the next delivered frame). """ - closed = asyncio.ensure_future(self._session.wait_closed()) + closed = asyncio.create_task(self._session.wait_closed()) try: while True: - get = asyncio.ensure_future(self._session.frames.get()) + get = asyncio.create_task(self._session.frames.get()) try: await asyncio.wait({get, closed}, return_when=asyncio.FIRST_COMPLETED) # get first: a frame buffered before close still drains. @@ -241,7 +241,7 @@ def __init__(self, client: RelayClient, ch: str, *, stale_after: float) -> None: self.sent = 0 self.resets = 0 self._mailbox: asyncio.Queue[tuple[bytes, dict[str, Any] | None]] = asyncio.Queue(maxsize=1) - self._task = asyncio.ensure_future(self._pump()) + self._task = asyncio.create_task(self._pump()) self._task.add_done_callback(self._on_pump_done) def _on_pump_done(self, task: asyncio.Future[None]) -> None: @@ -273,10 +273,10 @@ def stop(self) -> None: async def _pump(self) -> None: session = self._client._session - closed = asyncio.ensure_future(session.closed.wait()) + closed = asyncio.create_task(session.closed.wait()) try: while not session.closed.is_set(): - get = asyncio.ensure_future(self._mailbox.get()) + get = asyncio.create_task(self._mailbox.get()) try: await asyncio.wait({get, closed}, return_when=asyncio.FIRST_COMPLETED) if not get.done(): diff --git a/setup.py b/setup.py index f9083995ba..30acc36df7 100644 --- a/setup.py +++ b/setup.py @@ -58,6 +58,13 @@ def python_is_macos_universal_binary(executable: str | None = None) -> bool: TEST_MODULE_PATTERNS = ("test_*.py", "conftest.py") +# The Deno relay (repo-root web/) ships inside the wheel so a pip-installed +# dimos can run it without a checkout. Copied into build_lib below; editable +# installs skip the copy and locate.find_web_dir() resolves the checkout. +# MANIFEST.in grafts web/ so sdist->wheel builds can reproduce this. +RELAY_DIST_SOURCES = ("deno.json", "deno.lock", "relay", "shared") +RELAY_DIST_TARGET = os.path.join("dimos", "web", "relay_bridge", "_relay_dist") + class build_py(_build_py): def find_package_modules(self, package, package_dir): @@ -69,6 +76,24 @@ def find_package_modules(self, package, package_dir): ) ] + def run(self): + super().run() + if not getattr(self, "editable_mode", False): + self._copy_relay_dist() + + def _copy_relay_dist(self): + src = Path(__file__).parent / "web" + if not (src / "relay" / "main.ts").is_file(): + raise RuntimeError(f"relay sources missing at {src}; refusing to build the wheel") + dst = Path(self.build_lib) / RELAY_DIST_TARGET + for name in RELAY_DIST_SOURCES: + for path in sorted((src / name).rglob("*")) if (src / name).is_dir() else [src / name]: + if path.is_dir() or path.name.endswith("_test.ts"): + continue + target = dst / path.relative_to(src) + self.mkpath(str(target.parent)) + self.copy_file(str(path), str(target)) + extra_compile_args = [ "-O3", # Maximum optimization