import { createHash, timingSafeEqual } from "node:crypto"; import type { ServerWebSocket } from "bun"; import type { Broker } from "./broker"; const PROTOCOL_VERSION = 3; /** Maximum UTF-8 byte length accepted for a single broker JSON frame. */ const MAX_BROKER_JSON_FRAME_BYTES = 4 * 1024 * 1024; const BROKER_OPERATIONS = new Set([ "session.list", "session.get_endpoint", "session.create", "session.fork", "session.resume", "session.close", "session.delete", "broker.shutdown", "session.control", "broker.lookup_lifecycle", "model.resolve", ]); type RequestInput = Record; type BrokerRequest = { type: "broker_request"; id?: string; operation?: string; input?: RequestInput; idempotencyKey?: string; }; function digest(value: string): Buffer { return createHash("sha256").update(value).digest(); } function tokenMatches(expected: string, actual: string | null): boolean { return actual !== null && timingSafeEqual(digest(expected), digest(actual)); } function isInput(value: unknown): value is RequestInput { return typeof value === "object" && value !== null && !Array.isArray(value); } function send(socket: ServerWebSocket, frame: Record): number { return socket.send(JSON.stringify(frame)); } export type BrokerShutdownSendAction = "wait_for_drain" | "close"; export function brokerShutdownSendAction(status: number): BrokerShutdownSendAction { return status < 0 ? "wait_for_drain" : "close"; } function sendError(socket: ServerWebSocket, id: string | undefined, code: string, message: string): void { send(socket, { type: "broker_response", ...(id === undefined ? {} : { id }), ok: false, error: { code, message } }); } /** Loopback-only WebSocket transport for the agent-global SDK broker. */ export class BrokerTransport { readonly #broker: Broker; readonly #token: string; readonly #requestedPort: number; #server: Bun.Server | null = null; #port = 0; readonly #shutdownBackpressured = new WeakSet>(); #shutdownRequested = false; constructor(broker: Broker, token: string, port = 0) { this.#broker = broker; this.#token = token; this.#requestedPort = port; } get port(): number { if (!this.#server) throw new Error("Broker transport is not running"); return this.#port; } async start(): Promise { if (this.#server) return this.#port; this.#server = Bun.serve({ hostname: "127.0.0.1", port: this.#requestedPort, fetch: request => { const url = new URL(request.url); if (request.headers.get("upgrade")?.toLowerCase() !== "websocket") return new Response("Upgrade Required", { status: 426 }); if (!tokenMatches(this.#token, url.searchParams.get("token"))) return new Response("Unauthorized", { status: 401 }); if (!this.#server?.upgrade(request, { data: undefined })) return new Response("WebSocket upgrade failed", { status: 400 }); return undefined; }, websocket: { maxPayloadLength: MAX_BROKER_JSON_FRAME_BYTES * 2, open: socket => { send(socket, { type: "broker_hello", protocolVersion: PROTOCOL_VERSION }); }, message: (socket, message) => void this.#handleMessage(socket, message), drain: socket => { if (this.#shutdownBackpressured.delete(socket)) this.#scheduleStopAfterShutdownResponse(); }, }, }); this.#port = this.#server.port ?? 0; return this.#port; } async stop(): Promise { const server = this.#server; this.#server = null; if (server) await server.stop(true); } async #handleMessage(socket: ServerWebSocket, raw: string | Buffer): Promise { if (this.#shutdownRequested) { sendError(socket, undefined, "unavailable", "broker is shutting down"); return; } if (Buffer.byteLength(raw) > MAX_BROKER_JSON_FRAME_BYTES) { sendError(socket, undefined, "payload_too_large", "broker JSON frame exceeds 4 MiB limit"); return; } let frame: BrokerRequest; try { frame = JSON.parse(raw.toString()) as BrokerRequest; } catch { sendError(socket, undefined, "invalid_input", "malformed JSON"); return; } if (frame?.type !== "broker_request") { sendError(socket, undefined, "invalid_input", "invalid broker frame"); return; } if (typeof frame.id !== "string" || !frame.id) { sendError(socket, frame.id, "invalid_input", "request id is required"); return; } if (typeof frame.operation !== "string" || !BROKER_OPERATIONS.has(frame.operation)) { sendError(socket, frame.id, "unknown_operation", "unknown broker operation"); return; } if (frame.input === undefined || !isInput(frame.input)) { sendError(socket, frame.id, "invalid_input", "request input must be an object"); return; } if (frame.idempotencyKey !== undefined && typeof frame.idempotencyKey !== "string") { sendError(socket, frame.id, "invalid_input", "idempotencyKey must be a string"); return; } if (frame.operation === "broker.shutdown") { const action = brokerShutdownSendAction( send(socket, { type: "broker_response", id: frame.id, ok: true, result: { accepted: true } }), ); this.#shutdownRequested = true; if (action === "wait_for_drain") this.#shutdownBackpressured.add(socket); else this.#scheduleStopAfterShutdownResponse(); return; } try { const result = await this.#broker.handleRequest(frame.operation, frame.input, frame.idempotencyKey); send(socket, { type: "broker_response", id: frame.id, ...result }); } catch { sendError(socket, frame.id, "unavailable", "broker request failed"); } } #scheduleStopAfterShutdownResponse(): void { setTimeout(() => void this.#broker.stop(), 25); } }