import { createServer } from "node:http"; import { randomUUID } from "node:crypto"; import { DEV_ENDPOINT, ENVIRONMENT_DEV, ERROR_CODES } from "../protocol/index.mjs"; import { AUDIT_FIELD_NAMES, CLOUD_API_SERVICE_ID, createAuditRecord } from "../audit/index.mjs"; import { SUPPORTED_RPC_METHODS, createErrorEnvelope, handleJsonRpcRequest } from "./json-rpc.mjs"; import { describeCodeAgentAvailability, handleCodeAgentChat } from "./code-agent-chat.mjs"; import { applyRuntimeDbReadinessLayers, buildDbRuntimeReadiness } from "./db-contract.mjs"; import { buildCloudApiReadiness } from "./health-contract.mjs"; import { M3_IO_CONTROL_ROUTE, describeM3IoControl, describeM3IoControlLive, handleM3IoControl } from "./m3-io-control.mjs"; import { createConfiguredCloudRuntimeStore } from "../db/runtime-store.mjs"; const DEFAULT_BODY_LIMIT_BYTES = 1024 * 1024; export function createCloudApiServer(options = {}) { const env = options.env ?? process.env; const runtimeStore = options.runtimeStore || createConfiguredCloudRuntimeStore({ ...options, env }); return createServer(async (request, response) => { try { await routeRequest(request, response, { ...options, env, runtimeStore }); } catch (error) { sendJson(response, 500, { error: { code: "internal_error", message: "hwlab-cloud-api request handling failed", reason: error.message } }); } }); } export async function buildHealthPayload(options = {}) { const serviceId = CLOUD_API_SERVICE_ID; const env = options.env ?? process.env; const commitId = process.env.HWLAB_COMMIT_ID || process.env.HWLAB_GIT_SHA || "unknown"; const imageReference = process.env.HWLAB_IMAGE || "unknown"; const imageTag = process.env.HWLAB_IMAGE_TAG || commitId.slice(0, 7) || "unknown"; const dbProbe = await buildDbRuntimeReadiness(env, options.dbProbe); const codeAgent = describeCodeAgentAvailability(env, options); const runtime = await runtimeReadiness(options.runtimeStore ?? createConfiguredCloudRuntimeStore({ ...options, env })); const db = applyRuntimeDbReadinessLayers(dbProbe, runtime); const readiness = buildCloudApiReadiness({ db, codeAgent, runtime }); return { serviceId, environment: ENVIRONMENT_DEV, status: readiness.status, ready: readiness.ready, readiness, blockers: readiness.blockers, blockerCodes: readiness.blockerCodes, service: { id: serviceId, role: "cloud-api", healthPath: "/health", livePath: "/health/live" }, commit: { id: commitId, source: process.env.HWLAB_BUILD_SOURCE || "runtime-env" }, image: { reference: imageReference, tag: imageTag, digest: process.env.HWLAB_IMAGE_DIGEST || "unknown" }, endpoint: process.env.HWLAB_PUBLIC_ENDPOINT || DEV_ENDPOINT, observedAt: new Date().toISOString(), db, codeAgent, runtime }; } async function routeRequest(request, response, options) { const url = new URL(request.url || "/", "http://hwlab-cloud-api.local"); if (request.method === "GET" && (url.pathname === "/health" || url.pathname === "/health/live")) { sendJson(response, 200, await buildHealthPayload(options)); return; } if (request.method === "GET" && url.pathname === "/live") { sendJson(response, 200, { serviceId: CLOUD_API_SERVICE_ID, status: "live" }); return; } if (request.method === "POST" && (url.pathname === "/rpc" || url.pathname === "/json-rpc")) { await handleRpcHttpRequest(request, response, options); return; } if (url.pathname === "/v1" || url.pathname.startsWith("/v1/")) { await handleRestAdapter(request, response, url, options); return; } sendRestError(request, response, 404, "not_found", "Route is not part of the L1 cloud-api runtime", { operation: `${request.method || "GET"} ${url.pathname}`, target: { type: "route", id: url.pathname } }); } async function handleRpcHttpRequest(request, response, options) { const body = await readBody(request, options.bodyLimitBytes); let envelope; try { envelope = body ? JSON.parse(body) : {}; } catch (error) { sendJson( response, 400, createErrorEnvelope({ id: "req_unassigned", code: ERROR_CODES.parseError, message: "Invalid JSON body", data: { reason: error.message }, context: { operation: "json_rpc.parse", target: { type: "http_route", id: request.url || "/rpc" }, result: "rejected" } }) ); return; } const rpcResponse = await handleJsonRpcRequest(envelope, { adapter: "json-rpc", runtimeStore: options.runtimeStore, env: options.env, dbProbe: options.dbProbe, m3IoRequestJson: options.m3IoRequestJson }); sendJson(response, statusForRpcResponse(rpcResponse), rpcResponse); } async function handleRestAdapter(request, response, url, options) { if (request.method === "GET" && url.pathname === "/v1") { const dbProbe = await buildDbRuntimeReadiness(options.env ?? process.env, options.dbProbe); const codeAgent = describeCodeAgentAvailability(options.env ?? process.env, options); const runtime = await runtimeReadiness(options.runtimeStore); const db = applyRuntimeDbReadinessLayers(dbProbe, runtime); const readiness = buildCloudApiReadiness({ db, codeAgent, runtime }); sendJson(response, 200, { serviceId: CLOUD_API_SERVICE_ID, adapter: "rest", status: readiness.status, ready: readiness.ready, rpcBridge: "POST /v1/rpc/{method}", codeAgent, methods: SUPPORTED_RPC_METHODS, auditFields: AUDIT_FIELD_NAMES, db, runtime, readiness, blockers: readiness.blockers, blockerCodes: readiness.blockerCodes, m3IoControl: describeM3IoControl(options) }); return; } if (request.method === "GET" && url.pathname === M3_IO_CONTROL_ROUTE) { sendJson(response, 200, await describeM3IoControlLive(options)); return; } if (request.method === "POST" && url.pathname === M3_IO_CONTROL_ROUTE) { await handleM3IoControlHttp(request, response, options); return; } if (request.method === "POST" && url.pathname === "/v1/agent/chat") { await handleCodeAgentChatHttp(request, response, options); return; } if (request.method !== "POST" || !url.pathname.startsWith("/v1/rpc/")) { sendRestError(request, response, 404, "not_found", "REST route is not implemented in the L1 cloud-api runtime", { operation: `${request.method || "GET"} ${url.pathname}`, target: { type: "route", id: url.pathname } }); return; } const method = decodeURIComponent(url.pathname.slice("/v1/rpc/".length)); const body = await readBody(request, options.bodyLimitBytes); let params = {}; try { params = body ? JSON.parse(body) : {}; } catch (error) { sendRestError(request, response, 400, "parse_error", "Invalid JSON body", { operation: method, target: { type: "rpc_method", id: method }, result: "rejected", reason: error.message }); return; } if (!params || typeof params !== "object" || Array.isArray(params)) { sendRestError(request, response, 400, "invalid_params", "REST RPC bridge body must be a JSON object", { operation: method, target: { type: "rpc_method", id: method }, result: "rejected" }); return; } const rpcResponse = await handleJsonRpcRequest( { jsonrpc: "2.0", id: getHeader(request, "x-request-id") || `req_${randomUUID()}`, method, params, meta: { traceId: getHeader(request, "x-trace-id") || `trc_${randomUUID()}`, actorId: getHeader(request, "x-actor-id"), serviceId: getHeader(request, "x-source-service-id") || CLOUD_API_SERVICE_ID, environment: ENVIRONMENT_DEV } }, { adapter: "rest", runtimeStore: options.runtimeStore, env: options.env, dbProbe: options.dbProbe, m3IoRequestJson: options.m3IoRequestJson } ); sendJson(response, statusForRpcResponse(rpcResponse), rpcResponse); } async function handleCodeAgentChatHttp(request, response, options) { const body = await readBody(request, options.bodyLimitBytes); let params = {}; try { params = body ? JSON.parse(body) : {}; } catch (error) { sendJson(response, 400, { conversationId: "cnv_unassigned", sessionId: "cnv_unassigned", messageId: "msg_unassigned", status: "failed", createdAt: new Date().toISOString(), updatedAt: new Date().toISOString(), traceId: getHeader(request, "x-trace-id") || "trc_unassigned", provider: "unassigned", model: process.env.HWLAB_CODE_AGENT_MODEL || process.env.OPENAI_MODEL || "unknown", backend: "hwlab-cloud-api/code-agent-chat", error: { code: "parse_error", message: "Invalid JSON body", reason: error.message } }); return; } if (!params || typeof params !== "object" || Array.isArray(params)) { sendJson(response, 400, { conversationId: "cnv_unassigned", sessionId: "cnv_unassigned", messageId: "msg_unassigned", status: "failed", createdAt: new Date().toISOString(), updatedAt: new Date().toISOString(), traceId: getHeader(request, "x-trace-id") || "trc_unassigned", provider: "unassigned", model: process.env.HWLAB_CODE_AGENT_MODEL || process.env.OPENAI_MODEL || "unknown", backend: "hwlab-cloud-api/code-agent-chat", error: { code: "invalid_params", message: "Code Agent chat body must be a JSON object" } }); return; } const payload = await handleCodeAgentChat( { ...params, traceId: getHeader(request, "x-trace-id") || params.traceId }, { callProvider: options.callCodeAgentProvider, timeoutMs: options.codeAgentTimeoutMs, env: options.env, workspace: options.workspace, skillsDirs: options.skillsDirs, skillsDirsExact: options.skillsDirsExact } ); sendJson(response, payload.status === "failed" && payload.error?.code === "invalid_params" ? 400 : 200, payload); } async function handleM3IoControlHttp(request, response, options) { const body = await readBody(request, options.bodyLimitBytes); let params = {}; try { params = body ? JSON.parse(body) : {}; } catch (error) { sendJson(response, 400, { serviceId: CLOUD_API_SERVICE_ID, contractVersion: "m3-io-control-v1", status: "blocked", accepted: false, traceId: getHeader(request, "x-trace-id") || "trc_unassigned", error: { code: "parse_error", message: "Invalid JSON body", reason: error.message } }); return; } if (!params || typeof params !== "object" || Array.isArray(params)) { sendJson(response, 400, { serviceId: CLOUD_API_SERVICE_ID, contractVersion: "m3-io-control-v1", status: "blocked", accepted: false, traceId: getHeader(request, "x-trace-id") || "trc_unassigned", error: { code: "invalid_params", message: "M3 IO control body must be a JSON object" } }); return; } try { const payload = await handleM3IoControl( { ...params, traceId: getHeader(request, "x-trace-id") || params.traceId, requestId: getHeader(request, "x-request-id") || params.requestId, actorId: getHeader(request, "x-actor-id") || params.actorId }, { runtimeStore: options.runtimeStore, env: options.env, requestJson: options.m3IoRequestJson, traceId: getHeader(request, "x-trace-id"), requestId: getHeader(request, "x-request-id"), actorId: getHeader(request, "x-actor-id") } ); sendJson(response, payload.httpStatus ?? 200, payload); } catch (error) { const protocolCode = Number.isInteger(error?.code) ? error.code : ERROR_CODES.internalError; sendJson(response, protocolCode === ERROR_CODES.invalidParams ? 400 : 200, { serviceId: CLOUD_API_SERVICE_ID, contractVersion: "m3-io-control-v1", status: "blocked", accepted: false, traceId: getHeader(request, "x-trace-id") || params.traceId || "trc_unassigned", error: { code: protocolCode, message: error.message, data: error.data ?? {} } }); } } function sendRestError(request, response, statusCode, code, message, context = {}) { const requestId = getHeader(request, "x-request-id") || "req_unassigned"; sendJson(response, statusCode, { error: { code, message, reason: context.reason, audit: createAuditRecord({ requestId, actor: { type: "user", id: getHeader(request, "x-actor-id") || `system_${CLOUD_API_SERVICE_ID}` }, source: { serviceId: CLOUD_API_SERVICE_ID, environment: ENVIRONMENT_DEV, traceId: getHeader(request, "x-trace-id"), adapter: "rest" }, operation: context.operation || `${request.method || "GET"} ${request.url || "/"}`, target: context.target || { type: "route", id: request.url || "/" }, result: context.result || "failed" }) } }); } function statusForRpcResponse(rpcResponse) { if (!rpcResponse.error) { return 200; } if (rpcResponse.error.code === ERROR_CODES.invalidRequest || rpcResponse.error.code === ERROR_CODES.parseError) { return 400; } if (rpcResponse.error.code === ERROR_CODES.methodNotFound) { return 404; } return 200; } function readBody(request, limitBytes = DEFAULT_BODY_LIMIT_BYTES) { const limit = Number.isInteger(limitBytes) && limitBytes > 0 ? limitBytes : DEFAULT_BODY_LIMIT_BYTES; return new Promise((resolve, reject) => { let body = ""; request.setEncoding("utf8"); request.on("data", (chunk) => { body += chunk; if (Buffer.byteLength(body, "utf8") > limit) { request.destroy(new Error(`request body exceeds ${limit} bytes`)); } }); request.on("end", () => resolve(body)); request.on("error", reject); }); } function sendJson(response, statusCode, body) { response.writeHead(statusCode, { "content-type": "application/json; charset=utf-8", "cache-control": "no-store" }); response.end(`${JSON.stringify(body)}\n`); } function getHeader(request, name) { const value = request.headers[name.toLowerCase()]; if (Array.isArray(value)) { return value[0]; } return value; } async function runtimeReadiness(runtimeStore) { if (typeof runtimeStore?.readiness === "function") { return runtimeStore.readiness(); } return runtimeStore.summary(); }