539 lines
16 KiB
JavaScript
539 lines
16 KiB
JavaScript
import { createHash, randomUUID } from "node:crypto";
|
|
import { createServer as createHttpServer } from "node:http";
|
|
import { createServer as createTcpServer } from "node:net";
|
|
|
|
import {
|
|
assertDeviceAdapterSession,
|
|
} from "../../../packages/device-adapter-runtime/src/index.mjs";
|
|
import {
|
|
DEVICE_ADAPTER_MESSAGE_SCHEMA,
|
|
DEVICE_DISCOVERY_SIGNAL_SCHEMA,
|
|
normalizeAdapterAcceptance,
|
|
normalizeAdapterMessage,
|
|
} from "../../../packages/device-protocol-contract/src/index.mjs";
|
|
|
|
export function createDeviceGatewayRuntime(options = {}) {
|
|
const config = normalizeConfig(options);
|
|
const sessions = new Map();
|
|
const sessionsByAddress = new Map();
|
|
const connectionWindows = new Map();
|
|
let totalAccepted = 0;
|
|
let totalRejected = 0;
|
|
let totalDiscoveries = 0;
|
|
let totalMessagesAccepted = 0;
|
|
let totalPackagesAcknowledged = 0;
|
|
let totalBufferedBytes = 0;
|
|
|
|
const tcpServer = createTcpServer((socket) => {
|
|
const remoteAddress = normalizeRemoteAddress(socket.remoteAddress);
|
|
if (
|
|
sessions.size >= config.maxConcurrentSessions
|
|
|| currentAddressSessions(remoteAddress) >= config.maxSessionsPerAddress
|
|
|| !consumeConnectionPermit(remoteAddress)
|
|
) {
|
|
totalRejected += 1;
|
|
socket.destroy();
|
|
return;
|
|
}
|
|
|
|
const sessionRef = `session:${randomUUID()}`;
|
|
const session = {
|
|
sessionRef,
|
|
remoteAddress,
|
|
buffer: Buffer.alloc(0),
|
|
adapterSession: assertDeviceAdapterSession(
|
|
config.adapter.createSession({
|
|
profileRef: config.profile.profileRef,
|
|
}),
|
|
),
|
|
identifier: null,
|
|
sequence: 0,
|
|
state: "awaiting-header",
|
|
processing: false,
|
|
rejected: false,
|
|
closed: false,
|
|
};
|
|
sessions.set(socket, session);
|
|
incrementAddressSessions(remoteAddress);
|
|
totalAccepted += 1;
|
|
socket.setNoDelay(true);
|
|
socket.setTimeout(config.sessionTimeoutMs);
|
|
|
|
socket.on("data", (chunk) => {
|
|
if (session.closed || session.rejected) return;
|
|
socket.pause();
|
|
appendBuffer(session, chunk);
|
|
if (
|
|
session.buffer.length > config.maxBufferedBytes
|
|
|| totalBufferedBytes > config.maxAggregateBufferedBytes
|
|
) {
|
|
rejectSession(socket, session);
|
|
return;
|
|
}
|
|
if (session.processing) return;
|
|
session.processing = true;
|
|
void processSession(socket, session)
|
|
.catch(() => rejectSession(socket, session))
|
|
.finally(() => {
|
|
session.processing = false;
|
|
if (!session.closed && !session.rejected) socket.resume();
|
|
});
|
|
});
|
|
socket.on("timeout", () => rejectSession(socket, session));
|
|
socket.on("close", () => closeSession(socket, session));
|
|
socket.on("error", () => closeSession(socket, session));
|
|
});
|
|
|
|
const healthServer = createHttpServer((request, response) => {
|
|
response.setHeader("Content-Type", "application/json; charset=utf-8");
|
|
response.setHeader("Cache-Control", "no-store");
|
|
response.setHeader("X-Content-Type-Options", "nosniff");
|
|
if (request.method !== "GET" || request.url !== "/healthz") {
|
|
response.statusCode = 404;
|
|
return response.end('{"ok":false,"error":"device_gateway_route_not_found"}\n');
|
|
}
|
|
response.statusCode = 200;
|
|
return response.end(`${JSON.stringify({
|
|
ok: true,
|
|
service: "nodedc-device-gateway",
|
|
adapter: config.adapter?.adapterRef ?? "disabled",
|
|
protocolProfile: config.profile?.profileRef ?? "disabled",
|
|
framing: config.profile?.framing?.status ?? "disabled",
|
|
tcpListener: config.listenEnabled ? "telemetry-ingest" : "disabled",
|
|
publicIngress: config.publicIngressEnabled
|
|
? "telemetry-ingest"
|
|
: "disabled",
|
|
commandTransport: "disabled",
|
|
sessions: {
|
|
active: sessions.size,
|
|
accepted: totalAccepted,
|
|
rejected: totalRejected,
|
|
discoveries: totalDiscoveries,
|
|
messagesAccepted: totalMessagesAccepted,
|
|
packagesAcknowledged: totalPackagesAcknowledged,
|
|
bufferedBytes: totalBufferedBytes,
|
|
},
|
|
})}\n`);
|
|
});
|
|
|
|
return {
|
|
async start() {
|
|
await listen(healthServer, config.healthPort, config.healthHost);
|
|
if (config.listenEnabled) {
|
|
await listen(tcpServer, config.tcpPort, config.tcpHost);
|
|
}
|
|
return {
|
|
healthAddress: healthServer.address(),
|
|
tcpAddress: config.listenEnabled ? tcpServer.address() : null,
|
|
};
|
|
},
|
|
async stop() {
|
|
for (const [socket, session] of sessions) {
|
|
closeSession(socket, session);
|
|
socket.destroy();
|
|
}
|
|
await Promise.all([
|
|
closeServer(healthServer),
|
|
config.listenEnabled ? closeServer(tcpServer) : Promise.resolve(),
|
|
]);
|
|
},
|
|
status() {
|
|
return {
|
|
activeSessions: sessions.size,
|
|
totalAccepted,
|
|
totalRejected,
|
|
totalDiscoveries,
|
|
totalMessagesAccepted,
|
|
totalPackagesAcknowledged,
|
|
totalBufferedBytes,
|
|
commandTransport: "disabled",
|
|
publicIngress: config.publicIngressEnabled
|
|
? "telemetry-ingest"
|
|
: "disabled",
|
|
};
|
|
},
|
|
};
|
|
|
|
async function processSession(socket, session) {
|
|
while (!session.closed && !session.rejected) {
|
|
if (session.state === "awaiting-header") {
|
|
const parsed = session.adapterSession.parseHeader(session.buffer);
|
|
if (parsed.status === "incomplete") return;
|
|
|
|
const observedAt = config.now().toISOString();
|
|
const discovery = {
|
|
schemaVersion: DEVICE_DISCOVERY_SIGNAL_SCHEMA,
|
|
sessionRef: session.sessionRef,
|
|
...(config.routeRef ? { routeRef: config.routeRef } : {}),
|
|
modelProfileRef: config.profile.profileRef,
|
|
protocol: config.profile.protocol,
|
|
observedAt,
|
|
identifier: {
|
|
kind: parsed.identifier.kind,
|
|
value: parsed.identifier.value,
|
|
},
|
|
evidence: parsed.evidence,
|
|
};
|
|
const acceptedDiscovery = await config.onDiscovery(discovery);
|
|
assertDiscoveryAccepted(acceptedDiscovery);
|
|
session.identifier = discovery.identifier;
|
|
consumeBuffer(session, parsed.bytesConsumed);
|
|
session.state = "packages";
|
|
totalDiscoveries += 1;
|
|
await writeWithBackpressure(socket, session.adapterSession.buildHeaderAcknowledgement(
|
|
Math.floor(new Date(observedAt).getTime() / 1000),
|
|
));
|
|
continue;
|
|
}
|
|
|
|
const parsed = session.adapterSession.parseMessage(session.buffer);
|
|
if (parsed.status === "incomplete") return;
|
|
const observedAt = config.now().toISOString();
|
|
session.sequence += 1;
|
|
const message = normalizeAdapterMessage({
|
|
schemaVersion: DEVICE_ADAPTER_MESSAGE_SCHEMA,
|
|
edgeRef: config.edgeRef,
|
|
adapterRef: config.adapter.adapterRef,
|
|
protocolProfileRef: config.profile.profileRef,
|
|
protocol: config.profile.protocol,
|
|
sessionRef: session.sessionRef,
|
|
...(config.routeRef ? { routeRef: config.routeRef } : {}),
|
|
messageRef: messageRef(parsed),
|
|
messageType: parsed.messageType,
|
|
sequence: session.sequence,
|
|
observedAt,
|
|
idempotencyKey: messageIdempotencyKey({
|
|
adapterRef: config.adapter.adapterRef,
|
|
profileRef: config.profile.profileRef,
|
|
identifier: session.identifier,
|
|
payload: parsed.payload,
|
|
}),
|
|
identifier: session.identifier,
|
|
payloadSchemaRef: parsed.payloadSchemaRef,
|
|
payload: parsed.payload,
|
|
});
|
|
const acceptance = normalizeAdapterAcceptance(
|
|
await config.onMessage(message),
|
|
);
|
|
if (acceptance.idempotencyKey !== message.idempotencyKey) {
|
|
throw new TypeError("device_gateway_core_acceptance_mismatch");
|
|
}
|
|
consumeBuffer(session, parsed.bytesConsumed);
|
|
totalMessagesAccepted += 1;
|
|
totalPackagesAcknowledged += 1;
|
|
await writeWithBackpressure(
|
|
socket,
|
|
session.adapterSession.buildMessageAcknowledgement(parsed),
|
|
);
|
|
}
|
|
}
|
|
|
|
function currentAddressSessions(remoteAddress) {
|
|
return sessionsByAddress.get(remoteAddress) || 0;
|
|
}
|
|
|
|
function incrementAddressSessions(remoteAddress) {
|
|
sessionsByAddress.set(
|
|
remoteAddress,
|
|
currentAddressSessions(remoteAddress) + 1,
|
|
);
|
|
}
|
|
|
|
function decrementAddressSessions(remoteAddress) {
|
|
const current = currentAddressSessions(remoteAddress);
|
|
if (current <= 1) {
|
|
sessionsByAddress.delete(remoteAddress);
|
|
} else {
|
|
sessionsByAddress.set(remoteAddress, current - 1);
|
|
}
|
|
}
|
|
|
|
function consumeConnectionPermit(remoteAddress) {
|
|
const nowMs = config.now().getTime();
|
|
for (const [address, window] of connectionWindows) {
|
|
if (nowMs - window.startedAt >= 60_000) {
|
|
connectionWindows.delete(address);
|
|
}
|
|
}
|
|
const current = connectionWindows.get(remoteAddress);
|
|
if (!current || nowMs - current.startedAt >= 60_000) {
|
|
if (connectionWindows.size >= config.maxTrackedSourceAddresses) {
|
|
return false;
|
|
}
|
|
connectionWindows.set(remoteAddress, {
|
|
startedAt: nowMs,
|
|
count: 1,
|
|
});
|
|
return true;
|
|
}
|
|
if (current.count >= config.maxConnectionsPerMinutePerAddress) {
|
|
return false;
|
|
}
|
|
current.count += 1;
|
|
return true;
|
|
}
|
|
|
|
function rejectSession(socket, session) {
|
|
if (!session.rejected) {
|
|
session.rejected = true;
|
|
totalRejected += 1;
|
|
}
|
|
socket.destroy();
|
|
}
|
|
|
|
function closeSession(socket, session) {
|
|
if (session.closed) return;
|
|
session.closed = true;
|
|
totalBufferedBytes -= session.buffer.length;
|
|
session.buffer = Buffer.alloc(0);
|
|
sessions.delete(socket);
|
|
decrementAddressSessions(session.remoteAddress);
|
|
}
|
|
|
|
function appendBuffer(session, chunk) {
|
|
session.buffer = Buffer.concat(
|
|
[session.buffer, chunk],
|
|
session.buffer.length + chunk.length,
|
|
);
|
|
totalBufferedBytes += chunk.length;
|
|
}
|
|
|
|
function consumeBuffer(session, bytesConsumed) {
|
|
session.buffer = session.buffer.subarray(bytesConsumed);
|
|
totalBufferedBytes -= bytesConsumed;
|
|
}
|
|
}
|
|
|
|
function normalizeConfig(input) {
|
|
const listenEnabled = input.listenEnabled === true;
|
|
const publicIngressEnabled = input.publicIngressEnabled === true;
|
|
if (publicIngressEnabled && !listenEnabled) {
|
|
throw new TypeError("device_gateway_public_ingress_listener_required");
|
|
}
|
|
if (publicIngressEnabled && input.coreChannelAuthenticated !== true) {
|
|
throw new TypeError("device_gateway_authenticated_core_channel_required");
|
|
}
|
|
if (
|
|
listenEnabled
|
|
&& (typeof input.onDiscovery !== "function"
|
|
|| typeof input.onMessage !== "function")
|
|
) {
|
|
throw new TypeError("device_gateway_core_acceptance_sink_required");
|
|
}
|
|
const registration = listenEnabled
|
|
? resolveAdapterRegistration(input.adapterRegistry, input.protocolProfileRef)
|
|
: null;
|
|
return {
|
|
listenEnabled,
|
|
publicIngressEnabled,
|
|
adapter: registration?.adapter ?? null,
|
|
profile: registration?.profile ?? null,
|
|
edgeRef: listenEnabled
|
|
? normalizeOpaqueRef(input.edgeRef, "device_gateway_edge_ref_invalid")
|
|
: "edge:disabled",
|
|
routeRef: input.routeRef == null || input.routeRef === ""
|
|
? undefined
|
|
: normalizeRouteRef(input.routeRef),
|
|
healthHost: normalizeHealthHost(input.healthHost, "127.0.0.1"),
|
|
healthPort: parseInteger(
|
|
input.healthPort,
|
|
18121,
|
|
0,
|
|
65535,
|
|
"device_gateway_health_port_invalid",
|
|
),
|
|
tcpHost: normalizeTcpHost(
|
|
input.tcpHost,
|
|
publicIngressEnabled ? "0.0.0.0" : "127.0.0.1",
|
|
publicIngressEnabled,
|
|
),
|
|
tcpPort: parseInteger(
|
|
input.tcpPort,
|
|
9921,
|
|
0,
|
|
65535,
|
|
"device_gateway_tcp_port_invalid",
|
|
),
|
|
maxBufferedBytes: parseInteger(
|
|
input.maxBufferedBytes,
|
|
registration?.profile?.framing?.maxBufferedBytes ?? 256 * 1024,
|
|
1024,
|
|
registration?.profile?.framing?.maxBufferedBytes ?? 256 * 1024,
|
|
"device_gateway_buffer_limit_invalid",
|
|
),
|
|
maxAggregateBufferedBytes: parseInteger(
|
|
input.maxAggregateBufferedBytes,
|
|
32 * 1024 * 1024,
|
|
1024,
|
|
32 * 1024 * 1024,
|
|
"device_gateway_aggregate_buffer_limit_invalid",
|
|
),
|
|
maxConcurrentSessions: parseInteger(
|
|
input.maxConcurrentSessions,
|
|
128,
|
|
1,
|
|
10000,
|
|
"device_gateway_session_limit_invalid",
|
|
),
|
|
maxSessionsPerAddress: parseInteger(
|
|
input.maxSessionsPerAddress,
|
|
16,
|
|
1,
|
|
1000,
|
|
"device_gateway_address_session_limit_invalid",
|
|
),
|
|
maxConnectionsPerMinutePerAddress: parseInteger(
|
|
input.maxConnectionsPerMinutePerAddress,
|
|
60,
|
|
1,
|
|
10000,
|
|
"device_gateway_address_rate_limit_invalid",
|
|
),
|
|
maxTrackedSourceAddresses: parseInteger(
|
|
input.maxTrackedSourceAddresses,
|
|
2048,
|
|
1,
|
|
65536,
|
|
"device_gateway_source_tracking_limit_invalid",
|
|
),
|
|
sessionTimeoutMs: parseInteger(
|
|
input.sessionTimeoutMs,
|
|
10000,
|
|
100,
|
|
60000,
|
|
"device_gateway_session_timeout_invalid",
|
|
),
|
|
onDiscovery: typeof input.onDiscovery === "function"
|
|
? input.onDiscovery
|
|
: undefined,
|
|
onMessage: typeof input.onMessage === "function"
|
|
? input.onMessage
|
|
: undefined,
|
|
now: typeof input.now === "function" ? input.now : () => new Date(),
|
|
};
|
|
}
|
|
|
|
function resolveAdapterRegistration(registry, profileRef) {
|
|
if (!registry || typeof registry.resolveProfile !== "function") {
|
|
throw new TypeError("device_gateway_adapter_registry_required");
|
|
}
|
|
return registry.resolveProfile(profileRef);
|
|
}
|
|
|
|
function normalizeOpaqueRef(value, errorCode) {
|
|
if (
|
|
typeof value !== "string"
|
|
|| !/^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$/.test(value)
|
|
) {
|
|
throw new TypeError(errorCode);
|
|
}
|
|
return value;
|
|
}
|
|
|
|
function normalizeRouteRef(value) {
|
|
if (
|
|
typeof value !== "string"
|
|
|| !/^route:[0-9a-f]{8}-[0-9a-f]{4}-[1-5][0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/i.test(value)
|
|
) {
|
|
throw new TypeError("device_gateway_route_ref_invalid");
|
|
}
|
|
return value.toLowerCase();
|
|
}
|
|
|
|
function assertDiscoveryAccepted(value) {
|
|
if (
|
|
!value
|
|
|| typeof value !== "object"
|
|
|| !["quarantine", "claimed"].includes(value.lifecycleState)
|
|
) {
|
|
throw new TypeError("device_gateway_core_discovery_not_accepted");
|
|
}
|
|
}
|
|
|
|
function messageRef(parsed) {
|
|
const digest = String(parsed?.payload?.packageDigest || "");
|
|
if (!/^sha256:[a-f0-9]{64}$/.test(digest)) {
|
|
throw new TypeError("device_gateway_adapter_message_digest_invalid");
|
|
}
|
|
return `package:${parsed.packageNumber}:${digest.slice("sha256:".length)}`;
|
|
}
|
|
|
|
function messageIdempotencyKey({ adapterRef, profileRef, identifier, payload }) {
|
|
return `sha256:${createHash("sha256")
|
|
.update(JSON.stringify({ adapterRef, profileRef, identifier, payload }), "utf8")
|
|
.digest("hex")}`;
|
|
}
|
|
|
|
function writeWithBackpressure(socket, bytes) {
|
|
if (!Buffer.isBuffer(bytes) || bytes.length === 0 || bytes.length > 4096) {
|
|
throw new TypeError("device_gateway_adapter_ack_invalid");
|
|
}
|
|
if (socket.write(bytes)) return Promise.resolve();
|
|
return new Promise((resolve, reject) => {
|
|
const onDrain = () => {
|
|
socket.off("error", onError);
|
|
resolve();
|
|
};
|
|
const onError = (error) => {
|
|
socket.off("drain", onDrain);
|
|
reject(error);
|
|
};
|
|
socket.once("drain", onDrain);
|
|
socket.once("error", onError);
|
|
});
|
|
}
|
|
|
|
function normalizeHealthHost(value, fallback) {
|
|
const normalized = String(value || fallback).trim();
|
|
if (!["127.0.0.1", "::1", "0.0.0.0", "::"].includes(normalized)) {
|
|
throw new TypeError("device_gateway_health_host_invalid");
|
|
}
|
|
return normalized;
|
|
}
|
|
|
|
function normalizeTcpHost(value, fallback, publicIngressEnabled) {
|
|
const normalized = String(value || fallback).trim();
|
|
const allowed = publicIngressEnabled
|
|
? ["0.0.0.0", "::"]
|
|
: ["127.0.0.1", "::1"];
|
|
if (!allowed.includes(normalized)) {
|
|
throw new TypeError(
|
|
publicIngressEnabled
|
|
? "device_gateway_public_ingress_host_invalid"
|
|
: "device_gateway_baseline_loopback_only",
|
|
);
|
|
}
|
|
return normalized;
|
|
}
|
|
|
|
function normalizeRemoteAddress(value) {
|
|
const normalized = String(value || "unknown").trim();
|
|
return normalized.slice(0, 64) || "unknown";
|
|
}
|
|
|
|
function parseInteger(value, fallback, minimum, maximum, errorCode) {
|
|
const parsed = Number(value ?? fallback);
|
|
if (
|
|
!Number.isSafeInteger(parsed)
|
|
|| parsed < minimum
|
|
|| parsed > maximum
|
|
) {
|
|
throw new TypeError(errorCode);
|
|
}
|
|
return parsed;
|
|
}
|
|
|
|
function listen(server, port, host) {
|
|
return new Promise((resolve, reject) => {
|
|
server.once("error", reject);
|
|
server.listen(port, host, resolve);
|
|
});
|
|
}
|
|
|
|
function closeServer(server) {
|
|
if (!server.listening) return Promise.resolve();
|
|
return new Promise((resolve, reject) => {
|
|
server.close((error) => (error ? reject(error) : resolve()));
|
|
});
|
|
}
|