Files
pikasTech-HWLAB/internal/cloud/gateway-demo-registry.mjs
T
2026-05-25 03:16:55 +00:00

195 lines
6.2 KiB
JavaScript

import { randomUUID } from "node:crypto";
import { ENVIRONMENT_DEV, JSON_RPC_VERSION } from "../protocol/index.mjs";
import { CLOUD_API_SERVICE_ID } from "../audit/index.mjs";
const DEFAULT_STALE_MS = 30000;
const DEFAULT_DISPATCH_TIMEOUT_MS = 120000;
export function createGatewayDemoRegistry({
now = () => new Date().toISOString(),
clock = () => Date.now(),
staleMs = DEFAULT_STALE_MS,
dispatchTimeoutMs = DEFAULT_DISPATCH_TIMEOUT_MS
} = {}) {
const sessions = new Map();
const pending = new Map();
function updateSession(input = {}) {
const gatewaySessionId = input.gatewaySessionId || `gws_${input.gatewayId || "gateway_demo"}`;
const gatewayId = input.gatewayId || gatewaySessionId.replace(/^gws_/, "") || "gateway_demo";
const previous = sessions.get(gatewaySessionId);
const observedAt = now();
const session = {
...(previous ?? {}),
serviceId: input.serviceId || "hwlab-gateway",
projectId: input.projectId || previous?.projectId || "prj_mvp_topology",
gatewayId,
gatewaySessionId,
status: "connected",
endpoint: input.endpoint ?? previous?.endpoint ?? null,
resourceId: input.resourceId ?? previous?.resourceId ?? "res_windows_host",
boxId: input.boxId ?? previous?.boxId ?? "box_windows_host",
capabilities: Array.isArray(input.capabilities) ? input.capabilities : previous?.capabilities ?? [],
outbound: normalizeGatewayOutbound(input.outbound ?? previous?.outbound),
system: input.system ?? previous?.system ?? {},
firstSeenAt: previous?.firstSeenAt ?? observedAt,
lastSeenAt: observedAt,
lastSeenEpochMs: clock(),
queue: previous?.queue ?? []
};
sessions.set(gatewaySessionId, session);
return session;
}
function isOnline(gatewaySessionId, { maxAgeMs = staleMs } = {}) {
const session = sessions.get(gatewaySessionId);
if (!session) return false;
return clock() - session.lastSeenEpochMs <= maxAgeMs;
}
function enqueue({ gatewaySessionId, request, timeoutMs = dispatchTimeoutMs }) {
const session = sessions.get(gatewaySessionId);
if (!session || !isOnline(gatewaySessionId)) {
return Promise.resolve({
ok: false,
status: "not_connected",
error: `gatewaySessionId ${gatewaySessionId} is not connected`
});
}
const requestId = request.id ?? `req_gateway_demo_${randomUUID()}`;
const outbound = {
...request,
id: requestId
};
return new Promise((resolve) => {
const timer = setTimeout(() => {
pending.delete(requestId);
resolve({
ok: false,
status: "timed_out",
error: `gateway dispatch timed out after ${timeoutMs}ms`,
request: outbound
});
}, timeoutMs);
pending.set(requestId, {
gatewaySessionId,
resolve,
timer,
request: outbound,
createdAt: now()
});
session.queue.push(outbound);
});
}
function nextRequest(gatewaySessionId) {
const session = sessions.get(gatewaySessionId);
return session?.queue.shift() ?? null;
}
function complete(input = {}) {
const response = input.response ?? input;
const requestId = response.id;
const waiter = pending.get(requestId);
if (!waiter) {
return {
accepted: false,
status: "unknown_request",
requestId: requestId ?? null
};
}
clearTimeout(waiter.timer);
pending.delete(requestId);
waiter.resolve({
ok: !response.error,
status: response.error ? "failed" : "completed",
response
});
return {
accepted: true,
status: "completed",
requestId
};
}
function describe() {
return {
serviceId: CLOUD_API_SERVICE_ID,
environment: ENVIRONMENT_DEV,
mode: "gateway-outbound-poll-demo",
staleMs,
dispatchTimeoutMs,
sessions: [...sessions.values()].map((session) => ({
serviceId: session.serviceId,
projectId: session.projectId,
gatewayId: session.gatewayId,
gatewaySessionId: session.gatewaySessionId,
status: isOnline(session.gatewaySessionId) ? "online" : "stale",
resourceId: session.resourceId,
boxId: session.boxId,
capabilityCount: session.capabilities.length,
queueDepth: session.queue.length,
inflightCount: session.outbound.inflightCount,
maxInflightRequests: session.outbound.maxInflightRequests,
inflightRequests: session.outbound.inflightRequests,
firstSeenAt: session.firstSeenAt,
lastSeenAt: session.lastSeenAt,
system: session.system
})),
pendingCount: pending.size
};
}
return {
updateSession,
isOnline,
enqueue,
nextRequest,
complete,
describe,
getSession: (gatewaySessionId) => sessions.get(gatewaySessionId) ?? null
};
}
function normalizeGatewayOutbound(input = {}) {
const source = input && typeof input === "object" && !Array.isArray(input) ? input : {};
const inflightRequests = Array.isArray(source.inflightRequests)
? source.inflightRequests.map((request) => ({
requestId: request?.requestId ?? null,
method: request?.method ?? null,
operationId: request?.operationId ?? null,
traceId: request?.traceId ?? null,
startedAt: request?.startedAt ?? null,
durationMs: Number.isInteger(request?.durationMs) ? request.durationMs : null
}))
: [];
return {
pollIntervalMs: Number.isInteger(source.pollIntervalMs) ? source.pollIntervalMs : null,
commandExecutionEnabled: source.commandExecutionEnabled === true,
maxInflightRequests: Number.isInteger(source.maxInflightRequests) ? source.maxInflightRequests : null,
inflightCount: Number.isInteger(source.inflightCount) ? source.inflightCount : inflightRequests.length,
inflightRequests
};
}
export function createGatewayShellRequest({ id, params = {}, meta = {} }) {
return {
jsonrpc: JSON_RPC_VERSION,
id: id ?? `req_gateway_shell_${randomUUID()}`,
method: "hardware.invoke.shell",
params,
meta: {
traceId: meta.traceId ?? params.traceId ?? `trc_gateway_shell_${randomUUID()}`,
actorId: meta.actorId,
serviceId: CLOUD_API_SERVICE_ID,
environment: ENVIRONMENT_DEV
}
};
}