import test from "node:test"; import assert from "node:assert/strict"; import net from "node:net"; import WebSocket from "ws"; import { createBrokerServer } from "../server.mjs"; function collectMessages(socket) { const messages = []; socket.on("message", (raw) => { messages.push(JSON.parse(raw.toString())); }); return messages; } function waitForClose(socket) { return new Promise((resolve) => { socket.once("close", resolve); }); } function waitForCloseOrError(socket) { return new Promise((resolve) => { const onDone = () => { socket.off("error", onDone); socket.off("close", onDone); resolve(); }; socket.once("error", onDone); socket.once("close", onDone); }); } function rawUpgradeRequest(port, path) { return new Promise((resolve, reject) => { const socket = net.createConnection({ host: "127.0.0.1", port }, () => { socket.write( [ `GET ${path} HTTP/1.1`, `Host: 127.0.0.1:${port}`, "Connection: Upgrade", "Upgrade: websocket", "Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==", "Sec-WebSocket-Version: 13", "", "", ].join("\r\n") ); }); let response = ""; socket.setEncoding("utf8"); socket.on("data", (chunk) => { response += chunk; }); socket.on("end", () => resolve(response)); socket.on("error", reject); }); } async function waitFor(predicate, { timeoutMs = 1000, intervalMs = 10, description = "condition" } = {}) { const deadline = Date.now() + timeoutMs; while (Date.now() < deadline) { if (predicate()) { return; } await new Promise((resolve) => setTimeout(resolve, intervalMs)); } throw new Error(`Timed out waiting for ${description}`); } test("broker defaults to strict auth and fails closed when manager URL is missing", async () => { const previousEnv = { EDGE_AUTH_MODE: process.env.EDGE_AUTH_MODE, EDGE_MANAGER_URL: process.env.EDGE_MANAGER_URL, EDGE_PUBLIC_API_URL: process.env.EDGE_PUBLIC_API_URL, }; delete process.env.EDGE_AUTH_MODE; delete process.env.EDGE_MANAGER_URL; delete process.env.EDGE_PUBLIC_API_URL; let broker; try { broker = createBrokerServer({ sharedSecret: "secret" }); assert.equal(broker.state.authMode, "strict"); assert.equal(broker.state.managerUrl, ""); const address = await broker.listen(0); const port = address.port; const shellResponse = await rawUpgradeRequest(port, "/ws/browser-shell?token=session-token"); const agentResponse = await rawUpgradeRequest(port, "/ws/agent?gatewayId=701&token=agent-token"); assert.doesNotMatch(shellResponse, /101 Switching Protocols/); assert.match(shellResponse, /^HTTP\/1\.1 401 Unauthorized/m); assert.match(shellResponse, /"error_code":"shell_session_invalid"/); assert.match(shellResponse, /"message":"Shell session could not be validated\."/); assert.doesNotMatch(shellResponse, /Edge manager URL is not configured/); assert.doesNotMatch(agentResponse, /101 Switching Protocols/); assert.match(agentResponse, /^HTTP\/1\.1 503 Service Unavailable/m); assert.match(agentResponse, /"error_code":"agent_validation_failed"/); assert.match(agentResponse, /"stage":"agent_validate"/); assert.match(agentResponse, /"message":"Gateway agent could not be validated\."/); assert.doesNotMatch(agentResponse, /Edge manager URL is not configured/); } finally { if (broker) { await broker.close(); } for (const [key, value] of Object.entries(previousEnv)) { if (value === undefined) { delete process.env[key]; } else { process.env[key] = value; } } } }); test("broker rejects protected HTTP endpoints when shared secret is missing", async () => { const broker = createBrokerServer({ authMode: "stub", sharedSecret: "", commandTimeoutMs: 2000 }); const address = await broker.listen(0); const port = address.port; const agent = new WebSocket(`ws://127.0.0.1:${port}/ws/agent?gatewayId=701&token=agent-token`); await new Promise((resolve) => agent.once("open", resolve)); const agentMessages = collectMessages(agent); const commandResponse = await fetch(`http://127.0.0.1:${port}/api/gateways/701/commands`, { method: "POST", headers: { "content-type": "application/json", }, body: JSON.stringify({ commandType: "SET_RELAY_STATE", payload: { relayId: "M-7", on: true }, }), }); const commandJson = await commandResponse.json(); assert.equal(commandResponse.status, 503); assert.equal(commandJson.ok, false); assert.equal(commandJson.shared_secret_required, true); assert.match(commandJson.error, /shared secret is not configured/); assert.equal(agentMessages.some((message) => message.type === "COMMAND"), false); const diagnosticsResponse = await fetch(`http://127.0.0.1:${port}/api/diagnostics/shared-secret`, { method: "POST", }); const diagnosticsJson = await diagnosticsResponse.json(); assert.equal(diagnosticsResponse.status, 503); assert.equal(diagnosticsJson.shared_secret_required, true); const syncResponse = await fetch(`http://127.0.0.1:${port}/api/gateways/701/sync`, { method: "POST", }); const syncJson = await syncResponse.json(); assert.equal(syncResponse.status, 503); assert.equal(syncJson.shared_secret_required, true); agent.terminate(); await broker.close(); }); test("broker dispatches commands to connected agents", async () => { const broker = createBrokerServer({ authMode: "stub", sharedSecret: "secret", commandTimeoutMs: 2000 }); const address = await broker.listen(0); const port = address.port; const agent = new WebSocket(`ws://127.0.0.1:${port}/ws/agent?gatewayId=701&token=agent-token`); await new Promise((resolve) => agent.once("open", resolve)); agent.on("message", (raw) => { const message = JSON.parse(raw.toString()); if (message.type === "COMMAND") { agent.send(JSON.stringify({ type: "COMMAND_RESULT", commandId: message.commandId, ok: true, payload: { online: true, on: true }, })); } }); const response = await fetch(`http://127.0.0.1:${port}/api/gateways/701/commands`, { method: "POST", headers: { "content-type": "application/json", "x-edge-broker-secret": "secret", }, body: JSON.stringify({ commandType: "GET_RELAY_STATUS", payload: { relayId: "M-7" }, }), }); const json = await response.json(); assert.equal(response.status, 200); assert.equal(json.ok, true); assert.equal(json.payload.on, true); agent.terminate(); await broker.close(); }); test("broker exposes health and shared-secret diagnostics", async () => { const broker = createBrokerServer({ authMode: "manager", sharedSecret: "secret", managerUrl: "http://manager.test" }); const address = await broker.listen(0); const port = address.port; const healthResponse = await fetch(`http://127.0.0.1:${port}/api/health`); const healthJson = await healthResponse.json(); assert.equal(healthResponse.status, 200); assert.equal(healthJson.ok, true); assert.equal(healthJson.service, "edge-broker"); assert.equal(healthJson.auth_mode, "manager"); assert.equal(healthJson.manager_url_configured, true); assert.equal(healthJson.shared_secret_configured, true); const invalidSecretResponse = await fetch(`http://127.0.0.1:${port}/api/diagnostics/shared-secret`, { method: "POST", headers: { "x-edge-broker-secret": "wrong-secret", }, }); const invalidSecretJson = await invalidSecretResponse.json(); assert.equal(invalidSecretResponse.status, 403); assert.equal(invalidSecretJson.ok, false); assert.equal(invalidSecretJson.shared_secret_required, true); const validSecretResponse = await fetch(`http://127.0.0.1:${port}/api/diagnostics/shared-secret`, { method: "POST", headers: { "x-edge-broker-secret": "secret", }, }); const validSecretJson = await validSecretResponse.json(); assert.equal(validSecretResponse.status, 200); assert.equal(validSecretJson.ok, true); assert.equal(validSecretJson.shared_secret_required, true); await broker.close(); }); test("broker bridges browser shell sessions through the connected agent", async () => { const closedSessions = []; const broker = createBrokerServer({ authMode: "stub", validateShellSession: async () => ({ id: "shell-1", gateway_id: "701", reason: "diagnostic" }), closeShellSession: async (_id, _token, transcript, reason) => { closedSessions.push({ transcript, reason }); }, }); const address = await broker.listen(0); const port = address.port; const agent = new WebSocket(`ws://127.0.0.1:${port}/ws/agent?gatewayId=701&token=agent-token`); await new Promise((resolve) => agent.once("open", resolve)); agent.on("message", (raw) => { const message = JSON.parse(raw.toString()); if (message.type === "OPEN_ROOT_SHELL") { agent.send(JSON.stringify({ type: "SHELL_OPENED", sessionId: "shell-1" })); agent.send(JSON.stringify({ type: "SHELL_OUTPUT", sessionId: "shell-1", data: "root@pi:~# " })); agent.send(JSON.stringify({ type: "SHELL_EXIT", sessionId: "shell-1", code: 0 })); } }); const browser = new WebSocket(`ws://127.0.0.1:${port}/ws/browser-shell?token=session-token`); const browserMessages = collectMessages(browser); await new Promise((resolve) => browser.once("open", resolve)); await new Promise((resolve) => setTimeout(resolve, 100)); assert.ok(browserMessages.some((message) => message.type === "opened")); assert.ok(browserMessages.some((message) => message.type === "output" && /root@pi/.test(message.data))); assert.ok(browserMessages.some((message) => message.type === "closed" && message.reason === "agent_exit" && message.code === 0)); assert.equal(closedSessions.length, 1); assert.equal(closedSessions[0].reason, "agent_exit"); browser.terminate(); agent.terminate(); await broker.close(); }); test("broker forwards browser shell input, resize, and close events to the agent", async () => { const broker = createBrokerServer({ authMode: "stub", validateShellSession: async () => ({ id: "shell-2", gateway_id: "701", reason: "diagnostic" }), }); const address = await broker.listen(0); const port = address.port; const agent = new WebSocket(`ws://127.0.0.1:${port}/ws/agent?gatewayId=701&token=agent-token`); await new Promise((resolve) => agent.once("open", resolve)); const agentMessages = collectMessages(agent); const browser = new WebSocket(`ws://127.0.0.1:${port}/ws/browser-shell?token=session-token`); await new Promise((resolve) => browser.once("open", resolve)); await waitFor( () => agentMessages.some((message) => message.type === "OPEN_ROOT_SHELL" && message.payload.sessionId === "shell-2"), { description: "agent shell open request" } ); browser.send(JSON.stringify({ type: "input", data: "ls\r" })); browser.send(JSON.stringify({ type: "resize", cols: 140, rows: 44 })); browser.send(JSON.stringify({ type: "close" })); await new Promise((resolve) => setTimeout(resolve, 50)); assert.ok(agentMessages.some((message) => message.type === "OPEN_ROOT_SHELL" && message.payload.sessionId === "shell-2")); assert.ok( agentMessages.some( (message) => message.type === "SHELL_INPUT" && message.payload.sessionId === "shell-2" && message.payload.data === "ls\r" ) ); assert.ok( agentMessages.some( (message) => message.type === "RESIZE_ROOT_SHELL" && message.payload.sessionId === "shell-2" && message.payload.cols === 140 && message.payload.rows === 44 ) ); assert.ok(agentMessages.some((message) => message.type === "CLOSE_ROOT_SHELL" && message.payload.sessionId === "shell-2")); browser.terminate(); agent.terminate(); await broker.close(); }); test("broker closes browser shell sessions immediately when no agent is connected", async () => { const closedSessions = []; const broker = createBrokerServer({ authMode: "stub", validateShellSession: async () => ({ id: "shell-3", gateway_id: "701", reason: "diagnostic" }), closeShellSession: async (_id, _token, transcript, reason) => { closedSessions.push({ transcript, reason }); }, }); const address = await broker.listen(0); const port = address.port; const browser = new WebSocket(`ws://127.0.0.1:${port}/ws/browser-shell?token=session-token`); const browserMessages = collectMessages(browser); await new Promise((resolve) => browser.once("open", resolve)); await waitForClose(browser); await waitFor(() => closedSessions.length === 1, { description: "offline shell session close callback" }); assert.ok( browserMessages.some( (message) => message.type === "closed" && message.reason === "agent_offline" && /not connected to the broker/i.test(message.message) ) ); assert.deepEqual(closedSessions, [{ transcript: "", reason: "agent_offline" }]); await broker.close(); }); test("broker rejects browser shell upgrades without a token using HTTP diagnostics", async () => { const broker = createBrokerServer({ authMode: "stub" }); const address = await broker.listen(0); const port = address.port; const response = await rawUpgradeRequest(port, "/ws/browser-shell"); assert.match(response, /^HTTP\/1\.1 400 Bad Request/m); assert.match(response, /"error_code":"shell_session_token_missing"/); assert.match(response, /"message":"Missing shell session token\."/); await broker.close(); }); test("broker rejects invalid browser shell upgrades without leaking the token", async () => { const broker = createBrokerServer({ authMode: "stub", validateShellSession: async () => { const error = new Error("Shell session expired"); error.status = 401; error.code = "shell_session_expired"; throw error; }, }); const address = await broker.listen(0); const port = address.port; const rawToken = "session-token-secret"; const response = await rawUpgradeRequest(port, `/ws/browser-shell?token=${rawToken}`); assert.match(response, /^HTTP\/1\.1 401 Unauthorized/m); assert.match(response, /"error_code":"shell_session_expired"/); assert.match(response, /"message":"Shell session could not be validated\."/); assert.doesNotMatch(response, /Shell session expired/); assert.doesNotMatch(response, new RegExp(rawToken)); await broker.close(); }); test("broker rejects websocket upgrade errors without exposing exception text", async () => { const broker = createBrokerServer({ authMode: "stub", validateBrowserStream: async () => { throw new Error("UPSTREAM-SENSITIVE: redis://cache.internal:6379 timeout"); }, }); const address = await broker.listen(0); const port = address.port; const response = await rawUpgradeRequest(port, "/ws/browser-gateway-stream?token=session-token"); assert.match(response, /^HTTP\/1\.1 500 Internal Server Error/m); assert.match(response, /"error_code":"websocket_upgrade_failed"/); assert.match(response, /"message":"WebSocket upgrade failed\."/); assert.doesNotMatch(response, /UPSTREAM-SENSITIVE/); assert.doesNotMatch(response, /redis:\/\/cache\.internal/); await broker.close(); }); test("broker closes browser shell sessions when the agent never reports shell opened", async () => { const closedSessions = []; const broker = createBrokerServer({ authMode: "stub", shellOpenTimeoutMs: 30, validateShellSession: async () => ({ id: "shell-timeout", gateway_id: "701", reason: "diagnostic" }), closeShellSession: async (_id, _token, transcript, reason, details) => { closedSessions.push({ transcript, reason, details }); }, }); const address = await broker.listen(0); const port = address.port; const agent = new WebSocket(`ws://127.0.0.1:${port}/ws/agent?gatewayId=701&token=agent-token`); const agentMessages = collectMessages(agent); await new Promise((resolve) => agent.once("open", resolve)); const browser = new WebSocket(`ws://127.0.0.1:${port}/ws/browser-shell?token=session-token`); const browserMessages = collectMessages(browser); await new Promise((resolve) => browser.once("open", resolve)); await waitFor( () => agentMessages.some( (message) => message.type === "OPEN_ROOT_SHELL" && message.payload.sessionId === "shell-timeout" ), { description: "agent shell open request before timeout" } ); await waitForClose(browser); await waitFor(() => closedSessions.length === 1, { description: "timeout shell session close callback" }); assert.ok(browserMessages.some((message) => message.type === "closed" && message.reason === "shell_open_timeout")); assert.equal(closedSessions[0].reason, "shell_open_timeout"); assert.equal(closedSessions[0].details.stage, "shell_open"); assert.equal(closedSessions[0].details.code, 1011); agent.terminate(); await broker.close(); }); test("broker closes browser shell sessions when the agent disconnects before shell open", async () => { const closedSessions = []; const broker = createBrokerServer({ authMode: "stub", validateShellSession: async () => ({ id: "shell-4", gateway_id: "701", reason: "diagnostic" }), closeShellSession: async (_id, _token, transcript, reason) => { closedSessions.push({ transcript, reason }); }, }); const address = await broker.listen(0); const port = address.port; const agent = new WebSocket(`ws://127.0.0.1:${port}/ws/agent?gatewayId=701&token=agent-token`); await new Promise((resolve) => agent.once("open", resolve)); const agentMessages = collectMessages(agent); const browser = new WebSocket(`ws://127.0.0.1:${port}/ws/browser-shell?token=session-token`); const browserMessages = collectMessages(browser); await new Promise((resolve) => browser.once("open", resolve)); await waitFor( () => agentMessages.some((message) => message.type === "OPEN_ROOT_SHELL" && message.payload.sessionId === "shell-4"), { description: "agent shell open request before disconnect" } ); agent.terminate(); await waitForClose(browser); await waitFor(() => closedSessions.length === 1, { description: "disconnect shell session close callback" }); assert.ok(agentMessages.some((message) => message.type === "OPEN_ROOT_SHELL" && message.payload.sessionId === "shell-4")); assert.ok(browserMessages.some((message) => message.type === "closed" && message.reason === "agent_disconnected")); assert.deepEqual(closedSessions, [{ transcript: "", reason: "agent_disconnected" }]); await broker.close(); }); test("broker defaults to strict auth when no validators are configured", async () => { const broker = createBrokerServer(); const address = await broker.listen(0); const port = address.port; const agent = new WebSocket(`ws://127.0.0.1:${port}/ws/agent?gatewayId=701&token=agent-token`); await waitForCloseOrError(agent); await broker.close(); }); test("broker sends an agent welcome before connection progress and backlog dispatch", async () => { const broker = createBrokerServer({ authMode: "stub", reportGatewayPresence: async () => ({}), requestGatewayBacklog: async () => ({ dispatch: [ { type: "TASK_DISPATCH", taskType: "OPERATION", operation: { id: 91, type: "DISCOVERY", request: {}, }, }, ], }), }); const address = await broker.listen(0); const port = address.port; const agent = new WebSocket(`ws://127.0.0.1:${port}/ws/agent?gatewayId=701&token=agent-token`); const messages = collectMessages(agent); await new Promise((resolve) => agent.once("open", resolve)); await waitFor( () => messages.some((message) => message.type === "TASK_DISPATCH" && message.operation?.id === 91), { description: "backlog dispatch after connection welcome" } ); assert.equal(messages[0].type, "WELCOME"); assert.equal(messages[0].gatewayId, "701"); assert.equal(typeof messages[0].connectionId, "string"); assert.ok(messages[0].connectionId.length > 0); const firstProgressIndex = messages.findIndex((message) => message.type === "CONNECTION_PROGRESS"); const dispatchIndex = messages.findIndex((message) => message.type === "TASK_DISPATCH"); assert.ok(firstProgressIndex > 0); assert.ok(dispatchIndex > firstProgressIndex); assert.ok( messages.some( (message) => message.type === "CONNECTION_PROGRESS" && message.stage === "presence" && message.status === "started" ) ); assert.ok( messages.some( (message) => message.type === "CONNECTION_PROGRESS" && message.stage === "backlog" && message.status === "started" ) ); assert.ok( messages.some( (message) => message.type === "CONNECTION_PROGRESS" && message.stage === "backlog" && message.status === "succeeded" && message.dispatchCount === 1 ) ); agent.terminate(); await broker.close(); }); test("broker sends connection errors to agents after the welcome message", async () => { const broker = createBrokerServer({ authMode: "stub", reportGatewayPresence: async () => { throw new Error("presence callback unavailable"); }, requestGatewayBacklog: async () => { throw new Error("backlog sync unavailable"); }, }); const address = await broker.listen(0); const port = address.port; const agent = new WebSocket(`ws://127.0.0.1:${port}/ws/agent?gatewayId=701&token=agent-token`); const messages = collectMessages(agent); await new Promise((resolve) => agent.once("open", resolve)); await waitFor( () => messages.filter((message) => message.type === "CONNECTION_ERROR").length >= 2, { description: "agent connection error notifications" } ); assert.equal(messages[0].type, "WELCOME"); assert.ok( messages.some( (message) => message.type === "CONNECTION_ERROR" && message.stage === "presence" && message.error === "presence callback unavailable" ) ); assert.ok( messages.some( (message) => message.type === "CONNECTION_ERROR" && message.stage === "backlog" && message.error === "backlog sync unavailable" ) ); agent.terminate(); await broker.close(); }); test("broker syncs queued gateway backlog on agent connect and manual sync", async () => { const backlogRequests = []; const broker = createBrokerServer({ authMode: "stub", sharedSecret: "secret", requestGatewayBacklog: async (gatewayId, payload) => { backlogRequests.push({ gatewayId, payload }); return { dispatch: [ { type: "TASK_DISPATCH", taskType: "OPERATION", operation: { id: 91, type: "DISCOVERY", request: {}, }, }, ], }; }, }); const address = await broker.listen(0); const port = address.port; const agent = new WebSocket( `ws://127.0.0.1:${port}/ws/agent?gatewayId=701&token=agent-token&agentInstanceId=instance-1` ); const messages = collectMessages(agent); await new Promise((resolve) => agent.once("open", resolve)); await waitFor( () => messages.some((message) => message.type === "TASK_DISPATCH" && message.operation?.id === 91), { description: "initial backlog dispatch" } ); assert.equal(backlogRequests.length, 1); assert.equal(backlogRequests[0].gatewayId, "701"); assert.equal(backlogRequests[0].payload.agent_instance_id, "instance-1"); const response = await fetch(`http://127.0.0.1:${port}/api/gateways/701/sync`, { method: "POST", headers: { "content-type": "application/json", "x-edge-broker-secret": "secret", }, body: JSON.stringify({ gatewayId: 701 }), }); const json = await response.json(); assert.equal(response.status, 200); assert.equal(json.ok, true); await waitFor(() => backlogRequests.length >= 2, { description: "manual sync backlog request" }); agent.terminate(); await broker.close(); }); test("broker fans out telemetry, task, log, and presence updates to browser gateway streams", async () => { const telemetryPayloads = []; const broker = createBrokerServer({ authMode: "stub", validateAgent: async () => ({ id: "701", gateway_id: "701", label: "CPH Edge 01" }), validateBrowserStream: async () => ({ id: "stream-1", gateway_id: "701", scopes: ["overview", "tasks", "logs", "statistics"], }), ingestTelemetry: async (_gatewayId, payload) => { telemetryPayloads.push(payload); return { gateway: { id: 701, metadata: payload.metadata || {} } }; }, ingestTaskEvent: async (_gatewayId, operationId, payload) => ({ id: operationId, status: "IN_PROGRESS", latest_event: payload, }), ingestLogEntry: async (_gatewayId, payload) => ({ id: 5001, ...payload, }), }); const address = await broker.listen(0); const port = address.port; const browser = new WebSocket(`ws://127.0.0.1:${port}/ws/browser-gateway-stream?token=stream-token`); const browserMessages = collectMessages(browser); await new Promise((resolve) => browser.once("open", resolve)); const agent = new WebSocket(`ws://127.0.0.1:${port}/ws/agent?gatewayId=701&token=agent-token`); await new Promise((resolve) => agent.once("open", resolve)); agent.send( JSON.stringify({ type: "TELEMETRY", payload: { status: "ONLINE", metadata: { broker_connected: true, }, }, }) ); agent.send( JSON.stringify({ type: "TASK_EVENT", operationId: 41, payload: { level: "INFO", code: "DISCOVERY_RUNNING", message: "Discovery is running", }, }) ); agent.send( JSON.stringify({ type: "LOG_FRAME", payload: { level: "INFO", stream: "agent", source: "EDGE_AGENT", message: "Gateway heartbeat acknowledged", }, }) ); await waitFor( () => browserMessages.some((message) => message.type === "presence.changed" && message.status === "connected"), { description: "presence update" } ); assert.ok(browserMessages.some((message) => message.type === "gateway.telemetry")); assert.ok(browserMessages.some((message) => message.type === "stats.updated")); assert.ok(browserMessages.some((message) => message.type === "task.updated" && message.operationId === 41)); assert.ok(browserMessages.some((message) => message.type === "log.append" && /heartbeat/.test(message.entry?.message))); assert.equal(telemetryPayloads[0]?.broker_connection_id?.length > 0, true); browser.terminate(); agent.terminate(); await broker.close(); }); test("broker ingests self-serve machine signals from connected agents", async () => { const machineSignals = []; const broker = createBrokerServer({ authMode: "stub", validateAgent: async () => ({ id: "701", gateway_id: "701", label: "CPH Edge 01" }), ingestMachineSignal: async (gatewayId, payload) => { machineSignals.push({ gatewayId, payload }); return { recorded: true, lane_id: payload.lane_id }; }, }); const address = await broker.listen(0); const port = address.port; const agent = new WebSocket(`ws://127.0.0.1:${port}/ws/agent?gatewayId=701&token=agent-token`); await new Promise((resolve) => agent.once("open", resolve)); agent.send( JSON.stringify({ type: "MACHINE_SIGNAL", payload: { lane_id: 3, relay_id: "machine-relay", component: "input", channel: 0, event: "input.toggle_on", state: true, }, }) ); await waitFor(() => machineSignals.length === 1, { description: "machine signal ingestion" }); assert.equal(machineSignals[0].gatewayId, "701"); assert.deepEqual(machineSignals[0].payload, { lane_id: 3, relay_id: "machine-relay", component: "input", channel: 0, event: "input.toggle_on", state: true, }); agent.terminate(); await broker.close(); }); test("broker survives telemetry ingestion failures for stale gateways", async () => { const broker = createBrokerServer({ authMode: "stub", validateAgent: async () => ({ id: "701", gateway_id: "701", label: "CPH Edge 01" }), ingestTelemetry: async () => { throw new Error("Edge gateway not found"); }, }); const address = await broker.listen(0); const port = address.port; const agent = new WebSocket(`ws://127.0.0.1:${port}/ws/agent?gatewayId=701&token=agent-token`); await new Promise((resolve) => agent.once("open", resolve)); agent.send( JSON.stringify({ type: "TELEMETRY", payload: { status: "ONLINE", }, }) ); await new Promise((resolve) => setTimeout(resolve, 100)); assert.equal(broker.server.listening, true); agent.terminate(); await broker.close(); }); test("broker still fans out telemetry when manager ingestion fails", async () => { const broker = createBrokerServer({ authMode: "stub", validateAgent: async () => ({ id: "701", gateway_id: "701", label: "CPH Edge 01" }), validateBrowserStream: async () => ({ id: "stream-telemetry-fallback", gateway_id: "701", scopes: ["overview", "statistics"], }), ingestTelemetry: async () => { throw new Error("manager unavailable"); }, }); const address = await broker.listen(0); const port = address.port; const browser = new WebSocket(`ws://127.0.0.1:${port}/ws/browser-gateway-stream?token=stream-token`); const browserMessages = collectMessages(browser); await new Promise((resolve) => browser.once("open", resolve)); const agent = new WebSocket(`ws://127.0.0.1:${port}/ws/agent?gatewayId=701&token=agent-token`); await new Promise((resolve) => agent.once("open", resolve)); agent.send( JSON.stringify({ type: "TELEMETRY", payload: { status: "ONLINE", metadata: { system_metrics: { cpu_usage_pct: 31, }, container_health: { state: "ONLINE", summary: "6/6 containers healthy", }, }, }, }) ); await waitFor( () => browserMessages.some( (message) => message.type === "gateway.telemetry" && message.error === "Telemetry ingestion failed" ), { description: "telemetry fanout after ingest failure" } ); assert.ok( browserMessages.every((message) => message.error !== "manager unavailable"), "raw manager errors must not be sent to browser streams" ); assert.ok( browserMessages.some( (message) => message.type === "stats.updated" && message.statistics?.system_metrics?.cpu_usage_pct === 31 ) ); browser.terminate(); agent.terminate(); await broker.close(); });