Files
pikasTech-HWLAB/internal/cloud/server.ts
T
2026-06-05 01:35:35 +08:00

935 lines
32 KiB
TypeScript

import { createServer } from "node:http";
import { randomUUID } from "node:crypto";
import { existsSync, lstatSync, mkdirSync, symlinkSync } from "node:fs";
import path from "node:path";
import {
buildMetadataFromEnv,
normalizeIsoTimestamp
} from "../build-metadata.mjs";
import { DEV_ENDPOINT, ENVIRONMENT_DEV, ERROR_CODES, SERVICE_IDS } from "../protocol/index.mjs";
import {
AUDIT_FIELD_NAMES,
CLOUD_API_SERVICE_ID,
createAuditRecord
} from "../audit/index.mjs";
import {
SUPPORTED_RPC_METHODS,
createErrorEnvelope,
createResultEnvelope,
handleJsonRpcRequest
} from "./json-rpc.ts";
import {
describeCodeAgentAvailability,
handleCodeAgentChat
} from "./code-agent-chat.ts";
import { defaultCodeAgentTraceStore } from "./code-agent-trace-store.ts";
import { createCodeAgentSessionRegistry } from "./code-agent-session-registry.ts";
import { codeAgentSessionLifecycleSummary } from "./code-agent-session-lifecycle.ts";
import { createCodexStdioSessionManager } from "./codex-stdio-session.ts";
import { buildGateDiagnosticsRows } from "./gate-diagnostics.ts";
import {
applyRuntimeDbReadinessLayers,
buildDbRuntimeReadiness
} from "./db-contract.ts";
import { buildCloudApiReadiness } from "./health-contract.ts";
import {
M3_IO_CONTROL_ROUTE,
M3_STATUS_ROUTE,
describeM3IoControl,
describeM3IoControlLive,
describeM3StatusLive,
handleM3IoControl
} from "./m3-io-control.ts";
import { createConfiguredCloudRuntimeStore } from "../db/runtime-store.ts";
import { createAccessController } from "./access-control.ts";
import { createGatewayDemoRegistry } from "./gateway-demo-registry.ts";
import { buildLiveBuildsPayload } from "./server-rest-payloads.ts";
import {
createCodeAgentChatResultStore,
handleCodeAgentCancelHttp,
handleCodeAgentChatHttp,
handleCodeAgentChatResultHttp,
handleCodeAgentInspectHttp,
handleCodeAgentSessionsHttp,
handleCodeAgentSteerHttp,
handleCodeAgentTraceHttp
} from "./server-code-agent-http.ts";
import { handleM3IoControlHttp } from "./server-m3-http.ts";
import { handleSkillsHttp } from "./server-skills-http.ts";
import {
configureSkillRuntime,
ensureCodexSkillsAggregationSync
} from "./skills-store.ts";
const DEFAULT_BODY_LIMIT_BYTES = 1024 * 1024;
function runtimeEnvironment(env = process.env) {
const value = env.HWLAB_GITOPS_PROFILE || env.HWLAB_ENVIRONMENT || ENVIRONMENT_DEV;
return typeof value === "string" && value.trim() ? value.trim() : ENVIRONMENT_DEV;
}
export function createCloudApiServer(options = {}) {
const env = options.env ?? process.env;
ensureCodeAgentRuntimeBase(env);
const sessionRegistry = options.sessionRegistry || createCodeAgentSessionRegistry();
const traceStore = options.traceStore || defaultCodeAgentTraceStore;
const codexStdioManager = options.codexStdioManager || createCodexStdioSessionManager({ traceStore });
const codeAgentChatResults = options.codeAgentChatResults || createCodeAgentChatResultStore({
maxResults: parsePositiveInteger(env.HWLAB_CODE_AGENT_CHAT_RESULT_CACHE_SIZE, 256)
});
const runtimeStore = options.runtimeStore || createConfiguredCloudRuntimeStore({ ...options, env });
const gatewayRegistry = options.gatewayRegistry || createGatewayDemoRegistry({
staleMs: parsePositiveInteger(env.HWLAB_GATEWAY_DEMO_STALE_MS, 30000),
dispatchTimeoutMs: parsePositiveInteger(env.HWLAB_GATEWAY_DISPATCH_TIMEOUT_MS, 120000)
});
const accessController = options.accessController || createAccessController({ ...options, env, runtimeStore, gatewayRegistry, traceStore, codeAgentChatResults });
if (typeof accessController.configureCodeAgentWorkspaceContext === "function") {
accessController.configureCodeAgentWorkspaceContext({ env, fetchImpl: options.fetchImpl, traceStore, codeAgentChatResults });
}
return createServer(async (request, response) => {
try {
await routeRequest(request, response, { ...options, env, runtimeStore, gatewayRegistry, accessController, sessionRegistry, traceStore, codexStdioManager, codeAgentChatResults });
} catch (error) {
if (error?.alreadySent) return;
sendJson(response, 500, {
error: {
code: "internal_error",
message: "hwlab-cloud-api request handling failed",
reason: error.message
}
});
}
});
}
export function ensureCodeAgentRuntimeBase(env = process.env) {
const codeAgentAdapter = String(env.HWLAB_CODE_AGENT_ADAPTER ?? env.HWLAB_CODE_AGENT_PROVIDER ?? "").trim().toLowerCase();
const useCodexStdioRuntime = !codeAgentAdapter || codeAgentAdapter === "codex-stdio";
const skillRuntime = configureSkillRuntime(env);
if (!useCodexStdioRuntime) return;
const workspace = env.HWLAB_CODE_AGENT_CODEX_WORKSPACE || env.HWLAB_CODE_AGENT_WORKSPACE || "/workspace/hwlab";
const codexHome = env.CODEX_HOME || env.HWLAB_CODE_AGENT_CODEX_HOME || "/codex-home";
const appCodexCommand = "/app/node_modules/.bin/codex";
const localCodexCommand = path.join(process.cwd(), "node_modules", ".bin", "codex");
const codexCommand = existsSync(appCodexCommand)
? appCodexCommand
: existsSync(localCodexCommand)
? localCodexCommand
: appCodexCommand;
env.HWLAB_CODE_AGENT_WORKSPACE ||= workspace;
env.HWLAB_CODE_AGENT_CODEX_WORKSPACE ||= workspace;
env.HWLAB_CODE_AGENT_CODEX_SANDBOX ||= "danger-full-access";
env.HWLAB_CODE_AGENT_CODEX_STDIO_ENABLED ||= "1";
env.HWLAB_CODE_AGENT_CODEX_STDIO_SUPERVISOR ||= "repo-owned";
env.CODEX_HOME ||= codexHome;
env.HWLAB_CODE_AGENT_CODEX_COMMAND ||= codexCommand;
for (const target of [path.dirname(workspace), codexHome]) {
try {
mkdirSync(target, { recursive: true });
} catch {
// Runtime health reports the concrete missing/writable blocker.
}
}
try {
if (!existsSync(workspace)) {
if (workspace === "/workspace/hwlab" && existsSync("/app")) {
symlinkSync("/app", workspace, "dir");
} else {
mkdirSync(workspace, { recursive: true });
}
}
lstatSync(workspace);
ensureCodexSkillsAggregationSync(skillRuntime);
} catch {
// Runtime health reports the concrete missing/writable blocker.
}
}
export async function buildHealthPayload(options = {}) {
const serviceId = CLOUD_API_SERVICE_ID;
const env = options.env ?? process.env;
ensureCodeAgentRuntimeBase(env);
const metadata = buildMetadataFromEnv(env, {
serviceId,
fallbackImageRepository: "ghcr.io/pikastech/hwlab-cloud-api"
});
const dbProbe = await buildDbRuntimeReadiness(env, options.dbProbe);
const codeAgent = await 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: runtimeEnvironment(env),
status: readiness.status,
ready: readiness.ready,
revision: metadata.revision,
readiness,
blockers: readiness.blockers,
blockerCodes: readiness.blockerCodes,
service: {
id: serviceId,
role: "cloud-api",
healthPath: "/health",
livePath: "/health/live"
},
commit: metadata.commit,
image: metadata.image,
build: metadata.build,
runtimeIdentity: metadata.runtime,
endpoint: 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 (url.pathname === "/auth" || url.pathname.startsWith("/auth/")) {
await options.accessController.handleAuthRoute(request, response, url);
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,
gatewayRegistry: options.gatewayRegistry
});
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 = await 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,
devicePod: {
route: "/v1/device-pods",
statusRoute: "/v1/device-pods/{devicePodId}/status",
eventsRoute: "/v1/device-pods/{devicePodId}/events",
contractVersion: "device-pod-authority-v1",
authority: "hwlab-cloud-api",
fakeFallback: false,
jobsRoute: "/v1/device-pods/{devicePodId}/jobs",
legacyHardwareUi: "disabled"
},
m3IoControl: describeM3IoControl(options)
});
return;
}
if (url.pathname === "/v1/auth/session" || url.pathname === "/v1/auth/login" || url.pathname === "/v1/auth/logout") {
const authUrl = new URL(url.href);
authUrl.pathname = url.pathname.replace(/^\/v1\/auth/u, "/auth");
await options.accessController.handleAuthRoute(request, response, authUrl);
return;
}
if (url.pathname === "/v1/skills" || url.pathname.startsWith("/v1/skills/")) {
await handleSkillsHttp(request, response, url, options);
return;
}
if (url.pathname === "/v1/users/me") {
await options.accessController.handleUserRoute(request, response, url);
return;
}
if (url.pathname === "/v1/api-keys" || url.pathname === "/v1/api-keys/default" || url.pathname.startsWith("/v1/api-keys/")) {
await options.accessController.handleUserRoute(request, response, url);
return;
}
if (url.pathname === "/v1/agent/conversations" && request.method === "GET") {
await options.accessController.listAgentConversations(request, response, url);
return;
}
const agentConversationMatch = url.pathname.match(/^\/v1\/agent\/conversations\/([^/]+)$/u);
if (agentConversationMatch && request.method === "GET") {
await options.accessController.getAgentConversation(request, response, decodeURIComponent(agentConversationMatch[1]));
return;
}
if (agentConversationMatch && (request.method === "PUT" || request.method === "PATCH")) {
await options.accessController.updateAgentConversation(request, response, decodeURIComponent(agentConversationMatch[1]));
return;
}
if (agentConversationMatch && request.method === "DELETE") {
await options.accessController.deleteAgentConversation(request, response, decodeURIComponent(agentConversationMatch[1]));
return;
}
if (url.pathname === "/v1/workbench/workspace" && request.method === "GET") {
await options.accessController.getWorkbenchWorkspace(request, response, url);
return;
}
const workbenchWorkspaceMatch = url.pathname.match(/^\/v1\/workbench\/workspace\/([^/]+)(?:\/(events|reset|select-conversation))?$/u);
if (workbenchWorkspaceMatch && !workbenchWorkspaceMatch[2] && (request.method === "PATCH" || request.method === "PUT")) {
await options.accessController.updateWorkbenchWorkspace(request, response, decodeURIComponent(workbenchWorkspaceMatch[1]));
return;
}
if (workbenchWorkspaceMatch && workbenchWorkspaceMatch[2] === "select-conversation" && request.method === "POST") {
await options.accessController.selectWorkbenchConversation(request, response, decodeURIComponent(workbenchWorkspaceMatch[1]));
return;
}
if (workbenchWorkspaceMatch && workbenchWorkspaceMatch[2] === "reset" && request.method === "POST") {
await options.accessController.resetWorkbenchWorkspace(request, response, decodeURIComponent(workbenchWorkspaceMatch[1]));
return;
}
if (workbenchWorkspaceMatch && workbenchWorkspaceMatch[2] === "events" && request.method === "GET") {
await options.accessController.workbenchWorkspaceEvents(request, response, url, decodeURIComponent(workbenchWorkspaceMatch[1]));
return;
}
if (url.pathname === "/v1/access/status" || url.pathname === "/v1/setup/status") {
await options.accessController.handleAccessStatusRoute(request, response, url);
return;
}
if (url.pathname === "/v1/setup/first-admin") {
await options.accessController.handleSetupRoute(request, response, url);
return;
}
if (url.pathname.startsWith("/v1/admin/")) {
await options.accessController.handleAdminRoute(request, response, url);
return;
}
if (url.pathname === "/v1/device-pods" || url.pathname.startsWith("/v1/device-pods/")) {
await options.accessController.handleDevicePodRoute(request, response, url);
return;
}
if (request.method === "GET" && url.pathname === M3_IO_CONTROL_ROUTE) {
sendJson(response, 200, await describeM3IoControlLive(options));
return;
}
if (request.method === "GET" && url.pathname === M3_STATUS_ROUTE) {
sendJson(response, 200, await describeM3StatusLive(options));
return;
}
if (request.method === "GET" && url.pathname === "/v1/diagnostics/gate") {
sendJson(response, 200, await buildGateDiagnosticsRows(options));
return;
}
if (request.method === "GET" && url.pathname === "/v1/live-builds") {
sendJson(response, 200, await buildLiveBuildsPayload(options, { buildHealthPayload, runtimeEnvironment }));
return;
}
if (request.method === "GET" && url.pathname === "/v1/gateway/sessions") {
sendJson(response, 200, options.gatewayRegistry.describe());
return;
}
if (url.pathname === "/v1/internal/device-pod/gateway-dispatch") {
await handleInternalDevicePodGatewayDispatch(request, response, options);
return;
}
if (request.method === "POST" && url.pathname === "/v1/gateway/poll") {
await handleGatewayPollHttp(request, response, options);
return;
}
if (request.method === "POST" && url.pathname === "/v1/gateway/result") {
await handleGatewayResultHttp(request, response, 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") {
const nextOptions = await codeAgentOptions(request, response, options);
if (!nextOptions) return;
await handleCodeAgentChatHttp(request, response, nextOptions);
return;
}
if (url.pathname === "/v1/agent/sessions" || url.pathname.startsWith("/v1/agent/sessions/")) {
const nextOptions = await codeAgentOptions(request, response, options, { required: true });
if (!nextOptions) return;
await handleCodeAgentSessionsHttp(request, response, url, nextOptions);
return;
}
if (request.method === "GET" && url.pathname === "/v1/agent/chat/inspect") {
const nextOptions = await codeAgentOptions(request, response, options);
if (!nextOptions) return;
await handleCodeAgentInspectHttp(request, response, url, nextOptions);
return;
}
if (request.method === "GET" && url.pathname.startsWith("/v1/agent/chat/result/")) {
const nextOptions = await codeAgentOptions(request, response, options);
if (!nextOptions) return;
await handleCodeAgentChatResultHttp(request, response, url, nextOptions);
return;
}
if (request.method === "POST" && url.pathname === "/v1/agent/chat/cancel") {
const nextOptions = await codeAgentOptions(request, response, options);
if (!nextOptions) return;
await handleCodeAgentCancelHttp(request, response, nextOptions);
return;
}
if (request.method === "POST" && url.pathname === "/v1/agent/chat/steer") {
const nextOptions = await codeAgentOptions(request, response, options);
if (!nextOptions) return;
await handleCodeAgentSteerHttp(request, response, nextOptions);
return;
}
if (request.method === "GET" && url.pathname.startsWith("/v1/agent/chat/trace/")) {
const nextOptions = await codeAgentOptions(request, response, options);
if (!nextOptions) return;
await handleCodeAgentTraceHttp(request, response, url, nextOptions);
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: runtimeEnvironment(options.env ?? process.env)
}
},
{
adapter: "rest",
runtimeStore: options.runtimeStore,
env: options.env,
dbProbe: options.dbProbe,
m3IoRequestJson: options.m3IoRequestJson,
gatewayRegistry: options.gatewayRegistry
}
);
sendJson(response, statusForRpcResponse(rpcResponse), rpcResponse);
}
async function handleInternalDevicePodGatewayDispatch(request, response, options) {
if (request.method !== "POST") {
sendJson(response, 405, { error: { code: "method_not_allowed", message: "Internal device-pod gateway dispatch only supports POST" } });
return;
}
if (!isDevicePodInternalCaller(request, options.env ?? process.env)) {
sendJson(response, 403, { error: { code: "device_pod_internal_authority_required", message: "Only hwlab-device-pod may use the internal gateway dispatch route" } });
return;
}
const body = await readJsonObject(request, options.bodyLimitBytes);
if (!body.ok) {
sendJson(response, 400, body.error);
return;
}
const params = body.value?.params && typeof body.value.params === "object" && !Array.isArray(body.value.params) ? body.value.params : {};
const gatewaySessionId = typeof params.gatewaySessionId === "string" ? params.gatewaySessionId.trim() : "";
if (!gatewaySessionId || options.gatewayRegistry?.isOnline?.(gatewaySessionId) !== true) {
sendJson(response, 409, {
accepted: false,
status: "blocked",
blocker: {
code: "gateway_dispatch_unavailable",
layer: "device-pod",
retryable: true,
summary: gatewaySessionId ? "gateway session is not connected" : "profile route.gatewaySessionId is missing"
}
});
return;
}
const operationId = typeof params.operationId === "string" && params.operationId ? params.operationId : `op_devicepod_gateway_${randomUUID()}`;
const traceId = typeof params.traceId === "string" && params.traceId ? params.traceId : `trc_devicepod_gateway_${randomUUID()}`;
const rpcResponse = await handleJsonRpcRequest({
jsonrpc: "2.0",
id: typeof body.value.id === "string" && body.value.id ? body.value.id : getHeader(request, "x-request-id") || `req_devicepod_gateway_${randomUUID()}`,
method: "hardware.invoke.shell",
params: { ...params, gatewaySessionId, operationId, traceId },
meta: {
traceId,
actorId: typeof body.value.actorId === "string" ? body.value.actorId : "svc_hwlab-device-pod",
serviceId: "hwlab-device-pod",
environment: runtimeEnvironment(options.env ?? process.env)
}
}, {
adapter: "device-pod-internal",
runtimeStore: options.runtimeStore,
env: options.env,
dbProbe: options.dbProbe,
m3IoRequestJson: options.m3IoRequestJson,
gatewayRegistry: options.gatewayRegistry
});
if (rpcResponse.error) {
sendJson(response, statusForRpcResponse(rpcResponse), rpcResponse);
return;
}
const result = rpcResponse.result ?? {};
const dispatch = result.dispatch ?? {};
const completed = result.status === "succeeded" || result.status === "completed" || dispatch.dispatchStatus === "succeeded" || dispatch.dispatchStatus === "completed";
sendJson(response, completed ? 200 : 409, createResultEnvelope({
id: rpcResponse.id,
result: {
accepted: completed,
status: completed ? "completed" : "failed",
operationId: result.operationId ?? operationId,
traceId,
gatewaySessionId,
resourceId: result.resourceId ?? params.resourceId ?? null,
capabilityId: result.capabilityId ?? params.capabilityId ?? null,
dispatch,
auditId: result.auditId ?? null,
evidenceId: result.evidenceId ?? null,
gateway: result.gateway ?? null,
blocker: completed ? null : {
code: "gateway_dispatch_failed",
layer: "device-pod",
retryable: true,
summary: dispatch.message ?? dispatch.stderr ?? dispatch.status ?? result.status ?? "gateway dispatch failed"
}
},
requestMeta: { traceId, serviceId: "hwlab-device-pod", environment: runtimeEnvironment(options.env ?? process.env) }
}));
}
async function codeAgentOptions(request, response, options, authOptions = {}) {
const auth = await options.accessController.authenticate(request, { required: authOptions.required ?? options.accessController.required });
if (!auth.ok) {
sendJson(response, auth.status, auth);
return null;
}
return { ...options, actor: auth.actor ?? null, authSession: auth.session ?? null };
}
async function readJsonObject(request, limitBytes) {
const body = await readBody(request, limitBytes);
try {
const value = body ? JSON.parse(body) : {};
if (!value || typeof value !== "object" || Array.isArray(value)) {
return { ok: false, error: { accepted: false, error: { code: "invalid_params", message: "body must be a JSON object" } } };
}
return { ok: true, value };
} catch (error) {
return { ok: false, error: { accepted: false, error: { code: "parse_error", message: "Invalid JSON body", reason: error.message } } };
}
}
function internalServiceHeader(request) {
return String(getHeader(request, "x-hwlab-internal-service") ?? "").trim();
}
function isDevicePodInternalCaller(request, env = process.env) {
const token = String(env.HWLAB_DEVICE_POD_INTERNAL_TOKEN ?? "").trim();
return Boolean(token)
&& internalServiceHeader(request) === "hwlab-device-pod"
&& String(getHeader(request, "x-hwlab-internal-token") ?? "").trim() === token;
}
async function handleGatewayPollHttp(request, response, options) {
const body = await readBody(request, options.bodyLimitBytes);
let params = {};
try {
params = body ? JSON.parse(body) : {};
} catch (error) {
sendJson(response, 400, {
accepted: false,
error: {
code: "parse_error",
message: "Invalid JSON body",
reason: error.message
}
});
return;
}
if (!params || typeof params !== "object" || Array.isArray(params)) {
sendJson(response, 400, {
accepted: false,
error: {
code: "invalid_params",
message: "gateway poll body must be a JSON object"
}
});
return;
}
const session = options.gatewayRegistry.updateSession(params);
await persistGatewayDemoSession({
runtimeStore: options.runtimeStore,
session,
request,
env: options.env
});
const outbound = options.gatewayRegistry.nextRequest(session.gatewaySessionId);
sendJson(response, 200, {
accepted: true,
serviceId: CLOUD_API_SERVICE_ID,
mode: "gateway-outbound-poll-demo",
gatewayId: session.gatewayId,
gatewaySessionId: session.gatewaySessionId,
queueDepth: session.queue.length,
request: outbound,
type: outbound ? "request" : "noop",
observedAt: new Date().toISOString()
});
}
async function handleGatewayResultHttp(request, response, options) {
const body = await readBody(request, options.bodyLimitBytes);
let params = {};
try {
params = body ? JSON.parse(body) : {};
} catch (error) {
sendJson(response, 400, {
accepted: false,
error: {
code: "parse_error",
message: "Invalid JSON body",
reason: error.message
}
});
return;
}
const result = options.gatewayRegistry.complete(params);
sendJson(response, result.accepted ? 200 : 404, {
...result,
serviceId: CLOUD_API_SERVICE_ID,
mode: "gateway-outbound-poll-demo",
observedAt: new Date().toISOString()
});
}
async function persistGatewayDemoSession({ runtimeStore, session, request, env = process.env }) {
if (!runtimeStore) return;
const refreshMs = parsePositiveInteger(env.HWLAB_GATEWAY_REGISTRATION_REFRESH_MS, 60000);
if (session.runtimeRegisteredEpochMs && Date.now() - session.runtimeRegisteredEpochMs < refreshMs) {
return;
}
const meta = {
traceId: getHeader(request, "x-trace-id") || `trc_gateway_poll_${randomUUID()}`,
actorId: `svc_${session.serviceId}`,
serviceId: CLOUD_API_SERVICE_ID,
environment: runtimeEnvironment(env)
};
try {
await runtimeStore.registerGatewaySession?.({
projectId: session.projectId,
gatewaySessionId: session.gatewaySessionId,
serviceId: "hwlab-gateway",
gatewayId: session.gatewayId,
endpoint: session.endpoint,
status: "connected",
labels: {
demoOutboundPoll: "true",
os: String(session.system?.platform ?? "unknown")
}
}, meta);
await runtimeStore.registerBoxResource?.({
projectId: session.projectId,
gatewaySessionId: session.gatewaySessionId,
resourceId: session.resourceId,
boxId: session.boxId,
resourceType: "simulator_endpoint",
state: "available",
metadata: {
serviceId: "hwlab-gateway",
demoShellHost: true
}
}, meta);
if (session.capabilities.length > 0) {
await runtimeStore.reportBoxCapabilities?.({
capabilities: session.capabilities.map((capability) => ({
...capability,
resourceId: capability.resourceId ?? session.resourceId,
projectId: capability.projectId ?? session.projectId
}))
}, meta);
}
session.runtimeRegisteredAt = new Date().toISOString();
session.runtimeRegisteredEpochMs = Date.now();
} catch {
// Demo polling must keep the gateway online even when durable runtime writes are not ready.
}
}
function sendRestError(request, response, statusCode, code, message, context = {}) {
const requestId = getHeader(request, "x-request-id") || "req_unassigned";
const environment = context.environment || runtimeEnvironment(context.env ?? process.env);
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,
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;
}
function firstHeaderValue(request, name) {
return getHeader(request, name);
}
function truthyFlag(value) {
return /^(?:1|true|yes|on|enabled)$/iu.test(String(value ?? "").trim());
}
function safeTraceId(value) {
const text = String(value ?? "").trim();
return /^trc_[A-Za-z0-9_.:-]+$/u.test(text) ? text : null;
}
function safeSessionId(value) {
const text = String(value ?? "").trim();
return /^ses_[A-Za-z0-9_.:-]+$/u.test(text) ? text : null;
}
function safeConversationId(value) {
const text = String(value ?? "").trim();
return /^cnv_[A-Za-z0-9_.:-]+$/u.test(text) ? text : null;
}
function safeOpaqueId(value) {
const text = String(value ?? "").trim();
return /^[A-Za-z0-9_.:-]{3,180}$/u.test(text) ? text : null;
}
function parsePositiveInteger(value, fallback) {
const parsed = Number.parseInt(value ?? "", 10);
return Number.isInteger(parsed) && parsed > 0 ? parsed : fallback;
}
function positiveInteger(value, fallback) {
const parsed = Number.parseInt(value ?? "", 10);
return Number.isInteger(parsed) && parsed > 0 ? parsed : fallback;
}
async function runtimeReadiness(runtimeStore) {
if (typeof runtimeStore?.readiness === "function") {
return runtimeStore.readiness();
}
return runtimeStore.summary();
}