Files
pikasTech-HWLAB/internal/patchpanel/runtime.mjs
T
2026-05-21 18:36:39 +00:00

471 lines
14 KiB
JavaScript

import { readFile } from "node:fs/promises";
import { request as httpRequest } from "node:http";
import { request as httpsRequest } from "node:https";
import { ENVIRONMENT_DEV } from "../protocol/index.mjs";
import { DEFAULT_GATEWAY_SESSION_ID, DEFAULT_PROJECT_ID } from "../sim/model.mjs";
import {
PATCH_PANEL_SERVICE_ID,
createPatchPanel,
createWiringConfig,
endpointKey
} from "./model.mjs";
const DEFAULT_APPLIED_AT = "not_applied";
function isoNow() {
return new Date().toISOString();
}
function clone(value) {
return JSON.parse(JSON.stringify(value));
}
function isObject(value) {
return value !== null && typeof value === "object" && !Array.isArray(value);
}
function ensureObject(value, label) {
if (!isObject(value)) {
throw new Error(`${label} must be an object`);
}
}
function normalizeEndpointMap(raw = {}) {
ensureObject(raw, "endpoint map");
const byResourceId = new Map();
const diagnostics = [];
for (const [resourceId, endpoint] of Object.entries(raw)) {
if (typeof endpoint !== "string" || endpoint.trim().length === 0) {
diagnostics.push({
level: "error",
code: "endpoint_url_invalid",
message: `endpoint map entry ${resourceId} must be a non-empty URL string`,
path: `$.endpointMap.${resourceId}`
});
continue;
}
try {
const normalized = new URL(endpoint).toString();
byResourceId.set(resourceId, normalized);
} catch {
diagnostics.push({
level: "error",
code: "endpoint_url_invalid",
message: `endpoint map entry ${resourceId} is not a valid URL: ${endpoint}`,
path: `$.endpointMap.${resourceId}`
});
}
}
return { byResourceId, diagnostics };
}
function wiringFromRuntimeConfig(config, now) {
if (isObject(config.wiringConfig)) {
return clone(config.wiringConfig);
}
return createWiringConfig({
wiringConfigId: config.wiringConfigId ?? "wir_patch_panel_runtime",
projectId: config.projectId ?? DEFAULT_PROJECT_ID,
gatewaySessionId: config.gatewaySessionId ?? DEFAULT_GATEWAY_SESSION_ID,
name: config.name ?? "Patch panel runtime wiring",
status: config.status ?? "active",
connections: config.connections,
constraints: config.constraints,
now
});
}
export async function readPatchPanelRuntimeConfig(path) {
const raw = await readFile(path, "utf8");
return JSON.parse(raw);
}
export function createPatchPanelRuntime({
wiringConfig,
endpointMap = {},
configSource = "inline",
now = isoNow,
requestJson = postJson
}) {
let appliedAt = DEFAULT_APPLIED_AT;
let activeWiringConfig = null;
let patchPanel = null;
let endpoints = new Map();
let configDiagnostics = [];
let evidenceSummary = {
configApplied: false,
configSource,
appliedAt,
lastSyncAt: "not_run",
routeCount: 0,
tickCount: 0,
deliveryCount: 0,
failedDeliveryCount: 0
};
function applyConfig(config) {
const observedAt = now();
const source = config.configSource ?? configSource;
const runtimeConfig = isObject(config.wiringConfig) || Array.isArray(config.connections) ? config : {};
const nextWiringConfig = wiringFromRuntimeConfig(runtimeConfig, observedAt);
const endpointResult = normalizeEndpointMap(runtimeConfig.endpointMap ?? config.endpointMap ?? endpointMap);
const nextPatchPanel = createPatchPanel({
wiringConfig: nextWiringConfig,
now
});
const endpointDiagnostics = endpointResult.diagnostics;
const targetDiagnostics = [];
for (const connection of nextPatchPanel.normalizedConnections) {
const targetResourceId = connection.to.resourceId;
if (connection.syncable && !endpointResult.byResourceId.has(targetResourceId)) {
targetDiagnostics.push({
level: "error",
code: "target_endpoint_missing",
message: `sync target ${targetResourceId}:${connection.to.port} has no endpoint URL in endpointMap`,
path: `$.endpointMap.${targetResourceId}`
});
}
}
activeWiringConfig = nextWiringConfig;
patchPanel = nextPatchPanel;
endpoints = endpointResult.byResourceId;
configDiagnostics = [
...nextPatchPanel.diagnostics,
...endpointDiagnostics,
...targetDiagnostics
];
appliedAt = observedAt;
evidenceSummary = {
...evidenceSummary,
configApplied: true,
configSource: source,
appliedAt,
wiringConfigId: activeWiringConfig.wiringConfigId,
connectionCount: activeWiringConfig.connections.length,
syncableConnectionCount: nextPatchPanel.normalizedConnections.filter((connection) => connection.syncable).length,
diagnostics: diagnosticSummary(configDiagnostics)
};
return {
accepted: diagnosticSummary(configDiagnostics).errorCount === 0,
appliedAt,
wiringConfigId: activeWiringConfig.wiringConfigId,
diagnostics: configDiagnostics,
summary: evidenceSummary
};
}
async function reloadFromFile(path) {
const config = await readPatchPanelRuntimeConfig(path);
return applyConfig({
...config,
configSource: path
});
}
function diagnostics() {
return {
serviceId: PATCH_PANEL_SERVICE_ID,
environment: ENVIRONMENT_DEV,
state: diagnosticSummary(configDiagnostics).errorCount > 0 ? "faulted" : "active",
observedAt: now(),
configSource: evidenceSummary.configSource,
appliedAt,
wiringConfigId: activeWiringConfig?.wiringConfigId ?? null,
endpointMap: Object.fromEntries(endpoints.entries()),
diagnostics: configDiagnostics,
summary: evidenceSummary
};
}
function status() {
const state = patchPanel.status();
const summary = diagnosticSummary(configDiagnostics);
return {
...state,
state: summary.errorCount > 0 ? "faulted" : state.state,
metadata: {
...state.metadata,
configSource: evidenceSummary.configSource,
appliedAt,
diagnostics: summary,
evidenceSummary
}
};
}
async function routeSignal(signal) {
const result = patchPanel.routeSignal(signal);
const diagnostics = [...(result.diagnostics ?? [])];
if (!result.accepted) {
evidenceSummary = {
...evidenceSummary,
routeCount: evidenceSummary.routeCount + 1,
lastSyncAt: result.observedAt,
diagnostics: diagnosticSummary([...configDiagnostics, ...diagnostics])
};
return {
...result,
...(diagnostics.length > 0 ? { diagnostics } : {})
};
}
const hasEndpointDeliveries = result.deliveries.some((delivery) => {
const connection = findDeliveryConnection(delivery);
return connection?.syncable === true;
});
if (!hasEndpointDeliveries) {
evidenceSummary = {
...evidenceSummary,
routeCount: evidenceSummary.routeCount + 1,
lastSyncAt: result.observedAt,
deliveryCount: evidenceSummary.deliveryCount + result.deliveries.length,
diagnostics: diagnosticSummary([...configDiagnostics, ...diagnostics])
};
return {
...result,
...(diagnostics.length > 0 ? { diagnostics } : {})
};
}
const deliveries = [];
for (const delivery of result.deliveries) {
const connection = findDeliveryConnection(delivery);
const write = await deliverToEndpoint(delivery, { requireEndpoint: connection?.syncable === true });
deliveries.push({
...delivery,
deliveryStatus: write.deliveryStatus ?? (write.accepted ? "applied" : "failed"),
endpoint: write.endpoint ?? null,
diagnostic: write.diagnostic ?? null
});
if (write.diagnostic) {
diagnostics.push(write.diagnostic);
}
}
const failedDeliveryCount = deliveries.filter((delivery) => delivery.deliveryStatus === "failed").length;
evidenceSummary = {
...evidenceSummary,
routeCount: evidenceSummary.routeCount + 1,
lastSyncAt: result.observedAt,
deliveryCount: evidenceSummary.deliveryCount + deliveries.length,
failedDeliveryCount: evidenceSummary.failedDeliveryCount + failedDeliveryCount,
diagnostics: diagnosticSummary([...configDiagnostics, ...diagnostics])
};
return {
...result,
accepted: result.accepted && failedDeliveryCount === 0,
deliveryCount: deliveries.length,
deliveries,
...(diagnostics.length > 0 ? { diagnostics } : {})
};
}
async function tick({ signals = [] } = {}) {
if (!Array.isArray(signals)) {
throw new Error("tick signals must be an array");
}
const startedAt = now();
const routed = [];
for (const signal of signals) {
routed.push(await routeSignal(signal));
}
evidenceSummary = {
...evidenceSummary,
tickCount: evidenceSummary.tickCount + 1,
lastSyncAt: now()
};
return {
accepted: routed.every((result) => result.accepted),
propagatedBy: PATCH_PANEL_SERVICE_ID,
startedAt,
observedAt: evidenceSummary.lastSyncAt,
routeCount: routed.length,
deliveryCount: routed.reduce((sum, result) => sum + result.deliveryCount, 0),
routes: routed,
diagnostics: routed.flatMap((result) => result.diagnostics ?? []),
evidenceSummary
};
}
function findDeliveryConnection(delivery) {
return patchPanel.normalizedConnections.find(
(connection) =>
connection.from.resourceId === delivery.sourceResourceId &&
connection.from.port === delivery.sourcePort &&
connection.to.resourceId === delivery.resourceId &&
connection.to.port === delivery.port
);
}
async function deliverToEndpoint(delivery, { requireEndpoint }) {
const endpoint = endpoints.get(delivery.resourceId);
if (!endpoint) {
if (!requireEndpoint) {
return {
accepted: true,
deliveryStatus: "recorded"
};
}
return {
accepted: false,
diagnostic: {
level: "error",
code: "target_endpoint_missing",
message: `cannot deliver ${delivery.sourceResourceId}:${delivery.sourcePort} -> ${delivery.resourceId}:${delivery.port}; endpointMap is missing ${delivery.resourceId}`,
path: `$.endpointMap.${delivery.resourceId}`,
endpointKey: endpointKey(delivery.resourceId, delivery.port)
}
};
}
const url = new URL("/internal/ports/write", endpoint.endsWith("/") ? endpoint : `${endpoint}/`).toString();
try {
const response = await requestJson(url, {
port: delivery.port,
value: delivery.value,
source: "patch-panel",
sourceResourceId: delivery.sourceResourceId,
sourcePort: delivery.sourcePort,
propagatedBy: PATCH_PANEL_SERVICE_ID,
observedAt: delivery.observedAt
});
if (response.status < 200 || response.status >= 300 || response.json?.accepted !== true) {
return {
accepted: false,
endpoint: url,
diagnostic: {
level: "error",
code: "target_write_rejected",
message: `target endpoint rejected ${delivery.resourceId}:${delivery.port} write with status ${response.status}`,
path: `$.endpointMap.${delivery.resourceId}`,
response: response.json ?? response.bodyPreview
}
};
}
return {
accepted: true,
endpoint: url,
response: response.json
};
} catch (error) {
return {
accepted: false,
endpoint: url,
diagnostic: {
level: "error",
code: "target_write_failed",
message: `failed to deliver ${delivery.resourceId}:${delivery.port}: ${error instanceof Error ? error.message : String(error)}`,
path: `$.endpointMap.${delivery.resourceId}`
}
};
}
}
applyConfig({
wiringConfig,
endpointMap,
configSource
});
return {
applyConfig,
reloadFromFile,
diagnostics,
status,
routeSignal,
tick,
get wiringConfig() {
return activeWiringConfig;
},
get endpointMap() {
return Object.fromEntries(endpoints.entries());
},
get evidenceSummary() {
return evidenceSummary;
}
};
}
export async function createPatchPanelRuntimeFromFile({ path, now = isoNow, requestJson = postJson }) {
const config = await readPatchPanelRuntimeConfig(path);
return createPatchPanelRuntime({
...config,
configSource: path,
now,
requestJson
});
}
export function diagnosticSummary(diagnostics) {
return {
errorCount: diagnostics.filter((item) => item.level === "error").length,
warningCount: diagnostics.filter((item) => item.level === "warning").length,
infoCount: diagnostics.filter((item) => item.level === "info").length
};
}
export async function postJson(url, body, { timeoutMs = 5000 } = {}) {
return new Promise((resolve, reject) => {
const parsed = new URL(url);
const payload = JSON.stringify(body);
const requestFn = parsed.protocol === "https:" ? httpsRequest : httpRequest;
const req = requestFn(
parsed,
{
method: "POST",
headers: {
"content-type": "application/json",
"content-length": Buffer.byteLength(payload)
},
timeout: timeoutMs
},
(response) => {
let text = "";
response.setEncoding("utf8");
response.on("data", (chunk) => {
text += chunk;
});
response.on("end", () => {
let json = null;
try {
json = text ? JSON.parse(text) : null;
} catch {
json = null;
}
resolve({
status: response.statusCode ?? 0,
json,
bodyPreview: json ? undefined : text.slice(0, 500)
});
});
}
);
req.on("timeout", () => {
req.destroy(new Error(`request timed out after ${timeoutMs}ms`));
});
req.on("error", reject);
req.write(payload);
req.end();
});
}