feat: establish standalone Device Core repository
This commit is contained in:
@@ -0,0 +1,18 @@
|
||||
FROM node:22-alpine
|
||||
|
||||
WORKDIR /app
|
||||
|
||||
COPY package.json package-lock.json ./
|
||||
COPY packages/device-protocol-contract ./packages/device-protocol-contract
|
||||
COPY packages/device-adapter-runtime ./packages/device-adapter-runtime
|
||||
COPY packages/device-adapter-catalog ./packages/device-adapter-catalog
|
||||
COPY packages/arusnavi-b2-adapter ./packages/arusnavi-b2-adapter
|
||||
COPY services/device-gateway ./services/device-gateway
|
||||
COPY services/device-control-core/package.json ./services/device-control-core/package.json
|
||||
COPY services/device-edge-relay/package.json ./services/device-edge-relay/package.json
|
||||
|
||||
RUN npm ci --omit=dev --ignore-scripts
|
||||
|
||||
USER node
|
||||
|
||||
CMD ["node", "services/device-gateway/src/server.mjs"]
|
||||
@@ -0,0 +1,13 @@
|
||||
{
|
||||
"name": "@nodedc/device-gateway",
|
||||
"version": "0.1.0",
|
||||
"private": true,
|
||||
"type": "module",
|
||||
"scripts": {
|
||||
"start": "node src/server.mjs",
|
||||
"test": "node --test test/*.test.mjs"
|
||||
},
|
||||
"engines": {
|
||||
"node": ">=20"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,147 @@
|
||||
import {
|
||||
normalizeAdapterAcceptance,
|
||||
} from "../../../packages/device-protocol-contract/src/index.mjs";
|
||||
|
||||
export function createCoreGatewayClient(options = {}) {
|
||||
const {
|
||||
coreUrl,
|
||||
gatewayToken,
|
||||
timeoutMs = 5000,
|
||||
fetchImpl = fetch,
|
||||
} = options;
|
||||
const baseUrl = normalizeCoreBaseUrl(coreUrl);
|
||||
validateClientOptions({ gatewayToken, timeoutMs, fetchImpl });
|
||||
const discoveryEndpoint = endpointUrl(
|
||||
baseUrl,
|
||||
"/internal/v1/device-discoveries:observe",
|
||||
);
|
||||
const messageEndpoint = endpointUrl(
|
||||
baseUrl,
|
||||
"/internal/v1/gateway/messages:accept",
|
||||
);
|
||||
|
||||
return Object.freeze({
|
||||
async observeDiscovery(signal) {
|
||||
const body = await postJson({
|
||||
endpoint: discoveryEndpoint,
|
||||
gatewayToken,
|
||||
timeoutMs,
|
||||
fetchImpl,
|
||||
value: signal,
|
||||
maxResponseBytes: 32 * 1024,
|
||||
});
|
||||
if (
|
||||
!body.discovery
|
||||
|| !["quarantine", "claimed"].includes(body.discovery.lifecycleState)
|
||||
|| body.discovery.commandTransport !== "disabled"
|
||||
) {
|
||||
throw new Error("device_gateway_core_ingest_contract_invalid");
|
||||
}
|
||||
return body.discovery;
|
||||
},
|
||||
async acceptMessage(message) {
|
||||
const body = await postJson({
|
||||
endpoint: messageEndpoint,
|
||||
gatewayToken,
|
||||
timeoutMs,
|
||||
fetchImpl,
|
||||
value: message,
|
||||
maxResponseBytes: 16 * 1024,
|
||||
});
|
||||
const acceptance = normalizeAdapterAcceptance(body.acceptance);
|
||||
if (acceptance.idempotencyKey !== message?.idempotencyKey) {
|
||||
throw new Error("device_gateway_core_acceptance_mismatch");
|
||||
}
|
||||
return acceptance;
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
export function createCoreDiscoveryClient({
|
||||
coreUrl,
|
||||
gatewayToken,
|
||||
timeoutMs = 5000,
|
||||
fetchImpl = fetch,
|
||||
} = {}) {
|
||||
return createCoreGatewayClient({
|
||||
coreUrl,
|
||||
gatewayToken,
|
||||
timeoutMs,
|
||||
fetchImpl,
|
||||
}).observeDiscovery;
|
||||
}
|
||||
|
||||
function validateClientOptions({ gatewayToken, timeoutMs, fetchImpl }) {
|
||||
if (typeof gatewayToken !== "string" || gatewayToken.length < 32) {
|
||||
throw new TypeError("device_gateway_core_token_invalid");
|
||||
}
|
||||
const normalizedTimeout = Number(timeoutMs);
|
||||
if (
|
||||
!Number.isSafeInteger(normalizedTimeout)
|
||||
|| normalizedTimeout < 100
|
||||
|| normalizedTimeout > 30_000
|
||||
) {
|
||||
throw new TypeError("device_gateway_core_timeout_invalid");
|
||||
}
|
||||
if (typeof fetchImpl !== "function") {
|
||||
throw new TypeError("device_gateway_core_fetch_invalid");
|
||||
}
|
||||
}
|
||||
|
||||
async function postJson({
|
||||
endpoint,
|
||||
gatewayToken,
|
||||
timeoutMs,
|
||||
fetchImpl,
|
||||
value,
|
||||
maxResponseBytes,
|
||||
}) {
|
||||
const response = await fetchImpl(endpoint, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
Authorization: `Bearer ${gatewayToken}`,
|
||||
"Content-Type": "application/json",
|
||||
},
|
||||
body: JSON.stringify(value),
|
||||
signal: AbortSignal.timeout(Number(timeoutMs)),
|
||||
});
|
||||
const body = await readBoundedJson(response, maxResponseBytes);
|
||||
if (!response.ok || body?.ok !== true) {
|
||||
throw new Error("device_gateway_core_ingest_failed");
|
||||
}
|
||||
return body;
|
||||
}
|
||||
|
||||
function normalizeCoreBaseUrl(value) {
|
||||
let url;
|
||||
try {
|
||||
url = new URL(String(value || ""));
|
||||
} catch {
|
||||
throw new TypeError("device_gateway_core_url_invalid");
|
||||
}
|
||||
if (url.protocol !== "http:" || url.username || url.password) {
|
||||
throw new TypeError("device_gateway_core_url_invalid");
|
||||
}
|
||||
if (url.pathname !== "/" || url.search || url.hash) {
|
||||
throw new TypeError("device_gateway_core_url_invalid");
|
||||
}
|
||||
return url;
|
||||
}
|
||||
|
||||
function endpointUrl(baseUrl, pathname) {
|
||||
const url = new URL(baseUrl);
|
||||
url.pathname = pathname;
|
||||
return url.toString();
|
||||
}
|
||||
|
||||
async function readBoundedJson(response, maxBytes) {
|
||||
const text = await response.text();
|
||||
if (Buffer.byteLength(text, "utf8") > maxBytes) {
|
||||
throw new Error("device_gateway_core_response_too_large");
|
||||
}
|
||||
try {
|
||||
return JSON.parse(text);
|
||||
} catch {
|
||||
throw new Error("device_gateway_core_response_invalid");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,646 @@
|
||||
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,
|
||||
pendingCommand: null,
|
||||
};
|
||||
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: config.commandTransport,
|
||||
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 {
|
||||
adapter: config.adapter?.adapterRef ?? "disabled",
|
||||
protocolProfile: config.profile?.profileRef ?? "disabled",
|
||||
framing: config.profile?.framing?.status ?? "disabled",
|
||||
activeSessions: sessions.size,
|
||||
totalAccepted,
|
||||
totalRejected,
|
||||
totalDiscoveries,
|
||||
totalMessagesAccepted,
|
||||
totalPackagesAcknowledged,
|
||||
totalBufferedBytes,
|
||||
commandTransport: config.commandTransport,
|
||||
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),
|
||||
));
|
||||
if (acceptedDiscovery.commandOffer) {
|
||||
await dispatchCommandOffer(
|
||||
socket,
|
||||
session,
|
||||
acceptedDiscovery.commandOffer,
|
||||
);
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
||||
if (session.pendingCommand) {
|
||||
const result = session.adapterSession.parseTypedCommandResponse(
|
||||
session.buffer,
|
||||
session.pendingCommand,
|
||||
);
|
||||
if (result.status === "incomplete") return;
|
||||
if (result.status === "acknowledged") {
|
||||
const command = session.pendingCommand;
|
||||
consumeBuffer(session, result.bytesConsumed);
|
||||
session.pendingCommand = null;
|
||||
await config.onCommandStatus({
|
||||
commandRef: command.commandRef,
|
||||
transportMessageRef: command.transportMessageRef,
|
||||
lifecycleState: "acknowledged",
|
||||
resultCode: result.resultCode,
|
||||
observedAt: config.now().toISOString(),
|
||||
sessionRef: session.sessionRef,
|
||||
adapterProfileRef: config.profile.profileRef,
|
||||
});
|
||||
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 acceptedMessage = await config.onMessage(message);
|
||||
const {
|
||||
commandOffer,
|
||||
...acceptanceValue
|
||||
} = acceptedMessage;
|
||||
const acceptance = normalizeAdapterAcceptance(acceptanceValue);
|
||||
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),
|
||||
);
|
||||
if (commandOffer) {
|
||||
await dispatchCommandOffer(
|
||||
socket,
|
||||
session,
|
||||
commandOffer,
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async function dispatchCommandOffer(socket, session, value) {
|
||||
if (session.pendingCommand) {
|
||||
throw new Error("device_gateway_command_already_pending");
|
||||
}
|
||||
const command = normalizeCommandOffer(value);
|
||||
const bytes = session.adapterSession.buildTypedCommand(command);
|
||||
session.pendingCommand = command;
|
||||
await writeWithBackpressure(socket, bytes);
|
||||
}
|
||||
|
||||
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 = Math.max(0, totalBufferedBytes - session.buffer.length);
|
||||
session.buffer = Buffer.alloc(0);
|
||||
if (session.pendingCommand && config.onCommandStatus) {
|
||||
const command = session.pendingCommand;
|
||||
session.pendingCommand = null;
|
||||
void config.onCommandStatus({
|
||||
commandRef: command.commandRef,
|
||||
transportMessageRef: command.transportMessageRef,
|
||||
lifecycleState: "unknown",
|
||||
resultCode: "tracker_session_closed",
|
||||
observedAt: config.now().toISOString(),
|
||||
sessionRef: session.sessionRef,
|
||||
adapterProfileRef: config.profile.profileRef,
|
||||
}).catch(() => undefined);
|
||||
}
|
||||
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) {
|
||||
const consumed = Math.min(bytesConsumed, session.buffer.length);
|
||||
session.buffer = session.buffer.subarray(consumed);
|
||||
totalBufferedBytes = Math.max(0, totalBufferedBytes - consumed);
|
||||
}
|
||||
}
|
||||
|
||||
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;
|
||||
if (registration?.profile?.commandTransport?.status === "typed-service-ping-v1") {
|
||||
const probe = assertDeviceAdapterSession(registration.adapter.createSession({
|
||||
profileRef: registration.profile.profileRef,
|
||||
}));
|
||||
if (
|
||||
typeof probe.buildTypedCommand !== "function"
|
||||
|| typeof probe.parseTypedCommandResponse !== "function"
|
||||
|| typeof input.onCommandStatus !== "function"
|
||||
) {
|
||||
throw new TypeError("device_gateway_typed_command_runtime_required");
|
||||
}
|
||||
}
|
||||
return {
|
||||
listenEnabled,
|
||||
publicIngressEnabled,
|
||||
adapter: registration?.adapter ?? null,
|
||||
profile: registration?.profile ?? null,
|
||||
commandTransport:
|
||||
registration?.profile?.commandTransport?.status ?? "disabled",
|
||||
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,
|
||||
onCommandStatus: typeof input.onCommandStatus === "function"
|
||||
? input.onCommandStatus
|
||||
: undefined,
|
||||
now: typeof input.now === "function" ? input.now : () => new Date(),
|
||||
};
|
||||
}
|
||||
|
||||
function normalizeCommandOffer(value) {
|
||||
if (
|
||||
!value
|
||||
|| typeof value !== "object"
|
||||
|| Array.isArray(value)
|
||||
|| typeof value.commandRef !== "string"
|
||||
|| !/^command:[0-9a-f-]{36}$/i.test(value.commandRef)
|
||||
|| value.commandType !== "service.ping"
|
||||
|| typeof value.accessCode !== "string"
|
||||
|| !/^\d{6}$/.test(value.accessCode)
|
||||
|| typeof value.transportMessageRef !== "string"
|
||||
|| !/^edge-command:[0-9a-f-]{36}$/i.test(value.transportMessageRef)
|
||||
) {
|
||||
throw new TypeError("device_gateway_command_offer_invalid");
|
||||
}
|
||||
return Object.freeze({
|
||||
commandRef: value.commandRef.toLowerCase(),
|
||||
commandType: value.commandType,
|
||||
accessCode: value.accessCode,
|
||||
transportMessageRef: value.transportMessageRef.toLowerCase(),
|
||||
});
|
||||
}
|
||||
|
||||
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()));
|
||||
});
|
||||
}
|
||||
@@ -0,0 +1,175 @@
|
||||
import { readFile } from "node:fs/promises";
|
||||
|
||||
import {
|
||||
DEVICE_ADAPTER_CATALOG,
|
||||
} from "../../../packages/device-adapter-catalog/src/index.mjs";
|
||||
import { createCoreGatewayClient } from "./core-client.mjs";
|
||||
import { createDeviceGatewayRuntime } from "./runtime.mjs";
|
||||
|
||||
const config = await readConfig();
|
||||
const coreClient = config.listenEnabled
|
||||
? createCoreGatewayClient({
|
||||
coreUrl: config.coreUrl,
|
||||
gatewayToken: config.gatewayToken,
|
||||
timeoutMs: config.coreTimeoutMs,
|
||||
})
|
||||
: undefined;
|
||||
const runtime = createDeviceGatewayRuntime({
|
||||
listenEnabled: config.listenEnabled,
|
||||
publicIngressEnabled: config.publicIngressEnabled,
|
||||
coreChannelAuthenticated: false,
|
||||
adapterRegistry: DEVICE_ADAPTER_CATALOG.registry,
|
||||
protocolProfileRef: config.protocolProfileRef,
|
||||
edgeRef: config.edgeRef,
|
||||
routeRef: config.routeRef,
|
||||
healthHost: config.healthHost,
|
||||
healthPort: config.healthPort,
|
||||
tcpHost: config.tcpHost,
|
||||
tcpPort: config.tcpPort,
|
||||
maxBufferedBytes: config.maxBufferedBytes,
|
||||
maxAggregateBufferedBytes: config.maxAggregateBufferedBytes,
|
||||
maxConcurrentSessions: config.maxConcurrentSessions,
|
||||
maxSessionsPerAddress: config.maxSessionsPerAddress,
|
||||
maxConnectionsPerMinutePerAddress:
|
||||
config.maxConnectionsPerMinutePerAddress,
|
||||
maxTrackedSourceAddresses: config.maxTrackedSourceAddresses,
|
||||
sessionTimeoutMs: config.sessionTimeoutMs,
|
||||
onDiscovery: coreClient?.observeDiscovery,
|
||||
onMessage: coreClient?.acceptMessage,
|
||||
});
|
||||
|
||||
const addresses = await runtime.start();
|
||||
console.log(JSON.stringify({
|
||||
event: "device_gateway_started",
|
||||
health: addresses.healthAddress,
|
||||
tcp: addresses.tcpAddress,
|
||||
publicIngress: config.publicIngressEnabled
|
||||
? "telemetry-ingest"
|
||||
: "disabled",
|
||||
commandTransport: "disabled",
|
||||
}));
|
||||
|
||||
process.on("SIGTERM", shutdown);
|
||||
process.on("SIGINT", shutdown);
|
||||
|
||||
async function shutdown() {
|
||||
await runtime.stop();
|
||||
process.exit(0);
|
||||
}
|
||||
|
||||
async function readConfig() {
|
||||
const listenEnabled = parseBoolean(
|
||||
process.env.DEVICE_GATEWAY_LISTEN_ENABLED,
|
||||
false,
|
||||
);
|
||||
const publicIngressEnabled = parseBoolean(
|
||||
process.env.DEVICE_GATEWAY_PUBLIC_INGRESS_ENABLED,
|
||||
false,
|
||||
);
|
||||
return {
|
||||
listenEnabled,
|
||||
publicIngressEnabled,
|
||||
protocolProfileRef: String(
|
||||
process.env.DEVICE_GATEWAY_PROTOCOL_PROFILE_REF
|
||||
|| DEVICE_ADAPTER_CATALOG.defaultProfileRef,
|
||||
),
|
||||
edgeRef: listenEnabled
|
||||
? requiredValue(
|
||||
process.env.DEVICE_GATEWAY_EDGE_REF,
|
||||
"device_gateway_edge_ref_required",
|
||||
)
|
||||
: "edge:disabled",
|
||||
routeRef: String(process.env.DEVICE_GATEWAY_ROUTE_REF || ""),
|
||||
healthHost: String(
|
||||
process.env.DEVICE_GATEWAY_HEALTH_HOST || "127.0.0.1",
|
||||
),
|
||||
healthPort: parsePort(process.env.DEVICE_GATEWAY_HEALTH_PORT, 18121),
|
||||
tcpHost: String(
|
||||
process.env.DEVICE_GATEWAY_TCP_HOST
|
||||
|| (publicIngressEnabled ? "0.0.0.0" : "127.0.0.1"),
|
||||
),
|
||||
tcpPort: parsePort(process.env.DEVICE_GATEWAY_TCP_PORT, 9921),
|
||||
maxBufferedBytes: parsePositiveInt(
|
||||
process.env.DEVICE_GATEWAY_MAX_BUFFERED_BYTES,
|
||||
65536,
|
||||
),
|
||||
maxAggregateBufferedBytes: parsePositiveInt(
|
||||
process.env.DEVICE_GATEWAY_MAX_AGGREGATE_BUFFERED_BYTES,
|
||||
32 * 1024 * 1024,
|
||||
),
|
||||
maxConcurrentSessions: parsePositiveInt(
|
||||
process.env.DEVICE_GATEWAY_MAX_SESSIONS,
|
||||
128,
|
||||
),
|
||||
maxSessionsPerAddress: parsePositiveInt(
|
||||
process.env.DEVICE_GATEWAY_MAX_SESSIONS_PER_ADDRESS,
|
||||
16,
|
||||
),
|
||||
maxConnectionsPerMinutePerAddress: parsePositiveInt(
|
||||
process.env.DEVICE_GATEWAY_MAX_CONNECTIONS_PER_MINUTE_PER_ADDRESS,
|
||||
60,
|
||||
),
|
||||
maxTrackedSourceAddresses: parsePositiveInt(
|
||||
process.env.DEVICE_GATEWAY_MAX_TRACKED_SOURCE_ADDRESSES,
|
||||
2048,
|
||||
),
|
||||
sessionTimeoutMs: parsePositiveInt(
|
||||
process.env.DEVICE_GATEWAY_SESSION_TIMEOUT_MS,
|
||||
10000,
|
||||
),
|
||||
coreUrl: listenEnabled
|
||||
? requiredValue(
|
||||
process.env.DEVICE_GATEWAY_CORE_URL,
|
||||
"device_gateway_core_url_required",
|
||||
)
|
||||
: "",
|
||||
gatewayToken: listenEnabled
|
||||
? await readRequiredSecretFile(
|
||||
process.env.DEVICE_GATEWAY_CORE_TOKEN_FILE,
|
||||
"device_gateway_core_token_file_required",
|
||||
)
|
||||
: "",
|
||||
coreTimeoutMs: parsePositiveInt(
|
||||
process.env.DEVICE_GATEWAY_CORE_TIMEOUT_MS,
|
||||
5000,
|
||||
),
|
||||
};
|
||||
}
|
||||
|
||||
async function readRequiredSecretFile(path, errorCode) {
|
||||
const normalized = requiredValue(path, errorCode);
|
||||
const value = (await readFile(normalized, "utf8")).trim();
|
||||
if (value.length < 32) throw new Error(errorCode);
|
||||
return value;
|
||||
}
|
||||
|
||||
function requiredValue(value, errorCode) {
|
||||
if (typeof value !== "string" || value.trim() === "") {
|
||||
throw new Error(errorCode);
|
||||
}
|
||||
return value.trim();
|
||||
}
|
||||
|
||||
function parsePort(value, fallback) {
|
||||
const parsed = Number(value || fallback);
|
||||
if (!Number.isSafeInteger(parsed) || parsed < 1 || parsed > 65535) {
|
||||
throw new Error("device_gateway_port_invalid");
|
||||
}
|
||||
return parsed;
|
||||
}
|
||||
|
||||
function parsePositiveInt(value, fallback) {
|
||||
const parsed = Number(value || fallback);
|
||||
if (!Number.isSafeInteger(parsed) || parsed < 1) {
|
||||
throw new Error("device_gateway_positive_integer_invalid");
|
||||
}
|
||||
return parsed;
|
||||
}
|
||||
|
||||
function parseBoolean(value, fallback) {
|
||||
if (value === undefined || value === null || value === "") return fallback;
|
||||
const normalized = String(value).trim().toLowerCase();
|
||||
if (["1", "true", "yes", "on"].includes(normalized)) return true;
|
||||
if (["0", "false", "no", "off"].includes(normalized)) return false;
|
||||
throw new Error("device_gateway_boolean_invalid");
|
||||
}
|
||||
@@ -0,0 +1,118 @@
|
||||
import assert from "node:assert/strict";
|
||||
import test from "node:test";
|
||||
|
||||
import {
|
||||
createCoreDiscoveryClient,
|
||||
createCoreGatewayClient,
|
||||
} from "../src/core-client.mjs";
|
||||
|
||||
const gatewayToken = "test-only-gateway-token-with-32-bytes";
|
||||
|
||||
test("posts a discovery through the authenticated internal Core boundary", async () => {
|
||||
let captured;
|
||||
const observe = createCoreDiscoveryClient({
|
||||
coreUrl: "http://device-control-core:18120",
|
||||
gatewayToken,
|
||||
fetchImpl: async (url, options) => {
|
||||
captured = { url, options };
|
||||
return new Response(JSON.stringify({
|
||||
ok: true,
|
||||
discovery: {
|
||||
lifecycleState: "quarantine",
|
||||
commandTransport: "disabled",
|
||||
},
|
||||
}), {
|
||||
status: 201,
|
||||
headers: { "Content-Type": "application/json" },
|
||||
});
|
||||
},
|
||||
});
|
||||
const signal = {
|
||||
schemaVersion: "nodedc.device.discovery-signal.v1",
|
||||
sessionRef: "session:test",
|
||||
};
|
||||
const discovery = await observe(signal);
|
||||
assert.equal(
|
||||
captured.url,
|
||||
"http://device-control-core:18120/internal/v1/device-discoveries:observe",
|
||||
);
|
||||
assert.equal(
|
||||
captured.options.headers.Authorization,
|
||||
`Bearer ${gatewayToken}`,
|
||||
);
|
||||
assert.deepEqual(JSON.parse(captured.options.body), signal);
|
||||
assert.equal(discovery.lifecycleState, "quarantine");
|
||||
});
|
||||
|
||||
test("fails closed when Core does not return an accepted discovery view", async () => {
|
||||
const observe = createCoreDiscoveryClient({
|
||||
coreUrl: "http://device-control-core:18120",
|
||||
gatewayToken,
|
||||
fetchImpl: async () => new Response(JSON.stringify({
|
||||
ok: true,
|
||||
discovery: {
|
||||
lifecycleState: "observed",
|
||||
commandTransport: "disabled",
|
||||
},
|
||||
}), { status: 200 }),
|
||||
});
|
||||
await assert.rejects(
|
||||
() => observe({ schemaVersion: "test" }),
|
||||
/device_gateway_core_ingest_contract_invalid/,
|
||||
);
|
||||
});
|
||||
|
||||
test("accepts a package only through the explicit Core acceptance contract", async () => {
|
||||
let captured;
|
||||
const client = createCoreGatewayClient({
|
||||
coreUrl: "http://device-control-core:18120",
|
||||
gatewayToken,
|
||||
fetchImpl: async (url, options) => {
|
||||
captured = { url, options };
|
||||
const message = JSON.parse(options.body);
|
||||
return new Response(JSON.stringify({
|
||||
ok: true,
|
||||
acceptance: {
|
||||
schemaVersion: "nodedc.device-adapter-acceptance.v1",
|
||||
acceptanceRef: "acceptance:test-001",
|
||||
idempotencyKey: message.idempotencyKey,
|
||||
status: "accepted",
|
||||
replayed: false,
|
||||
acceptedAt: "2026-08-11T12:00:00.000Z",
|
||||
},
|
||||
}), { status: 201 });
|
||||
},
|
||||
});
|
||||
const message = { idempotencyKey: `sha256:${"a".repeat(64)}` };
|
||||
const acceptance = await client.acceptMessage(message);
|
||||
assert.equal(
|
||||
captured.url,
|
||||
"http://device-control-core:18120/internal/v1/gateway/messages:accept",
|
||||
);
|
||||
assert.equal(acceptance.status, "accepted");
|
||||
assert.equal(acceptance.idempotencyKey, message.idempotencyKey);
|
||||
});
|
||||
|
||||
test("rejects a mismatched or non-durable Core package response", async () => {
|
||||
const client = createCoreGatewayClient({
|
||||
coreUrl: "http://device-control-core:18120",
|
||||
gatewayToken,
|
||||
fetchImpl: async () => new Response(JSON.stringify({
|
||||
ok: true,
|
||||
acceptance: {
|
||||
schemaVersion: "nodedc.device-adapter-acceptance.v1",
|
||||
acceptanceRef: "acceptance:test-001",
|
||||
idempotencyKey: `sha256:${"b".repeat(64)}`,
|
||||
status: "accepted",
|
||||
replayed: false,
|
||||
acceptedAt: "2026-08-11T12:00:00.000Z",
|
||||
},
|
||||
}), { status: 201 }),
|
||||
});
|
||||
await assert.rejects(
|
||||
() => client.acceptMessage({
|
||||
idempotencyKey: `sha256:${"a".repeat(64)}`,
|
||||
}),
|
||||
/device_gateway_core_acceptance_mismatch/,
|
||||
);
|
||||
});
|
||||
@@ -0,0 +1,149 @@
|
||||
import assert from "node:assert/strict";
|
||||
import { connect } from "node:net";
|
||||
import test from "node:test";
|
||||
|
||||
import {
|
||||
DEVICE_ADAPTER_CATALOG,
|
||||
} from "../../../packages/device-adapter-catalog/src/index.mjs";
|
||||
import {
|
||||
createControlCoreApp,
|
||||
} from "../../device-control-core/src/app.mjs";
|
||||
import { createCoreGatewayClient } from "../src/core-client.mjs";
|
||||
import { createDeviceGatewayRuntime } from "../src/runtime.mjs";
|
||||
|
||||
const gatewayToken = "test-only-gateway-token-with-32-bytes";
|
||||
const identifierPepper = "test-only-identifier-pepper-with-32-bytes";
|
||||
const specificationHeader = Buffer.from(
|
||||
"FF23E9EF782DE7120300",
|
||||
"hex",
|
||||
);
|
||||
const specificationPackage = Buffer.from(
|
||||
"5B01010000FBDEC251EC5D",
|
||||
"hex",
|
||||
);
|
||||
|
||||
test("B2 HEADER2 becomes a persisted masked quarantine discovery before ACK", async () => {
|
||||
let stored;
|
||||
let storedMessage;
|
||||
const core = createControlCoreApp({
|
||||
discoveryIngestEnabled: true,
|
||||
gatewayToken,
|
||||
identifierPepper,
|
||||
repository: {
|
||||
health: async () => "ready",
|
||||
upsertQuarantineDiscovery: async (value) => {
|
||||
stored = value;
|
||||
return {
|
||||
created: true,
|
||||
value: {
|
||||
...value.safeView,
|
||||
discoveryRef: "discovery:integration-001",
|
||||
},
|
||||
};
|
||||
},
|
||||
acceptAdapterMessage: async (value) => {
|
||||
storedMessage = value;
|
||||
return {
|
||||
acceptance: {
|
||||
schemaVersion: "nodedc.device-adapter-acceptance.v1",
|
||||
acceptanceRef: "acceptance:integration-001",
|
||||
idempotencyKey: value.safeView.idempotencyKey,
|
||||
status: "accepted",
|
||||
replayed: false,
|
||||
acceptedAt: "2026-08-11T12:00:00.000Z",
|
||||
},
|
||||
};
|
||||
},
|
||||
},
|
||||
});
|
||||
await listen(core);
|
||||
const coreAddress = core.address();
|
||||
const observe = createCoreGatewayClient({
|
||||
coreUrl: `http://127.0.0.1:${coreAddress.port}`,
|
||||
gatewayToken,
|
||||
});
|
||||
const gateway = createDeviceGatewayRuntime({
|
||||
healthPort: 0,
|
||||
tcpHost: "0.0.0.0",
|
||||
tcpPort: 0,
|
||||
listenEnabled: true,
|
||||
publicIngressEnabled: true,
|
||||
coreChannelAuthenticated: true,
|
||||
adapterRegistry: DEVICE_ADAPTER_CATALOG.registry,
|
||||
protocolProfileRef: DEVICE_ADAPTER_CATALOG.defaultProfileRef,
|
||||
edgeRef: "edge:integration-001",
|
||||
now: () => new Date(0x52db95de * 1000),
|
||||
onDiscovery: observe.observeDiscovery,
|
||||
onMessage: observe.acceptMessage,
|
||||
onCommandStatus: async () => undefined,
|
||||
});
|
||||
const addresses = await gateway.start();
|
||||
try {
|
||||
const response = await exchange(
|
||||
addresses.tcpAddress.port,
|
||||
Buffer.concat([specificationHeader, specificationPackage]),
|
||||
13,
|
||||
);
|
||||
assert.equal(
|
||||
response.toString("hex").toUpperCase(),
|
||||
"7B0400A0DE95DB527D7B00017D",
|
||||
);
|
||||
assert.match(stored.identifierDigest, /^hmac-sha256:[a-f0-9]{64}$/);
|
||||
assert.equal(stored.safeView.lifecycleState, "quarantine");
|
||||
assert.equal(stored.safeView.identifier.masked, "***********7769");
|
||||
assert.equal(stored.safeView.commandTransport, "disabled");
|
||||
assert.equal(
|
||||
JSON.stringify(stored).includes("865209039777769"),
|
||||
false,
|
||||
);
|
||||
assert.equal(storedMessage.safeView.messageType, "telemetry.package");
|
||||
assert.equal(storedMessage.safeView.identifier.masked, "***********7769");
|
||||
assert.equal(
|
||||
JSON.stringify(storedMessage).includes("865209039777769"),
|
||||
false,
|
||||
);
|
||||
} finally {
|
||||
await gateway.stop();
|
||||
await close(core);
|
||||
}
|
||||
});
|
||||
|
||||
function listen(server) {
|
||||
return new Promise((resolve, reject) => {
|
||||
server.once("error", reject);
|
||||
server.listen(0, "127.0.0.1", resolve);
|
||||
});
|
||||
}
|
||||
|
||||
function close(server) {
|
||||
return new Promise((resolve, reject) => {
|
||||
server.close((error) => (error ? reject(error) : resolve()));
|
||||
server.closeAllConnections?.();
|
||||
});
|
||||
}
|
||||
|
||||
function exchange(port, payload, expectedBytes) {
|
||||
return new Promise((resolve, reject) => {
|
||||
const chunks = [];
|
||||
let byteLength = 0;
|
||||
const socket = connect({ host: "127.0.0.1", port }, () => {
|
||||
socket.write(payload);
|
||||
});
|
||||
socket.on("data", (chunk) => {
|
||||
chunks.push(chunk);
|
||||
byteLength += chunk.length;
|
||||
if (byteLength >= expectedBytes) {
|
||||
socket.destroy();
|
||||
resolve(Buffer.concat(chunks, byteLength));
|
||||
}
|
||||
});
|
||||
socket.on("error", reject);
|
||||
socket.on("close", () => {
|
||||
if (byteLength < expectedBytes) {
|
||||
reject(new Error(
|
||||
`device_gateway_test_socket_closed_early:${byteLength}/${expectedBytes}`,
|
||||
));
|
||||
}
|
||||
});
|
||||
});
|
||||
}
|
||||
@@ -0,0 +1,413 @@
|
||||
import assert from "node:assert/strict";
|
||||
import { connect } from "node:net";
|
||||
import test from "node:test";
|
||||
|
||||
import {
|
||||
DEVICE_ADAPTER_CATALOG,
|
||||
} from "../../../packages/device-adapter-catalog/src/index.mjs";
|
||||
import { createDeviceGatewayRuntime } from "../src/runtime.mjs";
|
||||
|
||||
const specificationHeader = Buffer.from(
|
||||
"FF23E9EF782DE7120300",
|
||||
"hex",
|
||||
);
|
||||
const specificationPackage = Buffer.from(
|
||||
"5B01010000FBDEC251EC5D",
|
||||
"hex",
|
||||
);
|
||||
|
||||
test("baseline health exposes no public ingress and no command transport", async () => {
|
||||
const runtime = createDeviceGatewayRuntime({
|
||||
healthPort: 0,
|
||||
listenEnabled: false,
|
||||
});
|
||||
const addresses = await runtime.start();
|
||||
try {
|
||||
const response = await fetch(
|
||||
`http://127.0.0.1:${addresses.healthAddress.port}/healthz`,
|
||||
);
|
||||
assert.equal(response.status, 200);
|
||||
const body = await response.json();
|
||||
assert.equal(body.publicIngress, "disabled");
|
||||
assert.equal(body.commandTransport, "disabled");
|
||||
assert.equal(body.tcpListener, "disabled");
|
||||
assert.equal(addresses.tcpAddress, null);
|
||||
} finally {
|
||||
await runtime.stop();
|
||||
}
|
||||
});
|
||||
|
||||
test("telemetry ingress persists HEADER2 before acknowledging packages", async () => {
|
||||
const captured = [];
|
||||
const runtime = createDeviceGatewayRuntime({
|
||||
healthPort: 0,
|
||||
tcpHost: "0.0.0.0",
|
||||
tcpPort: 0,
|
||||
listenEnabled: true,
|
||||
publicIngressEnabled: true,
|
||||
coreChannelAuthenticated: true,
|
||||
now: () => new Date(0x52db95de * 1000),
|
||||
...gatewayAdapterOptions(),
|
||||
onDiscovery: async (value) => {
|
||||
captured.push(value);
|
||||
return { lifecycleState: "quarantine" };
|
||||
},
|
||||
onMessage: async (message) => acceptanceFor(message),
|
||||
});
|
||||
const addresses = await runtime.start();
|
||||
const client = await connectAndCollect(addresses.tcpAddress.port);
|
||||
try {
|
||||
client.socket.write(specificationHeader.subarray(0, 4));
|
||||
await new Promise((resolve) => setImmediate(resolve));
|
||||
assert.equal(client.bytes().length, 0);
|
||||
|
||||
client.socket.write(specificationHeader.subarray(4));
|
||||
await client.waitForBytes(9);
|
||||
assert.equal(
|
||||
client.bytes().subarray(0, 9).toString("hex").toUpperCase(),
|
||||
"7B0400A0DE95DB527D",
|
||||
);
|
||||
assert.equal(captured.length, 1);
|
||||
assert.equal(captured[0].identifier.value, "865209039777769");
|
||||
assert.equal(captured[0].evidence.framingStatus, "verified");
|
||||
assert.equal(captured[0].commandTransport, undefined);
|
||||
|
||||
client.socket.write(specificationPackage);
|
||||
await client.waitForBytes(13);
|
||||
assert.equal(
|
||||
client.bytes().subarray(9).toString("hex").toUpperCase(),
|
||||
"7B00017D",
|
||||
);
|
||||
assert.equal(runtime.status().totalDiscoveries, 1);
|
||||
assert.equal(runtime.status().totalMessagesAccepted, 1);
|
||||
assert.equal(runtime.status().totalPackagesAcknowledged, 1);
|
||||
assert.equal(runtime.status().commandTransport, "typed-service-ping-v1");
|
||||
assert.equal(runtime.status().publicIngress, "telemetry-ingest");
|
||||
|
||||
const response = await fetch(
|
||||
`http://127.0.0.1:${addresses.healthAddress.port}/healthz`,
|
||||
);
|
||||
const body = await response.json();
|
||||
assert.equal(body.framing, "verified-read-only");
|
||||
assert.equal(body.tcpListener, "telemetry-ingest");
|
||||
assert.equal(body.publicIngress, "telemetry-ingest");
|
||||
assert.equal(body.commandTransport, "typed-service-ping-v1");
|
||||
} finally {
|
||||
client.socket.destroy();
|
||||
await runtime.stop();
|
||||
}
|
||||
});
|
||||
|
||||
test("dispatches a typed service ping on the next telemetry package and records SERV OK", async () => {
|
||||
const statuses = [];
|
||||
const runtime = createDeviceGatewayRuntime({
|
||||
healthPort: 0,
|
||||
tcpPort: 0,
|
||||
listenEnabled: true,
|
||||
...gatewayAdapterOptions(),
|
||||
onDiscovery: async () => ({ lifecycleState: "claimed" }),
|
||||
onMessage: async (message) => ({
|
||||
...acceptanceFor(message),
|
||||
commandOffer: {
|
||||
commandRef: "command:11111111-1111-4111-8111-111111111111",
|
||||
commandType: "service.ping",
|
||||
accessCode: "123456",
|
||||
transportMessageRef: "edge-command:22222222-2222-4222-8222-222222222222",
|
||||
},
|
||||
}),
|
||||
onCommandStatus: async (status) => statuses.push(status),
|
||||
});
|
||||
const addresses = await runtime.start();
|
||||
const client = await connectAndCollect(addresses.tcpAddress.port);
|
||||
try {
|
||||
client.socket.write(Buffer.concat([specificationHeader, specificationPackage]));
|
||||
await client.waitForBytes(28);
|
||||
assert.equal(
|
||||
client.bytes().subarray(13).toString("ascii"),
|
||||
"123456*SERV*1.1",
|
||||
);
|
||||
client.socket.write(Buffer.from("SERV OK", "ascii"));
|
||||
await waitFor(() => statuses.length === 1);
|
||||
assert.deepEqual(statuses[0], {
|
||||
commandRef: "command:11111111-1111-4111-8111-111111111111",
|
||||
transportMessageRef: "edge-command:22222222-2222-4222-8222-222222222222",
|
||||
lifecycleState: "acknowledged",
|
||||
resultCode: "serv_ok",
|
||||
observedAt: statuses[0].observedAt,
|
||||
sessionRef: statuses[0].sessionRef,
|
||||
adapterProfileRef: "arusnavi.b2.internal.v1",
|
||||
});
|
||||
} finally {
|
||||
client.socket.destroy();
|
||||
await runtime.stop();
|
||||
}
|
||||
});
|
||||
|
||||
test("does not acknowledge malformed or unverified initial bytes", async () => {
|
||||
const captured = [];
|
||||
const runtime = createDeviceGatewayRuntime({
|
||||
healthPort: 0,
|
||||
tcpPort: 0,
|
||||
listenEnabled: true,
|
||||
...gatewayAdapterOptions(),
|
||||
onDiscovery: async (value) => {
|
||||
captured.push(value);
|
||||
return { lifecycleState: "quarantine" };
|
||||
},
|
||||
onMessage: async (message) => acceptanceFor(message),
|
||||
});
|
||||
const addresses = await runtime.start();
|
||||
try {
|
||||
const received = await sendAndCollect(
|
||||
addresses.tcpAddress.port,
|
||||
Buffer.from("not-a-b2-header", "utf8"),
|
||||
);
|
||||
assert.equal(received.length, 0);
|
||||
assert.equal(captured.length, 0);
|
||||
assert.equal(runtime.status().totalRejected, 1);
|
||||
} finally {
|
||||
await runtime.stop();
|
||||
}
|
||||
});
|
||||
|
||||
test("does not acknowledge a header when Core rejects discovery", async () => {
|
||||
const runtime = createDeviceGatewayRuntime({
|
||||
healthPort: 0,
|
||||
tcpPort: 0,
|
||||
listenEnabled: true,
|
||||
...gatewayAdapterOptions(),
|
||||
onDiscovery: async () => {
|
||||
throw new Error("core_unavailable");
|
||||
},
|
||||
onMessage: async (message) => acceptanceFor(message),
|
||||
});
|
||||
const addresses = await runtime.start();
|
||||
try {
|
||||
const received = await sendAndCollect(
|
||||
addresses.tcpAddress.port,
|
||||
specificationHeader,
|
||||
);
|
||||
assert.equal(received.length, 0);
|
||||
assert.equal(runtime.status().totalDiscoveries, 0);
|
||||
assert.equal(runtime.status().totalRejected, 1);
|
||||
} finally {
|
||||
await runtime.stop();
|
||||
}
|
||||
});
|
||||
|
||||
test("does not acknowledge a package when durable Core acceptance fails", async () => {
|
||||
const runtime = createDeviceGatewayRuntime({
|
||||
healthPort: 0,
|
||||
tcpPort: 0,
|
||||
listenEnabled: true,
|
||||
now: () => new Date(0x52db95de * 1000),
|
||||
...gatewayAdapterOptions(),
|
||||
onDiscovery: async () => ({ lifecycleState: "quarantine" }),
|
||||
onMessage: async () => {
|
||||
throw new Error("core_commit_failed");
|
||||
},
|
||||
});
|
||||
const addresses = await runtime.start();
|
||||
try {
|
||||
const received = await sendAndCollect(
|
||||
addresses.tcpAddress.port,
|
||||
Buffer.concat([specificationHeader, specificationPackage]),
|
||||
);
|
||||
assert.equal(
|
||||
received.toString("hex").toUpperCase(),
|
||||
"7B0400A0DE95DB527D",
|
||||
);
|
||||
assert.equal(runtime.status().totalMessagesAccepted, 0);
|
||||
assert.equal(runtime.status().totalPackagesAcknowledged, 0);
|
||||
assert.equal(runtime.status().totalRejected, 1);
|
||||
} finally {
|
||||
await runtime.stop();
|
||||
}
|
||||
});
|
||||
|
||||
test("acknowledges an idempotent Core replay as accepted delivery", async () => {
|
||||
const runtime = createDeviceGatewayRuntime({
|
||||
healthPort: 0,
|
||||
tcpPort: 0,
|
||||
listenEnabled: true,
|
||||
now: () => new Date(0x52db95de * 1000),
|
||||
...gatewayAdapterOptions(),
|
||||
onDiscovery: async () => ({ lifecycleState: "quarantine" }),
|
||||
onMessage: async (message) => acceptanceFor(message, true),
|
||||
});
|
||||
const addresses = await runtime.start();
|
||||
try {
|
||||
const received = await sendAndCollect(
|
||||
addresses.tcpAddress.port,
|
||||
Buffer.concat([specificationHeader, specificationPackage]),
|
||||
);
|
||||
assert.equal(
|
||||
received.toString("hex").toUpperCase(),
|
||||
"7B0400A0DE95DB527D7B00017D",
|
||||
);
|
||||
assert.equal(runtime.status().totalMessagesAccepted, 1);
|
||||
assert.equal(runtime.status().totalPackagesAcknowledged, 1);
|
||||
} finally {
|
||||
await runtime.stop();
|
||||
}
|
||||
});
|
||||
|
||||
test("closes an oversized tracker buffer and releases its aggregate budget", async () => {
|
||||
const runtime = createDeviceGatewayRuntime({
|
||||
healthPort: 0,
|
||||
tcpPort: 0,
|
||||
listenEnabled: true,
|
||||
maxBufferedBytes: 1024,
|
||||
maxAggregateBufferedBytes: 1024,
|
||||
...gatewayAdapterOptions(),
|
||||
onDiscovery: async () => ({ lifecycleState: "quarantine" }),
|
||||
onMessage: async (message) => acceptanceFor(message),
|
||||
});
|
||||
const addresses = await runtime.start();
|
||||
try {
|
||||
const received = await sendAndCollect(
|
||||
addresses.tcpAddress.port,
|
||||
Buffer.alloc(1025, 0x01),
|
||||
);
|
||||
assert.equal(received.length, 0);
|
||||
assert.equal(runtime.status().totalRejected, 1);
|
||||
assert.equal(runtime.status().totalBufferedBytes, 0);
|
||||
assert.equal(runtime.status().activeSessions, 0);
|
||||
} finally {
|
||||
await runtime.stop();
|
||||
}
|
||||
});
|
||||
|
||||
test("public ingress requires both discovery and durable message acceptance sinks", () => {
|
||||
assert.throws(
|
||||
() => createDeviceGatewayRuntime({
|
||||
listenEnabled: true,
|
||||
publicIngressEnabled: true,
|
||||
coreChannelAuthenticated: true,
|
||||
tcpHost: "0.0.0.0",
|
||||
}),
|
||||
/device_gateway_core_acceptance_sink_required/,
|
||||
);
|
||||
});
|
||||
|
||||
test("public ingress cannot start on the legacy bearer HTTP Core client", () => {
|
||||
assert.throws(
|
||||
() => createDeviceGatewayRuntime({
|
||||
listenEnabled: true,
|
||||
publicIngressEnabled: true,
|
||||
tcpHost: "0.0.0.0",
|
||||
...gatewayAdapterOptions(),
|
||||
onDiscovery: async () => ({ lifecycleState: "quarantine" }),
|
||||
onMessage: async (message) => acceptanceFor(message),
|
||||
}),
|
||||
/device_gateway_authenticated_core_channel_required/,
|
||||
);
|
||||
});
|
||||
|
||||
test("baseline rejects non-loopback binding", () => {
|
||||
assert.throws(
|
||||
() => createDeviceGatewayRuntime({
|
||||
listenEnabled: true,
|
||||
tcpHost: "0.0.0.0",
|
||||
...gatewayAdapterOptions(),
|
||||
onDiscovery: async () => ({ lifecycleState: "quarantine" }),
|
||||
onMessage: async (message) => acceptanceFor(message),
|
||||
}),
|
||||
/device_gateway_baseline_loopback_only/,
|
||||
);
|
||||
});
|
||||
|
||||
test("container health may bind all interfaces while TCP stays disabled", async () => {
|
||||
const runtime = createDeviceGatewayRuntime({
|
||||
healthHost: "0.0.0.0",
|
||||
healthPort: 0,
|
||||
listenEnabled: false,
|
||||
});
|
||||
const addresses = await runtime.start();
|
||||
try {
|
||||
assert.equal(addresses.healthAddress.address, "0.0.0.0");
|
||||
assert.equal(addresses.tcpAddress, null);
|
||||
assert.equal(runtime.status().publicIngress, "disabled");
|
||||
} finally {
|
||||
await runtime.stop();
|
||||
}
|
||||
});
|
||||
|
||||
function connectAndCollect(port) {
|
||||
return new Promise((resolve, reject) => {
|
||||
const chunks = [];
|
||||
let byteLength = 0;
|
||||
const waiters = [];
|
||||
const socket = connect({ host: "127.0.0.1", port }, () => {
|
||||
resolve({
|
||||
socket,
|
||||
bytes: () => Buffer.concat(chunks, byteLength),
|
||||
waitForBytes: (minimum) => {
|
||||
if (byteLength >= minimum) return Promise.resolve();
|
||||
return new Promise((waitResolve, waitReject) => {
|
||||
waiters.push({ minimum, waitResolve, waitReject });
|
||||
});
|
||||
},
|
||||
});
|
||||
});
|
||||
socket.on("data", (chunk) => {
|
||||
chunks.push(chunk);
|
||||
byteLength += chunk.length;
|
||||
for (let index = waiters.length - 1; index >= 0; index -= 1) {
|
||||
if (byteLength >= waiters[index].minimum) {
|
||||
waiters[index].waitResolve();
|
||||
waiters.splice(index, 1);
|
||||
}
|
||||
}
|
||||
});
|
||||
socket.on("error", (error) => {
|
||||
for (const waiter of waiters.splice(0)) waiter.waitReject(error);
|
||||
reject(error);
|
||||
});
|
||||
socket.on("close", () => {
|
||||
for (const waiter of waiters.splice(0)) {
|
||||
waiter.waitReject(new Error("device_gateway_test_socket_closed"));
|
||||
}
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
function sendAndCollect(port, payload) {
|
||||
return new Promise((resolve, reject) => {
|
||||
const chunks = [];
|
||||
const socket = connect({ host: "127.0.0.1", port }, () => {
|
||||
socket.end(payload);
|
||||
});
|
||||
socket.on("data", (chunk) => chunks.push(chunk));
|
||||
socket.on("close", () => resolve(Buffer.concat(chunks)));
|
||||
socket.on("error", reject);
|
||||
});
|
||||
}
|
||||
|
||||
function gatewayAdapterOptions() {
|
||||
return {
|
||||
adapterRegistry: DEVICE_ADAPTER_CATALOG.registry,
|
||||
protocolProfileRef: DEVICE_ADAPTER_CATALOG.defaultProfileRef,
|
||||
edgeRef: "edge:test-001",
|
||||
onCommandStatus: async () => undefined,
|
||||
};
|
||||
}
|
||||
|
||||
function acceptanceFor(message, replayed = false) {
|
||||
return {
|
||||
schemaVersion: "nodedc.device-adapter-acceptance.v1",
|
||||
acceptanceRef: "acceptance:test-001",
|
||||
idempotencyKey: message.idempotencyKey,
|
||||
status: "accepted",
|
||||
replayed,
|
||||
acceptedAt: "2026-08-11T12:00:00.000Z",
|
||||
};
|
||||
}
|
||||
|
||||
async function waitFor(predicate, timeoutMs = 1_000) {
|
||||
const deadline = Date.now() + timeoutMs;
|
||||
while (!predicate()) {
|
||||
if (Date.now() >= deadline) throw new Error("test_wait_timeout");
|
||||
await new Promise((resolve) => setTimeout(resolve, 5));
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user