import { createCloudApiServer } from "./server.ts"; import { createHwpodNodeWsRegistry } from "./hwpod-node-ws-registry.ts"; import { drainWorkbenchRealtimeConnections } from "./server-workbench-http.ts"; import { classifyRuntimeDbError } from "../db/runtime-store.ts"; import { createConfiguredCloudRuntimeStore } from "../db/runtime-store.ts"; const POSTGRES_TRANSIENT_UNHANDLED_REJECTION_BOUNDARY = Symbol.for("hwlab.cloud.postgresTransientUnhandledRejectionBoundaryInstalled"); export async function createCloudApiBunServer(options: any = {}) { const env = options.env ?? process.env; const logger = options.logger ?? console; installPostgresTransientUnhandledRejectionBoundary({ env, logger }); const runtimeStore = options.runtimeStore || createConfiguredCloudRuntimeStore({ ...options, env }); const hwpodNodeWsRegistry = options.hwpodNodeWsRegistry || createHwpodNodeWsRegistry({ ...options, env, runtimeStore }); const innerServer = createCloudApiServer({ ...options, env, runtimeStore, hwpodNodeWsRegistry }); try { await innerServer.hwlabStartupReady; } catch (error) { await innerServer.hwlabAbortStartup?.(); throw error; } await new Promise((resolve) => innerServer.listen(0, "127.0.0.1", resolve)); const innerAddress = innerServer.address(); const innerPort = typeof innerAddress === "object" && innerAddress ? innerAddress.port : 0; if (!innerPort) throw new Error("cloud-api inner HTTP server did not expose a port"); const host = options.host ?? "0.0.0.0"; const port = options.port ?? 0; const idleTimeout = parsePositiveInteger(options.idleTimeout ?? env.HWLAB_CLOUD_API_IDLE_TIMEOUT_SECONDS, 120); const server = (options.bunServe ?? Bun.serve)({ hostname: host, port, idleTimeout, async fetch(request, bunServer) { const url = new URL(request.url); if (url.pathname === "/v1/hwpod-node/ws") { const expectedToken = String(env.HWLAB_HWPOD_NODE_WS_TOKEN ?? "").trim(); const token = String(url.searchParams.get("token") ?? request.headers.get("x-hwpod-node-token") ?? ""); if (expectedToken && token !== expectedToken) return json({ ok: false, error: { code: "hwpod_node_ws_token_invalid" } }, 401); const upgraded = bunServer.upgrade(request, { data: { hwpodNodeWsRegistry } }); return upgraded ? undefined : json({ ok: false, status: "upgrade_required", route: "/v1/hwpod-node/ws" }, 426); } return proxyHttpToInnerServer(request, innerPort); }, websocket: { open(socket: any) { socket.data.hwpodNodeWsRegistry.openBunSocket(socket); }, message(socket: any, message: string | Buffer) { socket.data.hwpodNodeWsRegistry.handleBunMessage(socket, message); }, close(socket: any) { socket.data.hwpodNodeWsRegistry.closeBunSocket(socket); } } }); let removeSignalHandlers = () => {}; let stopping = false; const stop = async (reason = "manual", signal: string | null = null) => { if (stopping) return { ok: true, alreadyStopping: true }; stopping = true; removeSignalHandlers(); const sseDrainTimeoutMs = parsePositiveInteger(options.sseDrainTimeoutMs ?? env.HWLAB_WORKBENCH_SSE_DRAIN_TIMEOUT_MS, 2500); const innerCloseTimeoutMs = parsePositiveInteger(options.innerCloseTimeoutMs ?? env.HWLAB_CLOUD_API_INNER_CLOSE_TIMEOUT_MS, 8000); logger?.info?.({ event: "cloud_api_shutdown_started", reason, signal, sseDrainTimeoutMs, innerCloseTimeoutMs, valuesRedacted: true }); const sseDrain = await drainWorkbenchRealtimeConnections({ reason, signal, timeoutMs: sseDrainTimeoutMs, env }); try { server.stop(false); } catch (error) { logger?.warn?.({ event: "cloud_api_bun_stop_failed", reason, signal, errorName: error?.name ?? "Error", message: error instanceof Error ? error.message : String(error ?? "unknown"), valuesRedacted: true }); } const innerClose = await closeNodeHttpServer(innerServer, innerCloseTimeoutMs); logger?.info?.({ event: "cloud_api_shutdown_completed", reason, signal, sseDrain, innerClose, valuesRedacted: true }); return { ok: sseDrain.ok === true && innerClose.ok === true, sseDrain, innerClose }; }; if (options.installSignalHandlers !== false) { removeSignalHandlers = installCloudApiSignalHandlers({ stop, logger }); } return { server, innerServer, hwpodNodeWsRegistry, port: server.port, url: `http://${host === "0.0.0.0" ? "127.0.0.1" : host}:${server.port}`, stop }; } function installCloudApiSignalHandlers({ stop, logger }: { stop: (reason?: string, signal?: string | null) => Promise, logger?: any }) { const handlers: Array<{ signal: string, handler: () => void }> = []; for (const signal of ["SIGINT", "SIGTERM"]) { const handler = () => { void stop("process_signal", signal).catch((error) => { logger?.error?.({ event: "cloud_api_shutdown_failed", signal, errorName: error?.name ?? "Error", message: error instanceof Error ? error.message : String(error ?? "unknown"), valuesRedacted: true }); }).finally(() => { process.exit(0); }); }; process.once(signal, handler); handlers.push({ signal, handler }); } return () => { for (const { signal, handler } of handlers.splice(0)) { process.off(signal, handler); } }; } function closeNodeHttpServer(server: any, timeoutMs: number): Promise<{ ok: boolean, timeout: boolean, errorName: string | null, message: string | null }> { return new Promise((resolve) => { let settled = false; const finish = (ok: boolean, timeout = false, error: any = null) => { if (settled) return; settled = true; if (timeoutHandle) clearTimeout(timeoutHandle); resolve({ ok, timeout, errorName: error?.name ?? null, message: error instanceof Error ? error.message : null }); }; const timeoutHandle = setTimeout(() => { try { server.closeAllConnections?.(); } catch {} finish(false, true, null); }, timeoutMs); timeoutHandle.unref?.(); try { server.close((error: any) => { finish(!error, false, error); }); server.closeIdleConnections?.(); } catch (error) { finish(false, false, error); } }); } export function installPostgresTransientUnhandledRejectionBoundary(options: any = {}) { const proc: any = options.process ?? process; if (!proc || typeof proc.on !== "function") return { ok: false, installed: false, reason: "process-unavailable" }; if (proc[POSTGRES_TRANSIENT_UNHANDLED_REJECTION_BOUNDARY]) return { ok: true, installed: false, reason: "already-installed" }; proc[POSTGRES_TRANSIENT_UNHANDLED_REJECTION_BOUNDARY] = true; const logger = options.logger ?? console; proc.on("unhandledRejection", (reason: any) => { const classified = classifyRuntimeDbError(reason); if (classified.retryable === true && classified.transient === true && classified.connection?.queryResult === "connect_timeout") { logger?.warn?.({ event: "postgres_runtime_unhandled_rejection_contained", blocker: classified.blocker, queryResult: classified.connection?.queryResult, errorCode: classified.connection?.errorCode ?? "UNKNOWN", retryable: true, transient: true, retryAfterMs: classified.retryAfterMs ?? null, valuesRedacted: true, endpointRedacted: true }); return; } queueMicrotask(() => { throw reason; }); }); return { ok: true, installed: true }; } async function proxyHttpToInnerServer(request: Request, innerPort: number) { const sourceUrl = new URL(request.url); const targetUrl = new URL(sourceUrl.pathname + sourceUrl.search, `http://127.0.0.1:${innerPort}`); const method = request.method.toUpperCase(); return fetch(targetUrl, { method, headers: request.headers, body: method === "GET" || method === "HEAD" ? undefined : request.body, redirect: "manual" }); } function parsePositiveInteger(value: unknown, fallback: number) { const parsed = Number(value); if (!Number.isFinite(parsed) || parsed <= 0) return fallback; return Math.floor(parsed); } function json(body: any, status = 200) { return new Response(`${JSON.stringify(body)}\n`, { status, headers: { "content-type": "application/json; charset=utf-8", "cache-control": "no-store" } }); }