import WebSocket from "ws"; import http from "node:http"; import type { Agent as HttpAgent } from "node:http"; import type { RosbridgeClientOptions, ConnectionStatus, RosbridgeMessage, MessageHandler, ConnectionHandler, } from "./types.js"; /** Pending request/response tracker for service calls and action goals. */ export interface PendingRequest { resolve: (value: unknown) => void; reject: (reason: Error) => void; timer: ReturnType; } /** * Build an http.Agent that tunnels every connect() through an HTTP CONNECT * proxy when the environment defines one. * * Why this exists: * * * Node 22+ honors HTTP(S)_PROXY for ``https.request`` (via the experimental * ``EnvHttpProxyAgent``) but ``ws://`` connections from the ``ws`` library * bypass that and try a direct TCP connect to the resolved IP. In a * NemoClaw / OpenShell sandbox the gateway runs in a restricted netns * where the only routable peer is the OPA policy proxy at * ``10.200.0.1:3128``. Direct connects to anything else fail with * ``ECONNREFUSED`` — that's exactly what happens when ``rosbridge.url`` * points at ``ws://host.docker.internal:9090`` from inside the sandbox. * * * Adding an http-proxy-agent dep would balloon the offline-deploy bundle * and is unnecessary — Node's built-in ``http.request`` with method * ``CONNECT`` does exactly the same handshake. * * Returns ``null`` (let ``ws`` use its default Agent) when: * * * No proxy env var is set, or * * The target host appears in ``NO_PROXY`` (case-insensitive match, * ``*`` and leading ``.`` wildcards honored RFC-style). * * The plugin owns its own reconnect/cleanup, so we keep the Agent * stateless and use ``keepAlive: false`` — every connect attempt opens a * fresh CONNECT tunnel. */ function buildProxyAgentForUrl(wsUrl: string): HttpAgent | null { // When NODE_USE_ENV_PROXY=1 (Node 22+ / NemoClaw sandbox), the runtime's // built-in EnvHttpProxyAgent handles ws:// tunneling natively and adds the // authentication headers the OpenShell proxy requires. A manual CONNECT // request (our TunnelAgent below) omits those headers and gets 403 from // the proxy. Skip our agent and let Node handle it. if (process.env.NODE_USE_ENV_PROXY === "1") return null; let parsed: URL; try { parsed = new URL(wsUrl); } catch { return null; } const targetHost = parsed.hostname; const targetPort = parsed.port || (parsed.protocol === "wss:" ? "443" : "80"); const proxyUrl = process.env.HTTPS_PROXY || process.env.https_proxy || process.env.HTTP_PROXY || process.env.http_proxy || ""; if (!proxyUrl) return null; let proxy: URL; try { proxy = new URL(proxyUrl); } catch { return null; } const noProxy = (process.env.NO_PROXY || process.env.no_proxy || "").toLowerCase(); if (noProxy) { const target = targetHost.toLowerCase(); const entries = noProxy.split(",").map((s) => s.trim()).filter(Boolean); for (const entry of entries) { if (entry === "*") return null; const e = entry.startsWith(".") ? entry.slice(1) : entry; if (target === e || target.endsWith("." + e)) return null; } } class TunnelAgent extends http.Agent { createConnection( opts: { hostname?: string; host?: string; port?: number | string }, cb: (err: Error | null, socket?: NodeJS.Socket) => void, ): undefined { const host = opts.hostname || opts.host || targetHost; const port = String(opts.port ?? targetPort); const req = http.request({ host: proxy.hostname, port: proxy.port ? Number(proxy.port) : proxy.protocol === "https:" ? 443 : 80, method: "CONNECT", path: `${host}:${port}`, headers: { Host: `${host}:${port}` }, }); req.once("connect", (res, socket) => { if (res.statusCode === 200) { cb(null, socket); } else { socket.destroy(); cb( new Error( `HTTP CONNECT to proxy ${proxy.host} failed: ${res.statusCode} ${res.statusMessage ?? ""}`.trim(), ), ); } }); req.once("error", (err) => cb(err)); req.end(); return undefined; } } return new TunnelAgent({ keepAlive: false }); } /** * WebSocket client for the rosbridge protocol. * Handles connection lifecycle, reconnection, and message routing. */ export class RosbridgeClient { private ws: WebSocket | null = null; private options: Required; private status: ConnectionStatus = "disconnected"; private messageHandlers = new Map>(); private connectionHandlers = new Set(); private pendingRequests = new Map(); private reconnectAttempts = 0; private reconnectTimer: ReturnType | null = null; private intentionalClose = false; private idCounter = 0; constructor(options: RosbridgeClientOptions) { this.options = { url: options.url, reconnect: options.reconnect ?? true, reconnectInterval: options.reconnectInterval ?? 3000, maxReconnectAttempts: options.maxReconnectAttempts ?? 10, }; } /** Connect to the rosbridge WebSocket server. */ async connect(): Promise { if (this.status === "connected") return; this.intentionalClose = false; this.setStatus("connecting"); return new Promise((resolve, reject) => { const connectTimeout = setTimeout(() => { if (this.ws) { this.ws.close(); this.ws = null; } reject(new Error(`Connection to ${this.options.url} timed out`)); }, 10_000); try { const proxyAgent = buildProxyAgentForUrl(this.options.url); this.ws = proxyAgent ? new WebSocket(this.options.url, { agent: proxyAgent }) : new WebSocket(this.options.url); } catch (err) { clearTimeout(connectTimeout); this.setStatus("disconnected"); reject(err); return; } this.ws.onopen = () => { clearTimeout(connectTimeout); this.reconnectAttempts = 0; this.setStatus("connected"); resolve(); }; this.ws.onmessage = (event) => { const data = typeof event.data === "string" ? event.data : event.data.toString(); this.handleMessage(data); }; this.ws.onerror = (event: { message?: string; error?: Error }) => { clearTimeout(connectTimeout); if (this.status === "connecting") { this.ws = null; this.setStatus("disconnected"); const underlying = (event?.error && (event.error.message || String(event.error))) || event?.message || "no underlying error message"; reject(new Error(`WebSocket error connecting to ${this.options.url}: ${underlying}`)); } }; this.ws.onclose = (event: { code?: number; reason?: Buffer | string }) => { clearTimeout(connectTimeout); this.ws = null; if (this.status === "connecting") { this.setStatus("disconnected"); const code = event?.code ?? "n/a"; const reasonRaw = event?.reason; const reasonStr = typeof reasonRaw === "string" ? reasonRaw : reasonRaw && typeof (reasonRaw as Buffer).toString === "function" ? (reasonRaw as Buffer).toString() : ""; reject( new Error( `WebSocket closed during connection to ${this.options.url} (code=${code}${reasonStr ? `, reason="${reasonStr}"` : ""})`, ), ); return; } this.setStatus("disconnected"); this.rejectAllPending(new Error("WebSocket connection closed")); if (!this.intentionalClose && this.options.reconnect) { this.attemptReconnect(); } }; }); } /** Disconnect from the rosbridge server. */ async disconnect(): Promise { this.intentionalClose = true; if (this.reconnectTimer) { clearTimeout(this.reconnectTimer); this.reconnectTimer = null; } this.rejectAllPending(new Error("Client disconnected")); if (this.ws) { const ws = this.ws; this.ws = null; if (ws.readyState === WebSocket.OPEN || ws.readyState === WebSocket.CONNECTING) { await new Promise((resolve) => { ws.onclose = () => resolve(); ws.close(); // Force-resolve after 2s if server doesn't ack the close setTimeout(resolve, 2000); }); } } this.setStatus("disconnected"); } /** Send a rosbridge protocol message. */ send(message: RosbridgeMessage & Record): void { if (!this.ws || this.status !== "connected") { throw new Error("Not connected to rosbridge server"); } this.ws.send(JSON.stringify(message)); } /** Generate a unique message ID. */ nextId(prefix = "agenticros"): string { return `${prefix}_${++this.idCounter}`; } /** Subscribe to messages on a specific topic. */ onMessage(topic: string, handler: MessageHandler): () => void { if (!this.messageHandlers.has(topic)) { this.messageHandlers.set(topic, new Set()); } this.messageHandlers.get(topic)!.add(handler); return () => { this.messageHandlers.get(topic)?.delete(handler); }; } /** Register a connection status change handler. */ onConnection(handler: ConnectionHandler): () => void { this.connectionHandlers.add(handler); return () => { this.connectionHandlers.delete(handler); }; } /** Get current connection status. */ getStatus(): ConnectionStatus { return this.status; } /** * Register a pending request that will be resolved when a response * with the matching ID arrives (service_response or action_result). */ registerPending(id: string, resolve: (value: unknown) => void, reject: (reason: Error) => void, timeoutMs = 30_000): void { const timer = setTimeout(() => { this.pendingRequests.delete(id); reject(new Error(`Request ${id} timed out after ${timeoutMs}ms`)); }, timeoutMs); this.pendingRequests.set(id, { resolve, reject, timer }); } /** Resolve a pending request by ID and clean up. */ resolvePending(id: string, result: unknown): void { const pending = this.pendingRequests.get(id); if (pending) { clearTimeout(pending.timer); this.pendingRequests.delete(id); pending.resolve(result); } } /** Reject a pending request by ID and clean up. */ rejectPending(id: string, error: Error): void { const pending = this.pendingRequests.get(id); if (pending) { clearTimeout(pending.timer); this.pendingRequests.delete(id); pending.reject(error); } } private setStatus(status: ConnectionStatus): void { this.status = status; for (const handler of this.connectionHandlers) { handler(status); } } /** Route an incoming rosbridge message to the appropriate handler. */ private handleMessage(data: string): void { let msg: Record; try { msg = JSON.parse(data) as Record; } catch { return; // Ignore malformed messages } const op = msg.op as string | undefined; if (!op) return; switch (op) { case "publish": { // Incoming topic message — route to topic subscribers const topic = msg.topic as string; const payload = msg.msg as Record; const handlers = this.messageHandlers.get(topic); if (handlers) { for (const handler of handlers) { handler(payload); } } break; } case "service_response": { // Response to a call_service request const id = msg.id as string | undefined; if (id) { this.resolvePending(id, msg); } break; } case "action_result": { // Final result of an action goal const id = msg.id as string | undefined; if (id) { this.resolvePending(id, msg); } break; } case "action_feedback": { // Intermediate feedback for an action goal — route to feedback handlers const id = msg.id as string | undefined; if (id) { const handlers = this.messageHandlers.get(`__action_feedback__${id}`); if (handlers) { for (const handler of handlers) { handler(msg as Record); } } } break; } } } /** Attempt to reconnect with exponential backoff. */ private attemptReconnect(): void { if (this.intentionalClose) return; if (this.reconnectAttempts >= this.options.maxReconnectAttempts) return; this.reconnectAttempts++; // Exponential backoff: interval * 2^(attempt-1), capped at 30s const delay = Math.min( this.options.reconnectInterval * Math.pow(2, this.reconnectAttempts - 1), 30_000, ); this.reconnectTimer = setTimeout(async () => { this.reconnectTimer = null; if (this.intentionalClose) return; try { await this.connect(); // Re-subscribe to all active topics on successful reconnect for (const topic of this.messageHandlers.keys()) { if (topic.startsWith("__action_feedback__")) continue; this.send({ op: "subscribe", id: this.nextId("resub"), topic, }); } } catch { // connect() failed — onclose will trigger another attemptReconnect } }, delay); } /** Reject all pending requests (used on disconnect/close). */ private rejectAllPending(error: Error): void { for (const [id, pending] of this.pendingRequests) { clearTimeout(pending.timer); pending.reject(error); } this.pendingRequests.clear(); } }