""" Bridge to the on-robot agent (agent/x2_agent.py) over plain TCP. This is the bridge that drives a real X2. The dashboard server itself never imports rclpy - it talks JSON-lines to the agent running on PC2. That keeps the dashboard usable on any machine (including Windows) and, more importantly, keeps the web UI up when the robot is powered off so it can say so instead of dying. Reconnection is automatic and continuous: switch the robot off and the UI shows the offline gate; switch it on and the dashboard reattaches on its own. """ from __future__ import annotations import asyncio import base64 import contextlib import json import time from . import netinfo, recovery, settings, x2_spec from .bridge_base import Bridge, CommandResult RECONNECT_MIN = 1.0 RECONNECT_MAX = 6.0 COMMAND_TIMEOUT = 12.0 # asyncio's StreamReader defaults to a 64 KiB limit per line, and readline() # raises once a line exceeds it. Camera frames arrive base64-encoded at roughly # 400 KB and a full ROS graph dump is over 100 KB, so the default silently tore # the link down mid-command. 32 MiB leaves generous headroom. STREAM_LIMIT = 32 * 1024 * 1024 class AgentBridge(Bridge): name = "agent" simulated = False def __init__(self, hub, config: dict): super().__init__(hub, config) self._reader = None self._writer = None self._task: asyncio.Task | None = None self._push_task: asyncio.Task | None = None self._pending: dict[int, asyncio.Future] = {} self._next_id = 1 self._connected = asyncio.Event() self._closing = False self._agent_info: dict = {} # key -> (ts, bytes, content_type, flip_used) self._frames: dict[str, tuple[float, bytes, str, bool | None]] = {} self._attempts = 0 self._last_error = "" self._recovering = False # -- lifecycle ---------------------------------------------------------- @property def host(self) -> str: return (self.config.get("robot_host") or "").strip() @property def port(self) -> int: return int(self.config.get("agent_port") or x2_spec.AGENT_PORT) async def start(self) -> None: if not self.host and not self.config.get("auto_discover"): raise RuntimeError( "No robot host configured. Set it in the Settings tab, run the " "server with --robot , or enable auto-discovery." ) self._closing = False self.state.connection.transport = "agent" self.state.connection.host = f"{self.host}:{self.port}" self.state.connection.simulated = False self.state.connection.online = False self._task = asyncio.create_task(self._link_loop(), name="agent-link") self._push_task = asyncio.create_task(self._push_loop(), name="agent-push") # Give the first connection a moment so the UI opens attached rather # than flashing the offline gate, but never block startup on it. with contextlib.suppress(asyncio.TimeoutError): await asyncio.wait_for(self._connected.wait(), timeout=4.0) async def stop(self) -> None: self._closing = True for task in (self._task, self._push_task): if task: task.cancel() with contextlib.suppress(asyncio.CancelledError, Exception): await task await self._drop_link() self.state.connection.online = False self.state.connection.transport = "none" async def reconfigure(self, config: dict) -> None: """ Adopt new settings, and re-dial immediately if the address changed. Without the re-dial an operator who corrects the robot address sees nothing happen until they restart the bridge by hand, because the link loop only reads the address when it next reconnects - and a healthy connection never reconnects. """ before = (self.host, self.port) self.config = config after = ((config.get("robot_host") or "").strip(), int(config.get("agent_port") or x2_spec.AGENT_PORT)) if before != after and self._writer is not None: await self.hub.emit("info", "robot", f"Robot address changed to {after[0]}:{after[1]} - reconnecting") self._attempts = 0 # Dropping the socket makes _connect_once return; the link loop then # dials the new address on its next pass. await self._drop_link() # -- link --------------------------------------------------------------- async def _link_loop(self) -> None: while not self._closing: try: await self._connect_once() except asyncio.CancelledError: raise except Exception as exc: self._last_error = str(exc) finally: await self._drop_link() if self._closing: return if self.state.connection.online: # We were up and lost it - that is worth saying out loud. self.state.connection.online = False self.state.connection.error = self._last_error or "link lost" await self.hub.emit("warn", "robot", f"Lost connection to the robot ({self._last_error or 'link closed'})") self._attempts += 1 # After a few plain retries, try harder: the robot may have come # back on a different address, or be up with the agent not running. if self._attempts % 4 == 0: await self._attempt_recovery() delay = min(RECONNECT_MAX, RECONNECT_MIN * (1.35 ** min(self._attempts, 8))) await asyncio.sleep(delay) async def _attempt_recovery(self) -> None: """ Try to get the robot back without the operator doing anything. Two failure modes are handled, in order of likelihood after a power cycle: the agent is not running (robot answers SSH but not the agent port), and the robot moved to a different DHCP address. """ if self._recovering: return self._recovering = True try: config = settings.load() port = int(config.get("agent_port") or x2_spec.AGENT_PORT) host = self.host # 1. Is the robot there at all, just without the agent running? if host and config.get("auto_start_agent"): info = await netinfo.probe_host( host, ports=[port, int(config.get("robot_ssh_port") or 22)], timeout=1.5) open_ports = info.get("open_ports") or [] if port in open_ports: return # agent is up; the next retry will land if 22 in open_ports: await self._start_agent_on(host, config, port) return # 2. Otherwise look for it somewhere else on this network. if not config.get("auto_discover"): return self._set_recovery("Looking for the robot on this network…") await self.hub.emit("info", "robot", "Robot not answering - scanning the network for it") candidates = await recovery.find_robot(port, int(config.get("robot_ssh_port") or 22)) candidates = [c for c in candidates if c["host"] != host] running = next((c for c in candidates if c["agent"]), None) if running: await self._adopt(running["host"], "found the agent running there") return if config.get("auto_start_agent"): for candidate in candidates: if not candidate["ssh"]: continue if await self._start_agent_on(candidate["host"], config, port, adopt=True): return self._set_recovery("") except asyncio.CancelledError: raise except Exception as exc: self._set_recovery("") await self.hub.emit("warn", "robot", f"Recovery attempt failed: {exc}") finally: self._recovering = False async def _start_agent_on(self, host: str, config: dict, port: int, adopt: bool = False) -> bool: self._set_recovery(f"Robot is on at {host} but the agent is not running — starting it…") await self.hub.emit("info", "robot", f"{host} is reachable but the agent is down - starting it over SSH") ok, message = await recovery.start_agent(host, config) if not ok: self._set_recovery(message) await self.hub.emit("warn", "robot", message) return False if await recovery.wait_for_agent(host, port, timeout=45.0): self._set_recovery("") if adopt: await self._adopt(host, "started the agent there") else: await self.hub.emit("info", "robot", "Agent started - reconnecting") return True self._set_recovery("Started the agent, but it has not come up yet…") return False async def _adopt(self, host: str, why: str) -> None: """Switch to a newly found address and remember it.""" settings.save({"robot_host": host}) self.config = settings.load() self.state.connection.host = f"{host}:{self.port}" self._attempts = 0 self._set_recovery("") await self.hub.emit("info", "robot", f"Robot found at {host} - {why}") await self.hub.broadcast("settings", settings.describe()) def _set_recovery(self, message: str) -> None: """Progress text for the offline gate, so the operator sees us working.""" self.state.custom["recovery"] = message if message: self.state.connection.error = message async def _connect_once(self) -> None: host, port = self.host, self.port if not host: # Nothing to dial yet - let the recovery pass find one. self.state.connection.error = "Looking for the robot on this network…" raise ConnectionError("no robot address yet") try: self._reader, self._writer = await asyncio.wait_for( asyncio.open_connection(host, port, limit=STREAM_LIMIT), timeout=5.0) except (asyncio.TimeoutError, OSError) as exc: self._last_error = f"{host}:{port} unreachable" self.state.connection.error = ( f"Cannot reach the robot agent at {host}:{port}. " "Is the robot powered on and on this network?" ) raise ConnectionError(self._last_error) from exc was_offline = not self.state.connection.online self._attempts = 0 self._last_error = "" self.state.connection.online = True self.state.connection.since = time.time() self.state.connection.error = "" self._connected.set() if was_offline: await self.hub.emit("info", "robot", f"Connected to robot agent at {host}:{port}") # Re-assert anything that does not survive a reconnect. self.state.source_registered = False asyncio.create_task(self._announce_self()) while not self._closing: line = await self._reader.readline() if not line: raise ConnectionError("agent closed the connection") try: message = json.loads(line) except json.JSONDecodeError: continue self._on_message(message) async def _announce_self(self) -> None: """ Tell the robot this dashboard exists as a control source. The motion controller ignores velocity from a source it has never heard of, so this has to happen before the joystick will do anything. It is not a takeover: registering only makes the dashboard *eligible*, and arbitration still hands control to the highest-priority source that is actively sending - the handheld RC (80) and mobile app (60) both outrank us at 30. Done automatically rather than left to a button, because there is no situation where an operator wants the dashboard connected but not allowed to drive, and "register input source" means nothing to someone who just wants to move the robot. """ source = x2_spec.DASHBOARD_INPUT_SOURCE for attempt in range(3): await asyncio.sleep(1.0 if attempt == 0 else 3.0) if not self.state.connection.online: return result = await self.register_input_source( source["name"], source["priority"], source["timeout"]) if result.ok: return if attempt == 2: await self.hub.emit( "warn", "control", f"Could not announce the dashboard to the robot: {result.message}. " "Driving will be refused until this succeeds - use the Control tab to retry.", ) async def _drop_link(self) -> None: self._connected.clear() writer, self._writer, self._reader = self._writer, None, None if writer is not None: with contextlib.suppress(Exception): writer.close() await writer.wait_closed() for future in self._pending.values(): if not future.done(): future.set_exception(ConnectionError("link dropped")) self._pending.clear() def _on_message(self, message: dict) -> None: kind = message.get("type") if kind == "state": self._apply_state(message.get("data") or {}) elif kind == "hello": self._agent_info = message.get("data") or {} self.state.custom["agent"] = self._agent_info # Report the domain the agent is actually on, rather than whatever # this machine has configured - only the agent's value is real. try: self.state.connection.ros_domain_id = int( self._agent_info.get("ros_domain_id", 0)) except (TypeError, ValueError): pass elif kind == "result": future = self._pending.pop(message.get("id"), None) if future and not future.done(): future.set_result(message) # -- state mapping ------------------------------------------------------ def _apply_state(self, data: dict) -> None: s = self.state s.mode = data.get("mode") or "UNKNOWN" spec = x2_spec.MC_MODE_BY_ID.get(s.mode) s.mode_desc = spec["label"] if spec else s.mode status = data.get("mode_status") s.mode_status = x2_spec.MC_ACTION_STATUS.get(status, str(status or "")) for field in ("battery_pct", "battery_voltage", "battery_current", "battery_temp", "battery_cycles", "pmu_temp", "fan_rpm", "fan_pct"): setattr(s, field, data.get(field)) s.charging = bool(data.get("charging")) rails = {} for spec_rail in x2_spec.PMU_RAILS: live = (data.get("rails") or {}).get(spec_rail["key"]) if not live: continue rails[spec_rail["key"]] = { "label": spec_rail["label"], "nominal": spec_rail["nominal"], "voltage": live.get("voltage"), "current": live.get("current"), "ok": live.get("ok", True), } s.rails = rails s.imu = data.get("imu") or {} s.joints = data.get("joints") or {} s.hand_type = data.get("hand_type") or "None" s.hand_state = data.get("hand_state") or {} s.touch_head = data.get("touch_head") or {} s.velocity = data.get("velocity") or s.velocity s.velocity_command = data.get("velocity_command") or s.velocity_command s.odom = data.get("odom") or s.odom s.volume = data.get("volume") if data.get("volume") is not None else s.volume s.muted = bool(data.get("muted")) s.emoji_id = data.get("emoji_id") s.led = data.get("led") or s.led s.input_source = data.get("input_source") or "" s.topic_stats = data.get("topic_stats") or {} s.custom["pmu_raw"] = data.get("pmu_raw") or {} s.custom["pmu_info"] = data.get("pmu_info") or {} s.custom["cameras"] = data.get("cameras") or {} # Which feeds are actually subscribed on the robot right now. The # Vision tab reads this rather than tracking its own idea of on/off, # so a second browser sees the true state instead of "off". s.custom["streams"] = data.get("streams") or {} s.custom["face_status"] = data.get("face_status") s.custom["battery_capacity_mah"] = data.get("battery_capacity_mah") s.custom["battery_power"] = data.get("battery_power") s.custom["agent_version"] = data.get("agent_version") s.custom["hand_left_type"] = data.get("hand_left_type") s.custom["hand_right_type"] = data.get("hand_right_type") s.updated = time.time() async def _push_loop(self) -> None: """Forward state to browsers and record chart series.""" hz = max(1, int(self.config.get("telemetry_hz", 10))) while True: try: await asyncio.sleep(1.0 / hz) now = time.time() if self.state.connection.online: self._record(now) await self.hub.broadcast("state", self.state.as_dict()) except asyncio.CancelledError: raise except Exception: await asyncio.sleep(1.0) def _record(self, now: float) -> None: for key, value in ( ("battery_pct", self.state.battery_pct), ("battery_voltage", self.state.battery_voltage), ("battery_current", self.state.battery_current), ("battery_temp", self.state.battery_temp), ("pmu_temp", self.state.pmu_temp), ("fan_pct", self.state.fan_pct), ("vel_forward", self.state.velocity.get("forward")), ("vel_lateral", self.state.velocity.get("lateral")), ("vel_angular", self.state.velocity.get("angular")), ): if value is not None: self.hub.record(key, value, now) chest = self.state.imu.get("chest") or {} for key, field in (("imu_roll", "roll"), ("imu_pitch", "pitch"), ("imu_yaw", "yaw")): if field in chest: self.hub.record(key, chest[field], now) # -- command plumbing --------------------------------------------------- async def _send(self, name: str, args: dict | None = None, timeout: float = COMMAND_TIMEOUT) -> CommandResult: if self._writer is None or not self.state.connection.online: return CommandResult.failure( "The robot is not connected. Power it on and wait for the dashboard to reattach." ) message_id = self._next_id self._next_id += 1 future: asyncio.Future = asyncio.get_running_loop().create_future() self._pending[message_id] = future payload = json.dumps({"type": "cmd", "id": message_id, "name": name, "args": args or {}}) + "\n" try: self._writer.write(payload.encode()) await self._writer.drain() except Exception as exc: self._pending.pop(message_id, None) return CommandResult.failure(f"Send failed: {exc}") try: reply = await asyncio.wait_for(future, timeout=timeout) except asyncio.TimeoutError: self._pending.pop(message_id, None) return CommandResult.failure(f"The robot did not answer '{name}' within {timeout:g} s") except ConnectionError: return CommandResult.failure("Connection to the robot dropped mid-command") return CommandResult(bool(reply.get("ok")), reply.get("message", ""), reply.get("detail")) # -- commands ----------------------------------------------------------- async def set_mode(self, mode_id: str) -> CommandResult: result = await self._send("set_mode", {"mode": mode_id}) if result.ok: await self.hub.emit("info", "mode", f"Mode set to {mode_id}") return result async def get_mode(self) -> CommandResult: return await self._send("get_mode") async def set_velocity(self, forward: float, lateral: float, angular: float) -> CommandResult: return await self._send("set_velocity", {"forward": forward, "lateral": lateral, "angular": angular}, timeout=4.0) async def stop_motion(self) -> CommandResult: return await self._send("stop", timeout=4.0) async def play_preset(self, motion: int, area: int, interrupt: bool = True) -> CommandResult: return await self._send("preset", {"motion": motion, "area": area, "interrupt": interrupt}, timeout=15.0) async def set_joints(self, group: str, mode: str, targets: dict, stiffness=None, damping=None) -> CommandResult: return await self._send("set_joints", { "group": group, "mode": mode, "targets": targets, "stiffness": stiffness, "damping": damping, }) async def set_hand(self, side: str, positions: list) -> CommandResult: return await self._send("set_hand", {"side": side, "positions": positions}) async def speak(self, text: str, priority: int = 6, interrupt: bool = False) -> CommandResult: return await self._send("speak", {"text": text, "priority": priority, "interrupt": interrupt}, timeout=15.0) async def set_volume(self, volume: int) -> CommandResult: return await self._send("set_volume", {"volume": volume}) async def set_mute(self, muted: bool) -> CommandResult: return await self._send("set_mute", {"muted": muted}) async def play_emoji(self, emotion_id: int, mode: int = 1, priority: int = 6) -> CommandResult: return await self._send("emoji", {"emotion_id": emotion_id, "mode": mode, "priority": priority}) async def set_led(self, mode: int, r: int, g: int, b: int, priority: int = 6, keep: bool = True) -> CommandResult: return await self._send("led", {"mode": mode, "r": r, "g": g, "b": b, "priority": priority, "keep": keep}) async def register_input_source(self, name: str, priority: int, timeout: int) -> CommandResult: result = await self._send("register_source", {"name": name, "priority": priority, "timeout": timeout}) if result.ok: self.state.source_registered = True self.state.input_source = name return result async def publish_raw(self, topic: str, message_type: str, payload: dict) -> CommandResult: return await self._send("publish_raw", {"topic": topic, "type": message_type, "fields": payload}) async def call_service(self, service: str, service_type: str, payload: dict) -> CommandResult: return CommandResult.failure( "Arbitrary service calls are not exposed by the agent. Use the typed commands, " "or add a handler to agent/x2_agent.py." ) async def list_topics(self) -> list[dict]: return list(self.state.topic_stats.values()) async def graph(self) -> dict: result = await self._send("graph", timeout=10.0) return result.detail if result.ok and result.detail else {"topics": [], "services": []} # -- on-demand streams -------------------------------------------------- async def set_stream(self, key: str, active: bool) -> CommandResult: """Switch a camera or the LiDAR on or off on the robot.""" result = await self._send("stream_set", {"key": key, "active": bool(active)}, timeout=8.0) if result.ok: # A stopped feed must not keep serving its last frame from cache. if not active: self._frames.pop(key, None) await self.hub.emit("info", "vision", f"{key} switched {'on' if active else 'off'}") return result async def list_streams(self) -> CommandResult: return await self._send("stream_list", timeout=6.0) async def lidar_points(self) -> CommandResult: return await self._send("lidar_points", timeout=8.0) # -- camera ------------------------------------------------------------- async def camera_frame(self, camera_key: str, stream: str = "rgb", flip: bool | None = None) -> bytes | None: cached = self._frames.get(camera_key) # The agent pushes on request; a 120 ms cache stops several browser tabs # from each pulling the same frame off the robot. if cached and (time.time() - cached[0]) < 0.12 and cached[3] == flip: return cached[1] args: dict = {"key": camera_key} if flip is not None: args["flip"] = bool(flip) result = await self._send("camera_frame", args, timeout=6.0) if not result.ok or not result.detail: return None try: data = base64.b64decode(result.detail["b64"]) except (KeyError, ValueError): return None fmt = result.detail.get("format", "jpeg") self._frames[camera_key] = (time.time(), data, "image/png" if fmt == "png" else "image/jpeg", flip) return data def frame_content_type(self, camera_key: str = "", stream: str = "rgb") -> str: cached = self._frames.get(camera_key) return cached[2] if cached else "image/jpeg"