Files
pikasTech-HWLAB/internal/cloud/hwpod-node-ws-registry.ts

446 lines
18 KiB
TypeScript

/*
* SPEC: PJ2026-010405 云端控制台;PJ2026-010103 HWPOD 服务。
* Implementation reference: draft-2026-07-13-p0-cloud-console。
* Responsibility: 管理主动出站 HWPOD Node 连接、并发门禁,以及用户操作与 readiness 探测的独立摘要。
*/
import { randomUUID } from "node:crypto";
import { createHwpodOperationLedger } from "./hwpod-operation-ledger.ts";
const DEFAULT_DISPATCH_TIMEOUT_MS = 30000;
type HwpodNodeConnection = {
id: string;
socket: any;
nodeId: string | null;
name: string | null;
capabilities: string[];
labels: Record<string, unknown>;
runtime: {
kind: string | null;
installed: boolean | null;
desktopVisible: boolean | null;
};
maxInFlight: number;
latestOperation: Record<string, unknown> | null;
latestReadinessProbe: Record<string, unknown> | null;
diagnostics: any[];
firstSeenAt: string;
lastSeenAt: string;
};
type PendingDispatch = {
planId: string;
nodeId: string;
requestId: string;
connection: HwpodNodeConnection;
operationKind: "user-operation" | "readiness-probe";
operation: Record<string, unknown>;
timer: ReturnType<typeof setTimeout>;
accepted: boolean;
ledgerWrite: Promise<boolean>;
promise: Promise<any>;
resolve: (value: any) => void;
};
type AuthoritativeOperation = {
planId: string;
nodeId: string;
operation: Record<string, unknown>;
promise: Promise<any>;
result: any | null;
};
export function createHwpodNodeWsRegistry(options: any = {}) {
const now = options.now ?? (() => new Date().toISOString());
const clock = options.clock ?? (() => Date.now());
const logger = options.logger ?? console;
const env = options.env ?? process.env;
const operationLedger = options.operationLedger ?? createHwpodOperationLedger({
runtimeStore: options.runtimeStore,
now,
retentionMs: positiveInteger(env.HWLAB_HWPOD_OPERATION_RETENTION_SECONDS, 0) * 1000
});
const connections = new Map<string, HwpodNodeConnection>();
const pending = new Map<string, PendingDispatch>();
const operations = new Map<string, AuthoritativeOperation>();
const admissions = new Map<string, { nodeId: string; promise: Promise<any> }>();
const socketConnections = new WeakMap<object, HwpodNodeConnection>();
function openBunSocket(socket: any) {
const connection: HwpodNodeConnection = {
id: `hwpod_node_conn_${randomUUID()}`,
socket,
nodeId: null,
name: null,
capabilities: [],
labels: {},
runtime: { kind: null, installed: null, desktopVisible: null },
maxInFlight: 1,
latestOperation: null,
latestReadinessProbe: null,
diagnostics: [],
firstSeenAt: now(),
lastSeenAt: now()
};
socketConnections.set(socket, connection);
return connection;
}
function handleBunMessage(socket: any, raw: unknown) {
const connection = socketConnections.get(socket);
if (!connection) return;
const text = typeof raw === "string" ? raw : Buffer.from(raw as any).toString("utf8");
handleText(connection, text);
}
function closeBunSocket(socket: any) {
const connection = socketConnections.get(socket);
if (!connection) return;
socketConnections.delete(socket);
removeConnection(connection);
}
function dispatch(plan: any, requestMeta: any = {}, dispatchOptions: any = {}) {
const planId = safeText(requestMeta.operationKind) === "readiness-probe" ? "" : safeText(plan?.planId);
const nodeId = safeNodeId(plan?.nodeId);
const admission = planId ? admissions.get(planId) : null;
if (admission) {
return admission.nodeId === nodeId
? admission.promise
: Promise.resolve(blockedDispatch(plan, requestMeta, `hwpod operation ${planId} is already owned by node ${admission.nodeId}`, "hwpod_operation_identity_conflict"));
}
const promise = dispatchOperation(plan, requestMeta, dispatchOptions);
if (planId) {
admissions.set(planId, { nodeId, promise });
void promise.finally(() => admissions.delete(planId));
}
return promise;
}
async function dispatchOperation(plan: any, requestMeta: any = {}, dispatchOptions: any = {}) {
const nodeId = safeNodeId(plan?.nodeId);
const operationKind = safeText(requestMeta.operationKind) === "readiness-probe" ? "readiness-probe" : "user-operation";
const planId = operationKind === "user-operation" ? safeText(plan?.planId) : "";
const authoritative = planId ? operations.get(planId) : null;
if (authoritative) {
if (authoritative.nodeId !== nodeId) {
return blockedDispatch(plan, requestMeta, `hwpod operation ${planId} is already owned by node ${authoritative.nodeId}`, "hwpod_operation_identity_conflict");
}
return authoritative.result === null ? authoritative.promise : authoritative.result;
}
const durable = planId ? await operationLedger.get(planId) : null;
if (durable?.nodeId && durable.nodeId !== nodeId) {
return blockedDispatch(plan, requestMeta, `hwpod operation ${planId} is already owned by node ${durable.nodeId}`, "hwpod_operation_identity_conflict");
}
if (durable?.status === "completed" && durable.result) {
return durable.result;
}
const connection = nodeId ? connections.get(nodeId) : null;
if (!nodeId || !connection) {
return blockedDispatch(plan, requestMeta, `hwpod-node ${nodeId || "<missing>"} is not connected through outbound WebSocket`);
}
if (inFlightCountForNode(nodeId, pending) >= connection.maxInFlight) {
return blockedDispatch(plan, requestMeta, `hwpod-node ${nodeId} reached its maxInFlight=${connection.maxInFlight} limit`, "hwpod_node_busy");
}
const requestId = safeText(requestMeta.requestId) || `req_hwpod_ws_${randomUUID()}`;
const timeoutMs = positiveInteger(dispatchOptions.timeoutMs, DEFAULT_DISPATCH_TIMEOUT_MS);
const operation = operationStarted(plan, requestId, now());
if (planId) await operationLedger.start({ planId, nodeId, operation });
if (operationKind === "readiness-probe") connection.latestReadinessProbe = operation;
else connection.latestOperation = operation;
let resolveDispatch: (value: any) => void = () => {};
const promise = new Promise((resolve) => { resolveDispatch = resolve; });
const waiter = {
planId,
nodeId,
requestId,
connection,
operationKind,
operation,
timer: null as unknown as ReturnType<typeof setTimeout>,
accepted: false,
ledgerWrite: Promise.resolve(true),
promise,
resolve: resolveDispatch
} satisfies PendingDispatch;
if (planId) operations.set(planId, { planId, nodeId, operation, promise, result: null });
waiter.timer = setTimeout(() => {
finishTransportFailure(waiter, blockedDispatch(plan, requestMeta, `hwpod-node ${nodeId} WebSocket dispatch timed out after ${timeoutMs}ms`), "timed-out");
}, timeoutMs);
pending.set(requestId, waiter);
try {
sendJson(connection, { type: "hwpod-node-ops", requestId, plan, requestMeta });
} catch (error) {
finishTransportFailure(waiter, blockedDispatch(plan, requestMeta, error instanceof Error ? error.message : String(error)), "failed");
}
void promise.then((result) => {
const authoritativeOperation = planId ? operations.get(planId) : null;
if (authoritativeOperation?.promise === promise) {
authoritativeOperation.result = result;
authoritativeOperation.operation = operationFinished(
authoritativeOperation.operation,
safeText(result?.status) || (result?.ok === false ? "failed" : "completed"),
now(),
safeText(result?.blocker?.code) || null
);
}
});
return promise;
}
async function lookup(planId: string) {
const authoritative = operations.get(safeText(planId));
if (authoritative) return { planId: authoritative.planId, nodeId: authoritative.nodeId, operation: authoritative.operation, result: authoritative.result };
const durable = await operationLedger.get(safeText(planId));
return durable ? { planId: durable.planId, nodeId: durable.nodeId, operation: durable.operation, result: durable.result } : null;
}
function describe() {
const observedAt = now();
return {
mode: "hwpod-node-outbound-native-ws",
route: "/v1/hwpod-node/ws",
connectedCount: connections.size,
pendingCount: pending.size,
nodes: [...connections.values()].map((connection) => ({
nodeId: connection.nodeId,
name: connection.name,
status: "online",
capabilityCount: connection.capabilities.length,
capabilities: connection.capabilities,
platform: safeText(connection.labels.platform) || "unknown",
version: safeText(connection.labels.version) || null,
runtimeKind: connection.runtime.kind,
installed: connection.runtime.installed,
desktopVisible: connection.runtime.desktopVisible,
maxInFlight: connection.maxInFlight,
inFlightCount: inFlightCountForNode(connection.nodeId, pending),
busy: inFlightCountForNode(connection.nodeId, pending) >= connection.maxInFlight,
latestOperation: connection.latestOperation,
latestReadinessProbe: connection.latestReadinessProbe,
diagnosticCount: connection.diagnostics.length,
lastDiagnostic: connection.diagnostics.at(-1) ?? null,
diagnostics: connection.diagnostics.slice(-10),
firstSeenAt: connection.firstSeenAt,
lastSeenAt: connection.lastSeenAt,
ageMs: Math.max(0, clock() - Date.parse(connection.lastSeenAt || observedAt))
}))
};
}
function hasNode(nodeId: string) {
return connections.has(nodeId);
}
function handleText(connection: HwpodNodeConnection, text: string) {
let message: any;
try {
message = JSON.parse(text);
} catch (error) {
sendJson(connection, { type: "ack", requestId: "message", ok: false, message: "invalid JSON message" });
return;
}
const messageType = safeText(message.type);
const nodeId = safeNodeId(message.nodeId);
if (messageType === "register") {
if (!nodeId) {
sendJson(connection, { type: "ack", requestId: "register", ok: false, message: "nodeId is required" });
return;
}
connection.nodeId = nodeId;
connection.name = safeText(message.name) || nodeId;
connection.capabilities = Array.isArray(message.capabilities) ? message.capabilities.map(String).filter(Boolean).slice(0, 64) : [];
connection.labels = message.labels && typeof message.labels === "object" && !Array.isArray(message.labels) ? message.labels : {};
connection.runtime = normalizeRuntime(message.runtime, connection.labels);
connection.maxInFlight = positiveInteger(message.maxInFlight ?? connection.labels.maxInFlight, 1);
connection.lastSeenAt = now();
const previous = connections.get(nodeId);
if (previous && previous !== connection) previous.socket.close();
connections.set(nodeId, connection);
sendJson(connection, { type: "ack", requestId: "register", ok: true, message: "registered", nodeId, operationRetentionSeconds: Math.floor(operationLedger.retentionMs / 1000) });
return;
}
if (messageType === "heartbeat") {
if (nodeId && connection.nodeId === nodeId) connection.lastSeenAt = now();
return;
}
if (messageType === "hwpod-node-diagnostic") {
if (nodeId && connection.nodeId === nodeId) connection.lastSeenAt = now();
const diagnostic = normalizeDiagnostic(message.diagnostic, now());
connection.diagnostics.push(diagnostic);
if (connection.diagnostics.length > 50) connection.diagnostics = connection.diagnostics.slice(-50);
logger.warn?.(JSON.stringify({ event: "hwpod-node.diagnostic", nodeId: connection.nodeId || nodeId || null, connectionId: connection.id, diagnostic, valuesPrinted: false }));
sendJson(connection, { type: "ack", requestId: safeText(message.requestId) || "diagnostic", ok: true, message: "diagnostic recorded" });
return;
}
if (messageType === "hwpod-node-ops-accepted") {
const waiter = pending.get(safeText(message.requestId));
if (waiter?.planId && !waiter.accepted) {
waiter.accepted = true;
waiter.ledgerWrite = waiter.ledgerWrite.then(async () => {
await operationLedger.accept(waiter.planId);
return true;
}).catch((error) => {
finishLedgerFailure(waiter, "accept", error);
return false;
});
}
return;
}
if (messageType === "hwpod-node-ops-result") void completeDispatch(message);
}
async function completeDispatch(message: any) {
const requestId = safeText(message.requestId);
const waiter = requestId ? pending.get(requestId) : null;
if (!waiter) return;
clearTimeout(waiter.timer);
const result = message.result && typeof message.result === "object" && !Array.isArray(message.result)
? message.result
: { ok: false, status: "failed", results: [], blocker: { code: "hwpod_node_result_invalid", layer: "hwpod-node", retryable: true, summary: "hwpod-node WebSocket result payload is invalid" } };
finishPendingOperation(
waiter,
safeText(result.status) || (result.ok === false ? "failed" : "completed"),
now(),
safeText(result.blocker?.code) || null
);
if (waiter.planId && !(await waiter.ledgerWrite)) return;
try {
if (waiter.planId) await operationLedger.complete(waiter.planId, result);
} catch (error) {
finishLedgerFailure(waiter, "complete", error);
return;
}
pending.delete(requestId);
if (waiter.planId) operations.delete(waiter.planId);
waiter.resolve(result);
}
function removeConnection(connection: HwpodNodeConnection) {
if (connection.nodeId && connections.get(connection.nodeId) === connection) connections.delete(connection.nodeId);
for (const [requestId, waiter] of pending) {
if (waiter.nodeId !== connection.nodeId) continue;
clearTimeout(waiter.timer);
finishTransportFailure(waiter, blockedDispatch({ nodeId: waiter.nodeId, ops: [] }, { requestId }, `hwpod-node ${waiter.nodeId} WebSocket disconnected before returning result`), "disconnected");
}
}
function finishTransportFailure(waiter: PendingDispatch, result: any, status: string) {
clearTimeout(waiter.timer);
pending.delete(waiter.requestId);
finishPendingOperation(waiter, status, now(), "hwpod_node_unavailable");
if (waiter.planId && operations.get(waiter.planId)?.promise === waiter.promise) operations.delete(waiter.planId);
waiter.resolve(result);
}
function finishLedgerFailure(waiter: PendingDispatch, phase: "accept" | "complete", error: unknown) {
clearTimeout(waiter.timer);
pending.delete(waiter.requestId);
finishPendingOperation(waiter, "failed", now(), "hwpod_operation_ledger_write_failed");
if (waiter.planId) operations.delete(waiter.planId);
const summary = `hwpod operation ledger ${phase} failed: ${error instanceof Error ? error.message : String(error)}`;
waiter.resolve(blockedDispatch({ nodeId: waiter.nodeId, ops: [] }, { requestId: waiter.requestId }, summary, "hwpod_operation_ledger_write_failed"));
}
return { openBunSocket, handleBunMessage, closeBunSocket, dispatch, lookup, describe, hasNode };
}
function sendJson(connection: HwpodNodeConnection, value: any) {
connection.socket.send(JSON.stringify(value));
}
function blockedDispatch(plan: any, requestMeta: any, summary: string, code = "hwpod_node_unavailable") {
const ops = Array.isArray(plan?.ops) ? plan.ops : [];
return {
ok: false,
status: "blocked",
results: ops.map((op: any, index: number) => ({
opId: safeText(op?.opId) || `op_${index + 1}`,
op: safeText(op?.op) || "unknown",
ok: false,
status: "blocked",
blocker: { code, layer: "hwpod-node", retryable: true, summary }
})),
blocker: { code, layer: "hwpod-node", retryable: true, summary },
requestMeta
};
}
function safeText(value: unknown) {
return typeof value === "string" ? value.trim() : "";
}
function safeNodeId(value: unknown) {
const text = safeText(value);
return /^[A-Za-z0-9._:-]{1,128}$/u.test(text) ? text : "";
}
function normalizeDiagnostic(value: any, observedAt: string) {
const input = value && typeof value === "object" && !Array.isArray(value) ? value : {};
return {
seq: positiveInteger(input.seq, 0),
level: safeText(input.level) || "ERROR",
source: safeText(input.source) || "unknown",
message: safeText(input.message).slice(0, 2000),
details: input.details && typeof input.details === "object" && !Array.isArray(input.details) ? input.details : {},
observedAt: safeText(input.observedAt) || observedAt,
receivedAt: observedAt
};
}
function normalizeRuntime(value: any, labels: Record<string, unknown>) {
const input = value && typeof value === "object" && !Array.isArray(value) ? value : {};
return {
kind: safeText(input.kind ?? labels.runtimeKind) || null,
installed: booleanOrNull(input.installed ?? labels.installed),
desktopVisible: booleanOrNull(input.desktopVisible ?? labels.desktopVisible)
};
}
function operationStarted(plan: any, requestId: string, startedAt: string) {
return {
requestId,
planId: safeText(plan?.planId) || null,
hwpodId: safeText(plan?.hwpodId) || null,
status: "running",
opCount: Array.isArray(plan?.ops) ? plan.ops.length : 0,
startedAt,
finishedAt: null,
blockerCode: null
};
}
function operationFinished(operation: Record<string, unknown> | null, status: string, finishedAt: string, blockerCode: string | null) {
return {
...(operation ?? {}),
status,
finishedAt,
blockerCode
};
}
function finishPendingOperation(waiter: PendingDispatch, status: string, finishedAt: string, blockerCode: string | null) {
const finished = operationFinished(waiter.operation, status, finishedAt, blockerCode);
if (waiter.operationKind === "readiness-probe") {
if (waiter.connection.latestReadinessProbe?.requestId === waiter.requestId) waiter.connection.latestReadinessProbe = finished;
return;
}
if (waiter.connection.latestOperation?.requestId === waiter.requestId) waiter.connection.latestOperation = finished;
}
function inFlightCountForNode(nodeId: string | null, pending: Map<string, PendingDispatch>) {
if (!nodeId) return 0;
let count = 0;
for (const waiter of pending.values()) if (waiter.nodeId === nodeId) count += 1;
return count;
}
function booleanOrNull(value: unknown) {
return typeof value === "boolean" ? value : null;
}
function positiveInteger(value: unknown, fallback: number) {
const numberValue = Number(value);
return Number.isInteger(numberValue) && numberValue > 0 ? numberValue : fallback;
}