feat(device-plane): add typed B2 service ping transport
This commit is contained in:
@@ -52,6 +52,7 @@ export function createDeviceGatewayRuntime(options = {}) {
|
||||
processing: false,
|
||||
rejected: false,
|
||||
closed: false,
|
||||
pendingCommand: null,
|
||||
};
|
||||
sessions.set(socket, session);
|
||||
incrementAddressSessions(remoteAddress);
|
||||
@@ -103,7 +104,7 @@ export function createDeviceGatewayRuntime(options = {}) {
|
||||
publicIngress: config.publicIngressEnabled
|
||||
? "telemetry-ingest"
|
||||
: "disabled",
|
||||
commandTransport: "disabled",
|
||||
commandTransport: config.commandTransport,
|
||||
sessions: {
|
||||
active: sessions.size,
|
||||
accepted: totalAccepted,
|
||||
@@ -149,7 +150,7 @@ export function createDeviceGatewayRuntime(options = {}) {
|
||||
totalMessagesAccepted,
|
||||
totalPackagesAcknowledged,
|
||||
totalBufferedBytes,
|
||||
commandTransport: "disabled",
|
||||
commandTransport: config.commandTransport,
|
||||
publicIngress: config.publicIngressEnabled
|
||||
? "telemetry-ingest"
|
||||
: "disabled",
|
||||
@@ -186,9 +187,39 @@ export function createDeviceGatewayRuntime(options = {}) {
|
||||
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();
|
||||
@@ -215,9 +246,12 @@ export function createDeviceGatewayRuntime(options = {}) {
|
||||
payloadSchemaRef: parsed.payloadSchemaRef,
|
||||
payload: parsed.payload,
|
||||
});
|
||||
const acceptance = normalizeAdapterAcceptance(
|
||||
await config.onMessage(message),
|
||||
);
|
||||
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");
|
||||
}
|
||||
@@ -228,9 +262,26 @@ export function createDeviceGatewayRuntime(options = {}) {
|
||||
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;
|
||||
}
|
||||
@@ -287,8 +338,21 @@ export function createDeviceGatewayRuntime(options = {}) {
|
||||
function closeSession(socket, session) {
|
||||
if (session.closed) return;
|
||||
session.closed = true;
|
||||
totalBufferedBytes -= session.buffer.length;
|
||||
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);
|
||||
}
|
||||
@@ -302,8 +366,9 @@ export function createDeviceGatewayRuntime(options = {}) {
|
||||
}
|
||||
|
||||
function consumeBuffer(session, bytesConsumed) {
|
||||
session.buffer = session.buffer.subarray(bytesConsumed);
|
||||
totalBufferedBytes -= bytesConsumed;
|
||||
const consumed = Math.min(bytesConsumed, session.buffer.length);
|
||||
session.buffer = session.buffer.subarray(consumed);
|
||||
totalBufferedBytes = Math.max(0, totalBufferedBytes - consumed);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -326,11 +391,25 @@ function normalizeConfig(input) {
|
||||
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",
|
||||
@@ -412,10 +491,36 @@ function normalizeConfig(input) {
|
||||
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");
|
||||
|
||||
@@ -44,12 +44,14 @@ test("B2 HEADER2 becomes a persisted masked quarantine discovery before ACK", as
|
||||
acceptAdapterMessage: async (value) => {
|
||||
storedMessage = value;
|
||||
return {
|
||||
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",
|
||||
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",
|
||||
},
|
||||
};
|
||||
},
|
||||
},
|
||||
@@ -73,6 +75,7 @@ test("B2 HEADER2 becomes a persisted masked quarantine discovery before ACK", as
|
||||
now: () => new Date(0x52db95de * 1000),
|
||||
onDiscovery: observe.observeDiscovery,
|
||||
onMessage: observe.acceptMessage,
|
||||
onCommandStatus: async () => undefined,
|
||||
});
|
||||
const addresses = await gateway.start();
|
||||
try {
|
||||
@@ -115,6 +118,7 @@ function listen(server) {
|
||||
function close(server) {
|
||||
return new Promise((resolve, reject) => {
|
||||
server.close((error) => (error ? reject(error) : resolve()));
|
||||
server.closeAllConnections?.();
|
||||
});
|
||||
}
|
||||
|
||||
@@ -134,5 +138,12 @@ function exchange(port, payload, expectedBytes) {
|
||||
}
|
||||
});
|
||||
socket.on("error", reject);
|
||||
socket.on("close", () => {
|
||||
if (byteLength < expectedBytes) {
|
||||
reject(new Error(
|
||||
`device_gateway_test_socket_closed_early:${byteLength}/${expectedBytes}`,
|
||||
));
|
||||
}
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
@@ -81,7 +81,7 @@ test("telemetry ingress persists HEADER2 before acknowledging packages", async (
|
||||
assert.equal(runtime.status().totalDiscoveries, 1);
|
||||
assert.equal(runtime.status().totalMessagesAccepted, 1);
|
||||
assert.equal(runtime.status().totalPackagesAcknowledged, 1);
|
||||
assert.equal(runtime.status().commandTransport, "disabled");
|
||||
assert.equal(runtime.status().commandTransport, "typed-service-ping-v1");
|
||||
assert.equal(runtime.status().publicIngress, "telemetry-ingest");
|
||||
|
||||
const response = await fetch(
|
||||
@@ -91,7 +91,52 @@ test("telemetry ingress persists HEADER2 before acknowledging packages", async (
|
||||
assert.equal(body.framing, "verified-read-only");
|
||||
assert.equal(body.tcpListener, "telemetry-ingest");
|
||||
assert.equal(body.publicIngress, "telemetry-ingest");
|
||||
assert.equal(body.commandTransport, "disabled");
|
||||
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();
|
||||
@@ -319,6 +364,11 @@ function connectAndCollect(port) {
|
||||
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"));
|
||||
}
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
@@ -339,6 +389,7 @@ function gatewayAdapterOptions() {
|
||||
adapterRegistry: DEVICE_ADAPTER_CATALOG.registry,
|
||||
protocolProfileRef: DEVICE_ADAPTER_CATALOG.defaultProfileRef,
|
||||
edgeRef: "edge:test-001",
|
||||
onCommandStatus: async () => undefined,
|
||||
};
|
||||
}
|
||||
|
||||
@@ -352,3 +403,11 @@ function acceptanceFor(message, replayed = false) {
|
||||
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