Files
NODEDC_PLATFORM/device-plane/services/device-edge-channel/test/channel-integration.test.mjs
T

619 lines
20 KiB
JavaScript

import assert from "node:assert/strict";
import { spawnSync } from "node:child_process";
import { randomUUID, X509Certificate } from "node:crypto";
import { mkdtemp, readFile, rm, writeFile } from "node:fs/promises";
import { connect as connectHttp2 } from "node:http2";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { after, before, test } from "node:test";
import {
DEVICE_ADAPTER_ACCEPTANCE_SCHEMA,
DEVICE_ADAPTER_MESSAGE_SCHEMA,
DEVICE_DISCOVERY_SIGNAL_SCHEMA,
} from "../../../packages/device-protocol-contract/src/index.mjs";
import { createDeviceGatewayCoreChannelClient } from
"../../device-gateway-core/src/runtime.mjs";
import { createDeviceEdgeChannelServer } from "../src/runtime.mjs";
let fixtureDirectory;
let certificates;
before(async () => {
fixtureDirectory = await mkdtemp(join(tmpdir(), "nodedc-edge-channel-"));
certificates = await generateCertificateFixture(fixtureDirectory);
});
after(async () => {
await rm(fixtureDirectory, { recursive: true, force: true });
});
test("accepts synthetic discovery and durable message results over Core-initiated mTLS", async () => {
const acceptedMessages = new Map();
let messageCalls = 0;
const pair = await startPair({
acceptMessage: async (message) => {
messageCalls += 1;
const previous = acceptedMessages.get(message.idempotencyKey);
if (previous) return { ...previous, replayed: true };
const acceptance = acceptanceFor(message, false);
acceptedMessages.set(message.idempotencyKey, acceptance);
return acceptance;
},
});
try {
const discovery = await pair.edge.submitDiscovery(discoverySignal());
assert.equal(discovery.lifecycleState, "quarantine");
assert.equal(discovery.identifier.masked, "***********1088");
const message = adapterMessage();
const first = await pair.edge.submitAdapterMessage(message);
const replay = await pair.edge.submitAdapterMessage(message);
assert.equal(first.replayed, false);
assert.equal(replay.replayed, true);
assert.equal(messageCalls, 2);
assert.equal(pair.edge.status().eventsAccepted, 3);
assert.equal(pair.core.status().eventsAccepted, 3);
assert.equal(pair.edge.status().trackerIngress, "disabled");
assert.equal(pair.core.status().commandTransport, "disabled");
} finally {
await stopPair(pair);
}
});
test("keeps the channel alive and reconnects without losing idempotency", async () => {
const acceptedMessages = new Map();
const pair = await startPair({
keepaliveMs: 20,
deadPeerMs: 60,
reconnectMinimumMs: 20,
reconnectMaximumMs: 40,
acceptMessage: async (message) => {
const previous = acceptedMessages.get(message.idempotencyKey);
if (previous) return { ...previous, replayed: true };
const acceptance = acceptanceFor(message, false);
acceptedMessages.set(message.idempotencyKey, acceptance);
return acceptance;
},
});
try {
await delay(120);
assert.equal(pair.edge.status().channel, "accepted");
const message = adapterMessage();
assert.equal((await pair.edge.submitAdapterMessage(message)).replayed, false);
pair.edge.disconnectActiveChannel();
await waitFor(() => pair.core.status().connectionAttempts >= 2
&& pair.edge.status().channel === "accepted", 2_000);
assert.equal((await pair.edge.submitAdapterMessage(message)).replayed, true);
assert.ok(pair.core.status().reconnects >= 1);
} finally {
await stopPair(pair);
}
});
test("rotates the Edge certificate through staged overlap and rejects retired identity", async () => {
const edge = createEdgeServer();
const address = await edge.start();
let registration = edgeRegistration(address, [
{
generationRef: "trust-generation:1",
fingerprint: certificates.edge.fingerprint,
status: "active",
},
{
generationRef: "trust-generation:2",
fingerprint: certificates.edgeNext.fingerprint,
status: "staged",
},
]);
const core = createCoreClient({
address,
registrationProvider: async () => registration,
});
try {
await core.start();
await core.waitForReady(2_000);
await waitFor(() => edge.status().channel === "accepted", 2_000);
assert.equal(core.status().edgeTrustGeneration, "trust-generation:1");
edge.rotateTrust({
generationRef: "trust-generation:2",
key: certificates.edgeNext.key,
cert: certificates.edgeNext.cert,
ca: certificates.ca,
allowedCoreFingerprints: [certificates.core.fingerprint],
});
await waitFor(() => core.status().connectionAttempts >= 2
&& core.status().channel === "accepted"
&& core.status().edgeTrustGeneration === "trust-generation:2"
&& edge.status().channel === "accepted", 2_000);
registration = edgeRegistration(address, [{
generationRef: "trust-generation:2",
fingerprint: certificates.edgeNext.fingerprint,
status: "active",
}]);
assert.equal((await edge.submitAdapterMessage(adapterMessage())).status, "accepted");
const failuresBeforeRollback = core.status().protocolFailures;
edge.rotateTrust({
generationRef: "trust-generation:3",
key: certificates.edge.key,
cert: certificates.edge.cert,
ca: certificates.ca,
allowedCoreFingerprints: [certificates.core.fingerprint],
});
await waitFor(() => core.status().protocolFailures > failuresBeforeRollback, 2_000);
assert.notEqual(core.status().channel, "accepted");
} finally {
await core.stop();
await edge.stop();
}
});
test("rejects a revoked Edge registration before opening a channel", async () => {
const edge = createEdgeServer();
const address = await edge.start();
const core = createCoreClient({
address,
lifecycleState: "revoked",
});
try {
await core.start();
await waitFor(() => core.status().connectionAttempts >= 1, 500);
assert.equal(core.status().channel, "absent");
assert.equal(edge.status().channelsAccepted, 0);
assert.match(core.status().lastErrorCode, /registration_inactive/);
} finally {
await core.stop();
await edge.stop();
}
});
test("rejects an authenticated but non-allowlisted Core certificate", async () => {
const edge = createEdgeServer();
const address = await edge.start();
const core = createCoreClient({
address,
clientCertificate: certificates.intruder,
});
try {
await core.start();
await waitFor(() => edge.status().channelsRejected >= 1, 2_000);
assert.equal(edge.status().channelsAccepted, 0);
assert.equal(edge.status().channel, "absent");
} finally {
await core.stop();
await edge.stop();
}
});
test("rejects an unknown Edge certificate fingerprint", async () => {
const edge = createEdgeServer();
const address = await edge.start();
const core = createCoreClient({
address,
expectedEdgeFingerprint: certificates.intruder.fingerprint,
});
try {
await core.start();
await waitFor(() => core.status().protocolFailures >= 1, 2_000);
assert.equal(core.status().channelsAccepted, 0);
assert.match(core.status().lastErrorCode, /edge_identity_mismatch/);
} finally {
await core.stop();
await edge.stop();
}
});
test("returns a conclusive rejection when Core cannot accept a package", async () => {
const pair = await startPair({
acceptMessage: async () => {
throw new Error("device_control_core_unavailable");
},
});
try {
await assert.rejects(
pair.edge.submitAdapterMessage(adapterMessage()),
/device_control_core_unavailable/,
);
assert.equal(pair.edge.status().eventsRejected, 1);
assert.equal(pair.core.status().eventsRejected, 1);
} finally {
await stopPair(pair);
}
});
test("isolates tracker session ordering while allowing cross-session progress", async () => {
let releaseSlow;
const slowGate = new Promise((resolve) => {
releaseSlow = resolve;
});
const calls = [];
const pair = await startPair({
acceptMessage: async (message) => {
calls.push(message.sessionRef);
if (message.sessionRef === "session:slow") await slowGate;
return acceptanceFor(message, false);
},
});
try {
let slowResolved = false;
const slow = pair.edge.submitAdapterMessage(adapterMessage({
sessionRef: "session:slow",
messageRef: "message:slow-1",
idempotencyKey: `sha256:${"c".repeat(64)}`,
})).then((value) => {
slowResolved = true;
return value;
});
await waitFor(() => calls.includes("session:slow"), 500);
const fast = await pair.edge.submitAdapterMessage(adapterMessage({
sessionRef: "session:fast",
messageRef: "message:fast-1",
idempotencyKey: `sha256:${"d".repeat(64)}`,
}));
assert.equal(fast.status, "accepted");
assert.equal(slowResolved, false);
releaseSlow();
assert.equal((await slow).status, "accepted");
} finally {
releaseSlow?.();
await stopPair(pair);
}
});
test("preserves per-session order and applies a bounded acceptance window", async () => {
let releaseFirst;
const firstGate = new Promise((resolve) => {
releaseFirst = resolve;
});
const calls = [];
const pair = await startPair({
maxPendingAcceptances: 2,
acceptMessage: async (message) => {
calls.push(message.messageRef);
if (message.messageRef === "message:ordered-1") await firstGate;
return acceptanceFor(message, false);
},
});
try {
const first = pair.edge.submitAdapterMessage(adapterMessage({
sessionRef: "session:ordered",
messageRef: "message:ordered-1",
idempotencyKey: `sha256:${"e".repeat(64)}`,
sequence: 1,
}));
await waitFor(() => calls.length === 1, 500);
const second = pair.edge.submitAdapterMessage(adapterMessage({
sessionRef: "session:ordered",
messageRef: "message:ordered-2",
idempotencyKey: `sha256:${"f".repeat(64)}`,
sequence: 2,
}));
await delay(30);
assert.deepEqual(calls, ["message:ordered-1"]);
await assert.rejects(
pair.edge.submitAdapterMessage(adapterMessage({
sessionRef: "session:overflow",
messageRef: "message:overflow-1",
idempotencyKey: `sha256:${"1".repeat(64)}`,
})),
/acceptance_window_full/,
);
releaseFirst();
await Promise.all([first, second]);
assert.deepEqual(calls, ["message:ordered-1", "message:ordered-2"]);
} finally {
releaseFirst?.();
await stopPair(pair);
}
});
test("closes the logical session on an unknown message kind", async () => {
const edge = createEdgeServer();
const address = await edge.start();
const session = connectHttp2(`https://127.0.0.1:${address.port}`, {
key: certificates.core.key,
cert: certificates.core.cert,
ca: certificates.ca,
servername: "localhost",
minVersion: "TLSv1.3",
maxVersion: "TLSv1.3",
rejectUnauthorized: true,
});
try {
await onceEvent(session, "connect");
const request = session.request({
":method": "POST",
":path": "/internal/v1/device-edge/channel",
}, { endStream: false });
const hello = await readFirstEnvelope(request);
const invalid = {
schemaVersion: hello.schemaVersion,
edgeRegistrationId: hello.edgeRegistrationId,
channelGeneration: hello.channelGeneration,
trackerSessionId: "channel:control",
adapterProfileRef: "channel.control.v1",
sequence: 1,
eventAt: new Date().toISOString(),
receivedAt: new Date().toISOString(),
payloadBytes: 2,
messageKind: "tcp.forward",
correlationId: `correlation:${randomUUID()}`,
payload: {},
};
request.write(`${JSON.stringify(invalid)}\n`);
await waitFor(() => edge.status().protocolFailures >= 1, 1_000);
assert.equal(edge.status().channel, "absent");
request.close();
} finally {
session.close();
await edge.stop();
}
});
async function startPair(options = {}) {
const edge = createEdgeServer(options);
const address = await edge.start();
const core = createCoreClient({ address, ...options });
await core.start();
await core.waitForReady(2_000);
await waitFor(() => edge.status().channel === "accepted", 2_000);
return { edge, core };
}
function createEdgeServer(options = {}) {
const edgeCertificate = options.edgeCertificate ?? certificates.edge;
return createDeviceEdgeChannelServer({
edgeRegistrationId: "edge:pilot-1",
channelGeneration: "generation:pilot-1",
trustGeneration: options.edgeTrustGeneration ?? "trust-generation:1",
host: "127.0.0.1",
port: 0,
tls: {
key: edgeCertificate.key,
cert: edgeCertificate.cert,
ca: certificates.ca,
allowedCoreFingerprints: options.allowedCoreFingerprints
?? [certificates.core.fingerprint],
},
keepaliveMs: options.keepaliveMs ?? 50,
deadPeerMs: options.deadPeerMs ?? 150,
acceptanceTimeoutMs: 500,
maxPendingAcceptances: options.maxPendingAcceptances,
});
}
function createCoreClient(options) {
const clientCertificate = options.clientCertificate ?? certificates.core;
const registration = edgeRegistration(options.address,
options.edgeCertificateIdentities ?? [{
generationRef: "trust-generation:1",
fingerprint: options.expectedEdgeFingerprint
?? certificates.edge.fingerprint,
status: "active",
}],
options.lifecycleState ?? "active");
return createDeviceGatewayCoreChannelClient({
...(options.registrationProvider
? { registrationProvider: options.registrationProvider }
: { registration }),
tls: {
key: clientCertificate.key,
cert: clientCertificate.cert,
ca: certificates.ca,
},
coreIdentity: "workload:device-gateway-core",
keepaliveMs: options.keepaliveMs ?? 50,
deadPeerMs: options.deadPeerMs ?? 150,
reconnectMinimumMs: options.reconnectMinimumMs ?? 20,
reconnectMaximumMs: options.reconnectMaximumMs ?? 80,
random: () => 0,
observeDiscovery: options.observeDiscovery ?? (async () => ({
schemaVersion: "nodedc.device.discovery-view.v1",
lifecycleState: "quarantine",
commandTransport: "disabled",
identifier: { kind: "imei", masked: "***********1088" },
})),
acceptMessage: options.acceptMessage ?? (async (message) =>
acceptanceFor(message, false)),
});
}
function edgeRegistration(address, certificateIdentities, lifecycleState = "active") {
return {
edgeRegistrationId: "edge:pilot-1",
endpoint: `https://127.0.0.1:${address.port}/`,
servername: "localhost",
certificateIdentities,
lifecycleState,
};
}
async function stopPair(pair) {
await pair.core.stop();
await pair.edge.stop();
}
function discoverySignal() {
return {
schemaVersion: DEVICE_DISCOVERY_SIGNAL_SCHEMA,
sessionRef: "session:pilot-1",
modelProfileRef: "arusnavi.internal.b2.v1",
protocol: "INTERNAL",
observedAt: new Date().toISOString(),
identifier: { kind: "imei", value: "863151070211088" },
evidence: {
transport: "tcp",
bytesObserved: 10,
framingStatus: "verified",
specificationRef: "arusnavi.internal.protocol.v1",
},
};
}
function adapterMessage(overrides = {}) {
return {
schemaVersion: DEVICE_ADAPTER_MESSAGE_SCHEMA,
edgeRef: "edge:pilot-1",
adapterRef: "arusnavi-b2",
protocolProfileRef: "arusnavi.internal.b2.v1",
protocol: "INTERNAL",
sessionRef: "session:pilot-1",
messageRef: "message:pilot-1",
messageType: "telemetry.package",
sequence: 1,
observedAt: new Date().toISOString(),
idempotencyKey: `sha256:${"a".repeat(64)}`,
identifier: { kind: "imei", value: "863151070211088" },
payloadSchemaRef: "arusnavi.internal.package-metadata.v1",
payload: {
packageNumber: 1,
packetCount: 1,
byteLength: 11,
packageDigest: `sha256:${"b".repeat(64)}`,
},
...overrides,
};
}
function acceptanceFor(message, replayed) {
return {
schemaVersion: DEVICE_ADAPTER_ACCEPTANCE_SCHEMA,
acceptanceRef: "acceptance:pilot-1",
idempotencyKey: message.idempotencyKey,
status: "accepted",
replayed,
acceptedAt: "2026-08-11T12:00:00.000Z",
};
}
async function generateCertificateFixture(directory) {
runOpenSsl(directory, [
"req", "-x509", "-newkey", "rsa:2048", "-nodes", "-sha256",
"-days", "1", "-subj", "/CN=NODEDC Test Device Edge CA",
"-addext", "basicConstraints=critical,CA:TRUE",
"-addext", "keyUsage=critical,keyCertSign,cRLSign",
"-keyout", "ca.key", "-out", "ca.crt",
]);
await issueCertificate(directory, "edge", "localhost", [
"subjectAltName=DNS:localhost,IP:127.0.0.1",
"extendedKeyUsage=serverAuth",
]);
await issueCertificate(directory, "edge-next", "localhost", [
"subjectAltName=DNS:localhost,IP:127.0.0.1",
"extendedKeyUsage=serverAuth",
]);
await issueCertificate(directory, "core", "nodedc-device-gateway-core", [
"extendedKeyUsage=clientAuth",
]);
await issueCertificate(directory, "intruder", "unapproved-core", [
"extendedKeyUsage=clientAuth",
]);
const ca = await readFile(join(directory, "ca.crt"));
return {
ca,
edge: await readCertificate(directory, "edge"),
edgeNext: await readCertificate(directory, "edge-next"),
core: await readCertificate(directory, "core"),
intruder: await readCertificate(directory, "intruder"),
};
}
async function issueCertificate(directory, name, commonName, extensions) {
runOpenSsl(directory, [
"req", "-newkey", "rsa:2048", "-nodes", "-sha256",
"-subj", `/CN=${commonName}`,
"-keyout", `${name}.key`, "-out", `${name}.csr`,
]);
const extensionFile = `${name}.ext`;
await writeFile(join(directory, extensionFile), [
"basicConstraints=critical,CA:FALSE",
"keyUsage=critical,digitalSignature,keyEncipherment",
...extensions,
].join("\n"));
runOpenSsl(directory, [
"x509", "-req", "-sha256", "-days", "1",
"-in", `${name}.csr`, "-CA", "ca.crt", "-CAkey", "ca.key",
"-CAcreateserial", "-extfile", extensionFile, "-out", `${name}.crt`,
]);
}
async function readCertificate(directory, name) {
const cert = await readFile(join(directory, `${name}.crt`));
return {
key: await readFile(join(directory, `${name}.key`)),
cert,
fingerprint: new X509Certificate(cert).fingerprint256,
};
}
function runOpenSsl(directory, arguments_) {
const result = spawnSync("openssl", arguments_, {
cwd: directory,
encoding: "utf8",
});
if (result.status !== 0) {
throw new Error(`openssl_failed:${result.stderr}`);
}
}
function readFirstEnvelope(stream) {
return new Promise((resolve, reject) => {
let buffered = "";
const onData = (chunk) => {
buffered += chunk.toString("utf8");
const newline = buffered.indexOf("\n");
if (newline < 0) return;
cleanup();
resolve(JSON.parse(buffered.slice(0, newline)));
};
const onError = (error) => {
cleanup();
reject(error);
};
const cleanup = () => {
stream.off("data", onData);
stream.off("error", onError);
};
stream.on("data", onData);
stream.once("error", onError);
});
}
function onceEvent(emitter, event) {
return new Promise((resolve, reject) => {
const onEvent = (...args) => {
cleanup();
resolve(args);
};
const onError = (error) => {
cleanup();
reject(error);
};
const cleanup = () => {
emitter.off(event, onEvent);
emitter.off("error", onError);
};
emitter.once(event, onEvent);
emitter.once("error", onError);
});
}
async function waitFor(predicate, timeoutMs) {
const deadline = Date.now() + timeoutMs;
while (Date.now() < deadline) {
if (predicate()) return;
await delay(10);
}
assert.fail("condition_timeout");
}
function delay(milliseconds) {
return new Promise((resolve) => setTimeout(resolve, milliseconds));
}