diff --git a/services/edge-broker/server.mjs b/services/edge-broker/server.mjs index 81576304..7bd16008 100644 --- a/services/edge-broker/server.mjs +++ b/services/edge-broker/server.mjs @@ -254,16 +254,27 @@ export function createBrokerServer(options = {}) { } }; + const closeBrowserShellSocket = (sessionRecord, reason, message = null, code = 1000) => { + sessionRecord.closedReason = reason; + if (sessionRecord.ws.readyState >= 2) { + return; + } + + sendJson(sessionRecord.ws, { + type: "closed", + reason, + ...(message ? { message } : {}), + }); + sessionRecord.ws.close(code, reason); + }; + const markGatewayShellSessionsClosed = (gatewayId, reason) => { for (const sessionRecord of browserShellSessions.values()) { if (String(sessionRecord.session.gateway_id) !== String(gatewayId)) { continue; } - sessionRecord.closedReason = reason; - if (sessionRecord.ws.readyState < 2) { - sessionRecord.ws.close(); - } + closeBrowserShellSocket(sessionRecord, reason, "Gateway agent disconnected from the broker.", 1011); } }; @@ -448,6 +459,7 @@ export function createBrokerServer(options = {}) { closedReason: null, }; browserShellSessions.set(String(session.id), sessionRecord); + wss.emit("connection", ws, req); const agent = agents.get(String(session.gateway_id)); if (agent && agent.readyState === 1) { @@ -464,10 +476,13 @@ export function createBrokerServer(options = {}) { }, }); } else { - sessionRecord.closedReason = "agent_offline"; - ws.close(); + closeBrowserShellSocket( + sessionRecord, + "agent_offline", + "Gateway agent is not connected to the broker.", + 1011 + ); } - wss.emit("connection", ws, req); }); return; } @@ -624,11 +639,11 @@ export function createBrokerServer(options = {}) { } if (message.type === "SHELL_EXIT") { - sendJson(sessionRecord.ws, { type: "closed", code: message.code ?? 0 }); + sendJson(sessionRecord.ws, { type: "closed", reason: "agent_exit", code: message.code ?? 0 }); await closeBrowserShellSession(sessionRecord, "agent_exit"); browserShellSessions.delete(String(message.sessionId)); if (sessionRecord.ws.readyState < 2) { - sessionRecord.ws.close(); + sessionRecord.ws.close(1000, "agent_exit"); } } } diff --git a/services/edge-broker/test/broker.test.mjs b/services/edge-broker/test/broker.test.mjs index d005df47..78bb7e05 100644 --- a/services/edge-broker/test/broker.test.mjs +++ b/services/edge-broker/test/broker.test.mjs @@ -109,7 +109,7 @@ test("broker bridges browser shell sessions through the connected agent", async 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.code === 0)); + 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"); @@ -177,10 +177,19 @@ test("broker closes browser shell sessions immediately when no agent is connecte 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(); @@ -203,6 +212,7 @@ test("broker closes browser shell sessions when the agent disconnects before she 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 waitForMessage(agent); @@ -211,6 +221,7 @@ test("broker closes browser shell sessions when the agent disconnects before she 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();