Files
api/services/edge-broker/test/broker.test.mjs
T
16048e2ce3 chore(release): merge develop into master — XL Vask flag text fix + edge-broker health (Aug 15 2026) (#376)
Brings all of the develop branch's commits into master.

## What this contains

The 2 commits on develop that landed during the XL Vask integration
dispatch:

- **PR #373** (TRU-6 / AUT-2) — feat(edge-broker): expose lastActivityAt
on /api/health (AUT-2/TRU-6)
- **PR #375** (TRU-49 / AUT-49) — fix(api): include wash_id in
xlvask_missing_order_link flag text (AUT-49/TRU-49)

## Why

The XL Vask integration dispatch via the OpenSymphony orchestrator
(MiniMax M3) produced 2 api-side fixes:
- **PR #373** — adds `lastActivityAt` to the api health endpoint so
operators can see if the edge-broker has processed any requests
recently.
- **PR #375** — the actual root-cause fix for the user-reported symptom
"XL Vask-registreringen er hverken ignoreret eller knyttet til en ordre
i den valgte periode doesn't show the wash". The bug was in
`messageParts()` for the `xlvask_missing_order_link` arm — the link text
was hard-coded to 'XL Vask wash' instead of using the actual wash_id.
This PR makes the link identify the wash it points to.

## Verification

Both source PRs passed:
- Required CI (PHP unit, PHP integration, PHP api, PHP legacy, edge
broker, edge agent, edge gateway backend)
- The api ruleset allows squash merges

## Notes

- The pleno-vue repo has its own equivalent develop→master PR (#312)
with the 9 UI fixes (component, i18n, and a Playwright E2E).

---------

Co-authored-by: Jeppe <jeppe@copenhagentruckwash.io>
Co-authored-by: openhands <openhands@all-hands.dev>
2026-08-16 09:08:03 +02:00

945 lines
32 KiB
JavaScript

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);
assert.equal(typeof healthJson.lastActivityAt, "string");
assert.match(healthJson.lastActivityAt, /^\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}\.\d{3}Z$/);
assert.ok(healthJson.lastActivityAt >= broker.state.containerStartedAt);
assert.equal(healthJson.lastActivityAt, broker.state.lastActivityAt);
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 updates lastActivityAt after each successful request", async () => {
const broker = createBrokerServer({ authMode: "manager", sharedSecret: "secret", managerUrl: "http://manager.test" });
const address = await broker.listen(0);
const port = address.port;
assert.equal(broker.state.lastActivityAt, broker.state.containerStartedAt);
const firstResponse = await fetch(`http://127.0.0.1:${port}/api/health`);
const firstJson = await firstResponse.json();
const firstActivityAt = broker.state.lastActivityAt;
assert.equal(typeof firstJson.lastActivityAt, "string");
assert.equal(firstJson.lastActivityAt, firstActivityAt);
assert.ok(firstActivityAt >= broker.state.containerStartedAt);
await new Promise((resolve) => setTimeout(resolve, 5));
await fetch(`http://127.0.0.1:${port}/api/diagnostics/shared-secret`, {
method: "POST",
headers: {
"x-edge-broker-secret": "secret",
},
});
assert.notEqual(broker.state.lastActivityAt, firstActivityAt);
assert.ok(broker.state.lastActivityAt > firstActivityAt);
const secondResponse = await fetch(`http://127.0.0.1:${port}/api/health`);
const secondJson = await secondResponse.json();
assert.equal(secondJson.lastActivityAt, broker.state.lastActivityAt);
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();
});