Files
api/services/edge-broker/server.mjs
T

863 lines
28 KiB
JavaScript

import http from "node:http";
import { randomUUID } from "node:crypto";
import { fileURLToPath } from "node:url";
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";
}
function parseScopes(value) {
if (!Array.isArray(value)) {
return [];
}
return Array.from(
new Set(
value
.map((scope) => String(scope || "").trim().toLowerCase())
.filter(Boolean)
)
);
}
function eventScopes(message) {
switch (message?.type) {
case "gateway.telemetry":
case "presence.changed":
return ["overview", "statistics"];
case "task.updated":
return ["tasks", "overview"];
case "log.append":
return ["logs"];
case "stats.updated":
return ["statistics"];
default:
return [];
}
}
function sessionAllowsScopes(sessionRecord, scopes) {
const subscriptions = sessionRecord.subscriptions || new Set();
if (subscriptions.has("*")) {
return true;
}
if (!scopes || scopes.length === 0) {
return true;
}
return scopes.some((scope) => subscriptions.has(scope));
}
function sendJson(ws, payload) {
if (!ws || ws.readyState !== 1) {
return false;
}
ws.send(JSON.stringify(payload));
return true;
}
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 browserShellSessions = new Map();
const browserStreamSessions = new Map();
const gatewayStreamSessions = new Map();
const inflightGatewaySyncs = new Map();
const managerRequest = async (path, body = {}, method = "POST") => {
if (!managerUrl) {
throw new Error("Edge manager URL is not configured");
}
const response = await fetch(`${managerUrl}${path}`, {
method,
headers: {
"content-type": "application/json",
...(sharedSecret ? { "x-edge-broker-secret": sharedSecret } : {}),
},
body: method === "GET" ? undefined : 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, label: `Gateway ${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", cols: 120, rows: 32 })
: async ({ token }) => managerRequest("/edge-agent/internal/shell-sessions/validate", { token }));
const markShellSessionOpened =
options.markShellSessionOpened ||
(authMode === "stub"
? async () => ({})
: async (token, connectionId) =>
managerRequest("/edge-agent/internal/shell-sessions/opened", {
token,
connection_id: connectionId,
}));
const closeShellSession =
options.closeShellSession ||
(authMode === "stub"
? async () => ({})
: async (_id, token, transcript, reason) =>
managerRequest("/edge-agent/internal/shell-sessions/close", {
token,
transcript,
reason,
}));
const validateBrowserStream =
options.validateBrowserStream ||
(authMode === "stub"
? async ({ token }) => ({
id: token,
gateway_id: 1,
scopes: ["overview", "tasks", "logs", "statistics"],
})
: async ({ token }) => managerRequest("/edge-agent/internal/browser-streams/validate", { token }));
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 requestGatewayBacklog =
options.requestGatewayBacklog ||
(authMode === "stub"
? async () => ({ gateway: {}, dispatch: [] })
: async (gatewayId, payload = {}) =>
managerRequest(`/edge-agent/internal/gateways/${gatewayId}/backlog`, payload));
const ingestTelemetry =
options.ingestTelemetry ||
(authMode === "stub"
? async (_gatewayId, payload = {}) => payload
: async (gatewayId, payload = {}) =>
managerRequest(`/edge-agent/internal/gateways/${gatewayId}/telemetry`, payload));
const ingestTaskEvent =
options.ingestTaskEvent ||
(authMode === "stub"
? async (_gatewayId, _operationId, payload = {}) => payload
: async (gatewayId, operationId, payload = {}) =>
managerRequest(`/edge-agent/internal/gateways/${gatewayId}/operations/${operationId}/events`, payload));
const ingestTaskResult =
options.ingestTaskResult ||
(authMode === "stub"
? async (_gatewayId, _operationId, payload = {}) => payload
: async (gatewayId, operationId, payload = {}) =>
managerRequest(`/edge-agent/internal/gateways/${gatewayId}/operations/${operationId}/complete`, payload));
const ingestLogEntry =
options.ingestLogEntry ||
(authMode === "stub"
? async (_gatewayId, payload = {}) => payload
: async (gatewayId, payload = {}) =>
managerRequest(`/edge-agent/internal/gateways/${gatewayId}/logs`, payload));
const broadcastGatewayEvent = (gatewayId, message) => {
const sessionIds = gatewayStreamSessions.get(String(gatewayId));
if (!sessionIds || sessionIds.size === 0) {
return;
}
const allowedScopes = eventScopes(message);
for (const sessionId of sessionIds.values()) {
const sessionRecord = browserStreamSessions.get(String(sessionId));
if (!sessionRecord) {
continue;
}
if (!sessionAllowsScopes(sessionRecord, allowedScopes)) {
continue;
}
sendJson(sessionRecord.ws, message);
}
};
const closeBrowserShellSession = 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 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;
}
closeBrowserShellSocket(sessionRecord, reason, "Gateway agent disconnected from the broker.", 1011);
}
};
const registerGatewayStreamSession = (sessionRecord) => {
const gatewayId = String(sessionRecord.session.gateway_id);
if (!gatewayStreamSessions.has(gatewayId)) {
gatewayStreamSessions.set(gatewayId, new Set());
}
gatewayStreamSessions.get(gatewayId).add(String(sessionRecord.session.id));
browserStreamSessions.set(String(sessionRecord.session.id), sessionRecord);
};
const removeGatewayStreamSession = (sessionRecord) => {
browserStreamSessions.delete(String(sessionRecord.session.id));
const gatewayId = String(sessionRecord.session.gateway_id);
const sessionIds = gatewayStreamSessions.get(gatewayId);
if (!sessionIds) {
return;
}
sessionIds.delete(String(sessionRecord.session.id));
if (sessionIds.size === 0) {
gatewayStreamSessions.delete(gatewayId);
}
};
const syncGatewayBacklog = async (gatewayId, explicitAgent = null) => {
const normalizedGatewayId = String(gatewayId);
const agent = explicitAgent || agents.get(normalizedGatewayId);
if (!agent || agent.readyState !== 1) {
return { queued: false };
}
if (inflightGatewaySyncs.has(normalizedGatewayId)) {
return inflightGatewaySyncs.get(normalizedGatewayId);
}
const syncPromise = (async () => {
const backlog = await requestGatewayBacklog(normalizedGatewayId, {
agent_instance_id: agent.agentInstanceId || null,
});
const dispatch = Array.isArray(backlog?.dispatch) ? backlog.dispatch : [];
for (const instruction of dispatch) {
sendJson(agent, instruction);
}
return {
queued: dispatch.length > 0,
dispatch,
};
})().finally(() => {
inflightGatewaySyncs.delete(normalizedGatewayId);
});
inflightGatewaySyncs.set(normalizedGatewayId, syncPromise);
return syncPromise;
};
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,
});
});
sendJson(agent, {
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;
}
if (req.method === "POST" && /^\/api\/gateways\/\d+\/sync$/.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 result = await syncGatewayBacklog(gatewayId);
jsonResponse(res, 200, { ok: true, ...result });
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") || "");
const agentInstanceId = String(url.searchParams.get("agentInstanceId") || "");
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.agentInstanceId = agentInstanceId || null;
ws.connectionId = randomUUID();
agents.set(gatewayId, ws);
reportGatewayPresence(gatewayId, {
status: "connected",
connectionId: ws.connectionId,
metadata: {
remote_address: req.socket.remoteAddress || null,
agent_instance_id: ws.agentInstanceId,
},
}).catch(() => {});
broadcastGatewayEvent(gatewayId, {
type: "presence.changed",
gatewayId,
status: "connected",
connectionId: ws.connectionId,
});
syncGatewayBacklog(gatewayId, ws).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;
const sessionRecord = {
ws,
session,
transcript: "",
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) {
sendJson(agent, {
type: "OPEN_ROOT_SHELL",
payload: {
sessionId: String(session.id),
reason: session.reason,
cols: session.cols ?? session.metadata?.cols ?? null,
rows: session.rows ?? session.metadata?.rows ?? null,
cwd: session.cwd ?? session.metadata?.cwd ?? null,
shellCommand: session.shell_command ?? session.metadata?.shell_command ?? null,
shellArgs: session.shell_args ?? session.metadata?.shell_args ?? [],
},
});
} else {
closeBrowserShellSocket(
sessionRecord,
"agent_offline",
"Gateway agent is not connected to the broker.",
1011
);
}
});
return;
}
if (url.pathname === "/ws/browser-gateway-stream") {
const token = String(url.searchParams.get("token") || "");
if (token === "") {
socket.destroy();
return;
}
const session = await validateBrowserStream({ token, headers: req.headers });
wss.handleUpgrade(req, socket, head, (ws) => {
ws.sessionToken = token;
ws.streamSessionInfo = session;
const sessionRecord = {
ws,
session,
subscriptions: new Set(parseScopes(session.scopes || ["overview", "tasks", "logs", "statistics"])),
};
registerGatewayStreamSession(sessionRecord);
sendJson(ws, {
type: "gateway.stream.ready",
gatewayId: String(session.gateway_id),
subscriptions: Array.from(sessionRecord.subscriptions.values()),
connected: Boolean(agents.get(String(session.gateway_id))?.readyState === 1),
});
wss.emit("connection", ws, req);
});
return;
}
} catch {
socket.destroy();
return;
}
socket.destroy();
});
wss.on("connection", (ws) => {
ws.on("message", async (raw) => {
let message;
try {
message = JSON.parse(raw.toString());
} catch {
return;
}
try {
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 (message.type === "TELEMETRY") {
const payload = message.payload || {};
let ingested = null;
let ingestError = null;
try {
ingested = await ingestTelemetry(String(ws.gatewayId), payload);
} catch (error) {
ingestError = error instanceof Error ? error.message : String(error);
}
const fallbackStatistics = {
system_metrics: payload?.metadata?.system_metrics || {},
container_health: payload?.metadata?.container_health || {},
};
broadcastGatewayEvent(String(ws.gatewayId), {
type: "gateway.telemetry",
gatewayId: String(ws.gatewayId),
telemetry: payload,
gateway: ingested?.gateway || ingested || null,
error: ingestError,
});
broadcastGatewayEvent(String(ws.gatewayId), {
type: "stats.updated",
gatewayId: String(ws.gatewayId),
statistics: ingested?.statistics || ingested || fallbackStatistics,
error: ingestError,
});
return;
}
if (message.type === "TASK_EVENT") {
const operationId = Number(message.operationId ?? message.payload?.operation_id ?? 0);
if (!Number.isFinite(operationId) || operationId <= 0) {
return;
}
const operation = await ingestTaskEvent(String(ws.gatewayId), operationId, message.payload || {});
broadcastGatewayEvent(String(ws.gatewayId), {
type: "task.updated",
gatewayId: String(ws.gatewayId),
operationId,
operation,
});
return;
}
if (message.type === "TASK_RESULT") {
const operationId = Number(message.operationId ?? message.payload?.operation_id ?? 0);
if (!Number.isFinite(operationId) || operationId <= 0) {
return;
}
const operation = await ingestTaskResult(String(ws.gatewayId), operationId, message.payload || {});
broadcastGatewayEvent(String(ws.gatewayId), {
type: "task.updated",
gatewayId: String(ws.gatewayId),
operationId,
operation,
});
await syncGatewayBacklog(String(ws.gatewayId), ws).catch(() => {});
return;
}
if (message.type === "LOG_FRAME") {
const logEntry = await ingestLogEntry(String(ws.gatewayId), message.payload || {});
broadcastGatewayEvent(String(ws.gatewayId), {
type: "log.append",
gatewayId: String(ws.gatewayId),
entry: logEntry,
});
return;
}
if (["SHELL_OUTPUT", "SHELL_OPENED", "SHELL_EXIT"].includes(message.type)) {
const sessionRecord = browserShellSessions.get(String(message.sessionId));
if (!sessionRecord) {
return;
}
if (message.type === "SHELL_OUTPUT") {
sessionRecord.transcript += String(message.data || "");
sendJson(sessionRecord.ws, { type: "output", data: String(message.data || "") });
return;
}
if (message.type === "SHELL_OPENED") {
await markShellSessionOpened(sessionRecord.ws.sessionToken, ws.connectionId || null).catch(() => {});
sendJson(sessionRecord.ws, { type: "opened" });
return;
}
if (message.type === "SHELL_EXIT") {
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(1000, "agent_exit");
}
}
}
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") {
sendJson(agent, {
type: "SHELL_INPUT",
payload: {
sessionId,
data: String(message.data || ""),
},
});
return;
}
if (message.type === "resize") {
sendJson(agent, {
type: "RESIZE_ROOT_SHELL",
payload: {
sessionId,
cols: Number(message.cols || 0),
rows: Number(message.rows || 0),
},
});
return;
}
if (message.type === "close") {
sendJson(agent, {
type: "CLOSE_ROOT_SHELL",
payload: { sessionId },
});
}
return;
}
if (ws.streamSessionInfo) {
const sessionRecord = browserStreamSessions.get(String(ws.streamSessionInfo.id));
if (!sessionRecord) {
return;
}
if (message.type === "SUBSCRIBE") {
for (const scope of parseScopes(message.scopes || message.subscriptions || [])) {
sessionRecord.subscriptions.add(scope);
}
sendJson(sessionRecord.ws, {
type: "subscribed",
subscriptions: Array.from(sessionRecord.subscriptions.values()),
});
return;
}
if (message.type === "UNSUBSCRIBE") {
for (const scope of parseScopes(message.scopes || message.subscriptions || [])) {
sessionRecord.subscriptions.delete(scope);
}
sendJson(sessionRecord.ws, {
type: "unsubscribed",
subscriptions: Array.from(sessionRecord.subscriptions.values()),
});
return;
}
if (message.type === "PING") {
sendJson(sessionRecord.ws, { type: "PONG" });
}
}
} catch {
// Ignore stale gateway/session delivery errors without killing the broker process.
}
});
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));
}
markGatewayShellSessionsClosed(String(ws.gatewayId), "agent_disconnected");
broadcastGatewayEvent(String(ws.gatewayId), {
type: "presence.changed",
gatewayId: String(ws.gatewayId),
status: "disconnected",
reason: closeReason || "agent_disconnected",
});
reportGatewayPresence(String(ws.gatewayId), {
status: "disconnected",
connectionId: ws.connectionId || null,
reason: closeReason || "agent_disconnected",
metadata: {
agent_instance_id: ws.agentInstanceId || null,
},
}).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) {
sendJson(agent, {
type: "CLOSE_ROOT_SHELL",
payload: { sessionId },
});
}
const sessionRecord = browserShellSessions.get(sessionId);
if (sessionRecord) {
await closeBrowserShellSession(sessionRecord, sessionRecord.closedReason || "browser_closed");
browserShellSessions.delete(sessionId);
}
return;
}
if (ws.streamSessionInfo) {
const sessionRecord = browserStreamSessions.get(String(ws.streamSessionInfo.id));
if (sessionRecord) {
removeGatewayStreamSession(sessionRecord);
}
}
});
});
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 browserShellSessions.values()) {
session.ws.terminate();
}
for (const session of browserStreamSessions.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,
browserShellSessions,
browserStreamSessions,
gatewayStreamSessions,
pendingCommands,
managerUrl,
authMode,
},
};
}
async function runBrokerFromCli() {
const broker = createBrokerServer();
let shuttingDown = false;
const shutdown = async (signal) => {
if (shuttingDown) {
return;
}
shuttingDown = true;
try {
await broker.close();
process.exit(0);
} catch (error) {
console.error(`Failed to shut down broker after ${signal}:`, error);
process.exit(1);
}
};
process.on("SIGINT", () => {
void shutdown("SIGINT");
});
process.on("SIGTERM", () => {
void shutdown("SIGTERM");
});
const requestedPort = Number(process.env.PORT || 4300);
const address = await broker.listen(Number.isFinite(requestedPort) ? requestedPort : 4300);
const normalizedPort =
typeof address === "object" && address !== null && "port" in address
? address.port
: requestedPort;
console.log(`TruckWash edge broker listening on ${normalizedPort}`);
}
if (process.argv[1] && fileURLToPath(import.meta.url) === process.argv[1]) {
runBrokerFromCli().catch((error) => {
console.error("TruckWash edge broker failed to start:", error);
process.exit(1);
});
}