import { once } from "node:events"; import { createServer, type IncomingMessage, type ServerResponse, } from "node:http"; import type { AddressInfo, Socket } from "node:net"; import type { EvalLanguage } from "../tool/types.ts"; import { type BridgeError, generateBridgeToken, verifyBridgeToken, } from "./protocol.ts"; const DEFAULT_BODY_LIMIT_BYTES = 1024 * 1024; const LOOPBACK_HOST = "127.0.0.1"; export interface BridgeHttpCallRequest { args: unknown; callId: string; /** * Language of the calling prelude (py/rb). Lets the host attribute a * bridge call to that language's current active cell (the kernels serialize * per language, so one active cell per language is the invariant the * detached manager enforces). Absent for legacy callers, which then get * plain execution with no enrichment. */ language?: EvalLanguage; signal: AbortSignal; toolName: string; } export interface BridgeHttpCompletionRequest { opts?: unknown; prompt: string; signal: AbortSignal; } export interface BridgeServerOptions { bodyLimitBytes?: number; onCall: (request: BridgeHttpCallRequest) => Promise; onCompletion: (request: BridgeHttpCompletionRequest) => Promise; token?: string; } export interface BridgeServerHandle { close: () => Promise; port: number; token: string; } type JsonReply = | { ok: true; value: unknown } | { ok: false; error: BridgeError }; type Route = "/call" | "/completion"; export async function startBridgeServer( options: BridgeServerOptions ): Promise { const token = options.token ?? generateBridgeToken(); const sockets = new Set(); let closing: Promise | undefined; const server = createServer((request, response) => { void handleRequest(request, response, token, options).catch(() => { if (!response.writableFinished) { response.destroy(); } }); }); server.on("connection", (socket) => { sockets.add(socket); socket.on("close", () => sockets.delete(socket)); }); server.listen(0, LOOPBACK_HOST); await once(server, "listening"); const address = server.address(); if (!address || typeof address === "string") { throw new Error("Bridge server did not bind to a TCP port"); } return { port: (address as AddressInfo).port, token, close: async () => { closing ??= closeServer(server, sockets); await closing; }, }; } async function handleRequest( request: IncomingMessage, response: ServerResponse, token: string, options: BridgeServerOptions ): Promise { const abortController = new AbortController(); // IncomingMessage "close" fires on normal message completion in Node >= 16, // so premature disconnect must be detected on the response side instead: // its "close" without a finished response means the connection died mid-call. response.on("close", () => { if (!response.writableFinished) { abortController.abort(); } }); if (request.method !== "POST") { sendJson(response, 404, { ok: false, error: transportError("not_found", "Bridge route was not found"), }); return; } const route = routeFromUrl(request.url ?? ""); if (!route) { sendJson(response, 404, { ok: false, error: transportError("not_found", "Bridge route was not found"), }); return; } const auth = parseBearerToken(request.headers.authorization); if (!(auth && verifyBridgeToken(token, auth).ok)) { sendJson(response, 401, { ok: false, error: transportError("unauthorized", "Bridge authorization failed"), }); return; } const parsed = await readJsonBody( request, options.bodyLimitBytes ?? DEFAULT_BODY_LIMIT_BYTES ); if (!parsed.ok) { sendJson(response, parsed.status, { ok: false, error: transportError(parsed.code, parsed.message), }); return; } const reply = route === "/call" ? await dispatchCall(parsed.value, options, abortController.signal) : await dispatchCompletion(parsed.value, options, abortController.signal); sendJson(response, 200, reply); } async function dispatchCall( body: unknown, options: BridgeServerOptions, signal: AbortSignal ): Promise { if ( !isRecord(body) || typeof body.callId !== "string" || typeof body.toolName !== "string" || !("args" in body) ) { return { ok: false, error: transportError( "invalid_request", "Bridge call request was invalid" ), }; } try { const language = isEvalLanguage(body.language) ? body.language : undefined; return { ok: true, value: await options.onCall({ callId: body.callId, toolName: body.toolName, args: body.args, ...(language === undefined ? {} : { language }), signal, }), }; } catch (error) { return { ok: false, error: bridgeError(error) }; } } async function dispatchCompletion( body: unknown, options: BridgeServerOptions, signal: AbortSignal ): Promise { if (!isRecord(body) || typeof body.prompt !== "string") { return { ok: false, error: transportError( "invalid_request", "Bridge completion request was invalid" ), }; } try { return { ok: true, value: await options.onCompletion({ prompt: body.prompt, opts: body.opts, signal, }), }; } catch (error) { return { ok: false, error: bridgeError(error) }; } } async function readJsonBody( request: IncomingMessage, limit: number ): Promise< | { ok: true; value: unknown } | { ok: false; status: number; code: string; message: string } > { let raw = ""; try { for await (const chunk of request) { raw += String(chunk); if (Buffer.byteLength(raw, "utf8") > limit) { return { ok: false, status: 413, code: "body_too_large", message: `Bridge request body exceeds ${limit} bytes`, }; } } } catch { // The client socket died mid-body (ERR_STREAM_PREMATURE_CLOSE / aborted / // ECONNRESET). Do not rethrow: the reply goes nowhere (the response is // destroyed) but the rejection must not escape handleRequest. return { ok: false, status: 400, code: "disconnect", message: "Bridge request stream closed before completion", }; } try { return { ok: true, value: JSON.parse(raw) as unknown }; } catch { return { ok: false, status: 400, code: "invalid_json", message: "Bridge request body was not valid JSON", }; } } function routeFromUrl(rawUrl: string): Route | undefined { const path = new URL(rawUrl, "http://127.0.0.1").pathname; if (path === "/call" || path === "/completion") { return path; } } function parseBearerToken(header: string | undefined): string | undefined { const prefix = "Bearer "; if (!header?.startsWith(prefix)) { return; } return header.slice(prefix.length); } function sendJson( response: ServerResponse, status: number, body: JsonReply ): void { if (response.destroyed) { return; } try { response.writeHead(status, { "content-type": "application/json; charset=utf-8", }); response.end(JSON.stringify(body)); } catch { // The socket can die between the destroyed check and the write; the // connection-level catch destroys the response in that case. } } function bridgeError(error: unknown): BridgeError { if (error instanceof Error) { return { name: error.name, message: error.message, stack: error.stack }; } return { message: String(error) }; } function transportError(code: string, message: string): BridgeError { return { code, message }; } function isRecord(value: unknown): value is Record { return typeof value === "object" && value !== null; } function isEvalLanguage(value: unknown): value is EvalLanguage { return value === "js" || value === "py" || value === "rb" || value === "ts"; } async function closeServer( server: ReturnType, sockets: Set ): Promise { server.closeAllConnections?.(); for (const socket of sockets) { socket.destroy(); } await new Promise((resolve, reject) => { server.close((error) => { if (error && "code" in error && error.code !== "ERR_SERVER_NOT_RUNNING") { reject(error); } else { resolve(); } }); }); }