Files
api/services/edge-broker/server.mjs
T
Jeppe Bundgaard d30a006457 Add delivery metadata support, preferred channels, and enhanced agent validation
This commit introduces delivery metadata tracking for gateway commands, updates, and shells. Adds preferred delivery channel handling, refined validation for edge agents, improved relay management logic, and broker presence reporting. Includes schema changes, enhanced shell handling, and test coverage.
2026-04-16 13:44:28 +02:00

436 lines
13 KiB
JavaScript

import http from "node:http";
import { randomUUID } from "node:crypto";
import { WebSocketServer } from "ws";
function parseJsonBody(req) {
return new Promise((resolve, reject) => {
let raw = "";
req.on("data", (chunk) => {
raw += chunk.toString("utf8");
});
req.on("end", () => {
try {
resolve(raw === "" ? {} : JSON.parse(raw));
} catch (error) {
reject(error);
}
});
req.on("error", reject);
});
}
function jsonResponse(res, statusCode, body) {
res.writeHead(statusCode, { "content-type": "application/json" });
res.end(JSON.stringify(body));
}
function trimTrailingSlash(value) {
return String(value || "").replace(/\/+$/, "");
}
async function parseJsonResponse(response) {
const text = await response.text();
if (text === "") {
return {};
}
try {
return JSON.parse(text);
} catch {
return { error: text };
}
}
function resolveManagerUrl(options = {}) {
return trimTrailingSlash(options.managerUrl || process.env.EDGE_MANAGER_URL || process.env.EDGE_PUBLIC_API_URL || "");
}
function resolveAuthMode(options = {}, managerUrl = "") {
if (options.authMode) {
return options.authMode;
}
if (process.env.EDGE_AUTH_MODE) {
return process.env.EDGE_AUTH_MODE;
}
return managerUrl ? "manager" : "stub";
}
export function createBrokerServer(options = {}) {
const sharedSecret = options.sharedSecret ?? process.env.EDGE_BROKER_SHARED_SECRET ?? "";
const managerUrl = resolveManagerUrl(options);
const authMode = resolveAuthMode(options, managerUrl);
const commandTimeoutMs = options.commandTimeoutMs ?? 10000;
const agents = new Map();
const pendingCommands = new Map();
const browserSessions = new Map();
const managerRequest = async (path, body = {}) => {
if (!managerUrl) {
throw new Error("Edge manager URL is not configured");
}
const response = await fetch(`${managerUrl}${path}`, {
method: "POST",
headers: {
"content-type": "application/json",
...(sharedSecret ? { "x-edge-broker-secret": sharedSecret } : {}),
},
body: JSON.stringify(body),
});
const json = await parseJsonResponse(response);
if (!response.ok) {
throw new Error(json?.data?.message || json?.message || json?.error || `HTTP ${response.status}`);
}
return json?.data ?? json;
};
const validateAgent =
options.validateAgent ||
(authMode === "stub"
? async ({ gatewayId }) => ({ id: gatewayId, gateway_id: gatewayId })
: async ({ gatewayId, token }) =>
managerRequest(`/edge-agent/internal/gateways/${gatewayId}/validate`, { token }));
const validateShellSession =
options.validateShellSession ||
(authMode === "stub"
? async ({ token }) => ({ id: token, gateway_id: 1, reason: "stub" })
: async ({ token }) => managerRequest("/edge-agent/internal/shell-sessions/validate", { token }));
const closeShellSession =
options.closeShellSession ||
(authMode === "stub"
? async () => ({})
: async (_id, token, transcript, reason) =>
managerRequest("/edge-agent/internal/shell-sessions/close", {
token,
transcript,
reason,
}));
const reportGatewayPresence =
options.reportGatewayPresence ||
(authMode === "stub"
? async () => ({})
: async (gatewayId, { status, connectionId, reason = null, metadata = {} } = {}) =>
managerRequest(`/edge-agent/internal/gateways/${gatewayId}/presence`, {
status,
connection_id: connectionId,
reason,
metadata,
}));
const closeBrowserSession = async (sessionRecord, reason) => {
try {
await closeShellSession(
sessionRecord.session.id,
sessionRecord.ws.sessionToken,
sessionRecord.transcript,
reason
);
} catch {
// Preserve socket teardown even when the manager callback is unavailable.
}
};
const markBrowserSessionsClosed = (gatewayId, reason) => {
for (const sessionRecord of browserSessions.values()) {
if (String(sessionRecord.session.gateway_id) !== String(gatewayId)) {
continue;
}
sessionRecord.closedReason = reason;
if (sessionRecord.ws.readyState < 2) {
sessionRecord.ws.close();
}
}
};
const server = http.createServer(async (req, res) => {
try {
const url = new URL(req.url, "http://localhost");
if (req.method === "POST" && /^\/api\/gateways\/\d+\/commands$/.test(url.pathname)) {
if (sharedSecret && req.headers["x-edge-broker-secret"] !== sharedSecret) {
jsonResponse(res, 403, { error: "Forbidden" });
return;
}
const gatewayId = url.pathname.split("/")[3];
const agent = agents.get(String(gatewayId));
if (!agent || agent.readyState !== 1) {
jsonResponse(res, 503, { error: "Gateway agent is offline" });
return;
}
const body = await parseJsonBody(req);
const commandId = randomUUID();
const promise = new Promise((resolve, reject) => {
const timeout = setTimeout(() => {
pendingCommands.delete(commandId);
reject(new Error("Agent command timed out"));
}, commandTimeoutMs);
pendingCommands.set(commandId, {
resolve,
reject,
timeout,
});
});
agent.send(JSON.stringify({
type: "COMMAND",
commandId,
commandType: body.commandType,
payload: body.payload || {},
jobId: body.jobId ?? null,
}));
try {
const result = await promise;
jsonResponse(res, 200, result);
} catch (error) {
jsonResponse(res, 504, { ok: false, error: error instanceof Error ? error.message : String(error) });
}
return;
}
jsonResponse(res, 404, { error: "Not found" });
} catch (error) {
jsonResponse(res, 500, { error: error instanceof Error ? error.message : String(error) });
}
});
const wss = new WebSocketServer({ noServer: true });
server.on("upgrade", async (req, socket, head) => {
const url = new URL(req.url, "http://localhost");
try {
if (url.pathname === "/ws/agent") {
const gatewayId = String(url.searchParams.get("gatewayId") || "");
const token = String(url.searchParams.get("token") || "");
if (gatewayId === "" || token === "") {
socket.destroy();
return;
}
const gatewayInfo = await validateAgent({ gatewayId, token, headers: req.headers });
wss.handleUpgrade(req, socket, head, (ws) => {
const existing = agents.get(gatewayId);
if (existing && existing.readyState < 2) {
existing.close();
}
ws.gatewayId = gatewayId;
ws.gatewayInfo = gatewayInfo;
ws.connectionId = randomUUID();
agents.set(gatewayId, ws);
reportGatewayPresence(gatewayId, {
status: "connected",
connectionId: ws.connectionId,
metadata: {
remote_address: req.socket.remoteAddress || null,
},
}).catch(() => {});
wss.emit("connection", ws, req);
});
return;
}
if (url.pathname === "/ws/browser-shell") {
const token = String(url.searchParams.get("token") || "");
if (token === "") {
socket.destroy();
return;
}
const session = await validateShellSession({ token, headers: req.headers });
wss.handleUpgrade(req, socket, head, (ws) => {
ws.sessionToken = token;
ws.sessionInfo = session;
browserSessions.set(String(session.id), {
ws,
session,
transcript: "",
closedReason: null,
});
const agent = agents.get(String(session.gateway_id));
if (agent && agent.readyState === 1) {
agent.send(JSON.stringify({
type: "OPEN_ROOT_SHELL",
payload: {
sessionId: String(session.id),
reason: session.reason,
cols: session.metadata?.cols ?? null,
rows: session.metadata?.rows ?? null,
},
}));
} else {
const sessionRecord = browserSessions.get(String(session.id));
if (sessionRecord) {
sessionRecord.closedReason = "agent_offline";
}
ws.close();
}
wss.emit("connection", ws, req);
});
return;
}
} catch {
socket.destroy();
return;
}
socket.destroy();
});
wss.on("connection", (ws) => {
ws.on("message", async (raw) => {
const message = JSON.parse(raw.toString());
if (ws.gatewayId) {
if (message.type === "COMMAND_RESULT") {
const pending = pendingCommands.get(message.commandId);
if (!pending) {
return;
}
clearTimeout(pending.timeout);
pendingCommands.delete(message.commandId);
pending.resolve({
ok: Boolean(message.ok),
payload: message.payload,
error: message.error,
});
return;
}
if (["SHELL_OUTPUT", "SHELL_OPENED", "SHELL_EXIT"].includes(message.type)) {
const sessionRecord = browserSessions.get(String(message.sessionId));
if (!sessionRecord) {
return;
}
if (message.type === "SHELL_OUTPUT") {
sessionRecord.transcript += String(message.data || "");
sessionRecord.ws.send(JSON.stringify({ type: "output", data: String(message.data || "") }));
}
if (message.type === "SHELL_OPENED") {
sessionRecord.ws.send(JSON.stringify({ type: "opened" }));
}
if (message.type === "SHELL_EXIT") {
sessionRecord.ws.send(JSON.stringify({ type: "closed", code: message.code ?? 0 }));
await closeBrowserSession(sessionRecord, "agent_exit");
browserSessions.delete(String(message.sessionId));
if (sessionRecord.ws.readyState < 2) {
sessionRecord.ws.close();
}
}
}
return;
}
if (ws.sessionInfo) {
const sessionId = String(ws.sessionInfo.id);
const agent = agents.get(String(ws.sessionInfo.gateway_id));
if (!agent || agent.readyState !== 1) {
return;
}
if (message.type === "input") {
agent.send(JSON.stringify({
type: "SHELL_INPUT",
payload: {
sessionId,
data: String(message.data || ""),
},
}));
}
if (message.type === "resize") {
agent.send(JSON.stringify({
type: "RESIZE_ROOT_SHELL",
payload: {
sessionId,
cols: Number(message.cols || 0),
rows: Number(message.rows || 0),
},
}));
}
if (message.type === "close") {
agent.send(JSON.stringify({
type: "CLOSE_ROOT_SHELL",
payload: { sessionId },
}));
}
}
});
ws.on("close", async (_code, buffer) => {
const closeReason = buffer?.toString?.("utf8") || null;
if (ws.gatewayId) {
if (agents.get(String(ws.gatewayId)) === ws) {
agents.delete(String(ws.gatewayId));
}
markBrowserSessionsClosed(String(ws.gatewayId), "agent_disconnected");
reportGatewayPresence(String(ws.gatewayId), {
status: "disconnected",
connectionId: ws.connectionId || null,
reason: closeReason || "agent_disconnected",
}).catch(() => {});
return;
}
if (ws.sessionInfo) {
const sessionId = String(ws.sessionInfo.id);
const agent = agents.get(String(ws.sessionInfo.gateway_id));
if (agent && agent.readyState === 1) {
agent.send(JSON.stringify({
type: "CLOSE_ROOT_SHELL",
payload: { sessionId },
}));
}
const sessionRecord = browserSessions.get(sessionId);
if (sessionRecord) {
await closeBrowserSession(sessionRecord, sessionRecord.closedReason || "browser_closed");
browserSessions.delete(sessionId);
}
}
});
});
return {
server,
listen(port = Number(process.env.PORT || 4300)) {
return new Promise((resolve) => {
server.listen(port, () => resolve(server.address()));
});
},
close() {
return new Promise((resolve, reject) => {
for (const agent of agents.values()) {
agent.terminate();
}
for (const session of browserSessions.values()) {
session.ws.terminate();
}
for (const pending of pendingCommands.values()) {
clearTimeout(pending.timeout);
pending.reject(new Error("Broker shutting down"));
}
pendingCommands.clear();
wss.close(() => {
server.close((error) => {
if (error) {
reject(error);
return;
}
resolve();
});
});
});
},
state: {
agents,
browserSessions,
pendingCommands,
managerUrl,
authMode,
},
};
}