Files
NODEDC_PLATFORM/device-plane/services/device-edge-relay/src/runtime.mjs
T

385 lines
11 KiB
JavaScript

import { createServer as createHttpServer } from "node:http";
import { connect, createServer as createTcpServer } from "node:net";
export function createDeviceEdgeRelayRuntime(options = {}) {
const config = normalizeConfig(options);
const sessions = new Map();
const sessionsByAddress = new Map();
const connectionWindows = new Map();
let totalAccepted = 0;
let totalRejected = 0;
let totalForwarded = 0;
const tcpServer = createTcpServer({ allowHalfOpen: true }, (socket) => {
const remoteAddress = config.resolveRemoteAddress(socket.remoteAddress);
if (
!allowsSource(remoteAddress)
|| sessions.size >= config.maxConcurrentSessions
|| currentAddressSessions(remoteAddress) >= config.maxSessionsPerAddress
|| !consumeConnectionPermit(remoteAddress)
) {
totalRejected += 1;
socket.destroy();
return;
}
const session = {
remoteAddress,
socket,
upstream: null,
closed: false,
forwarded: false,
inboundBytes: 0,
outboundBytes: 0,
};
sessions.set(socket, session);
incrementAddressSessions(remoteAddress);
totalAccepted += 1;
socket.setNoDelay(true);
socket.setTimeout(config.sessionTimeoutMs);
socket.pause();
socket.on("timeout", () => rejectSession(session));
socket.on("close", () => closeSession(session));
socket.on("error", () => rejectSession(session));
socket.on("data", (chunk) => {
session.inboundBytes += chunk.length;
if (session.inboundBytes > config.maxBytesPerDirection) {
rejectSession(session);
}
});
const upstream = connect({
host: config.upstreamHost,
port: config.upstreamPort,
});
session.upstream = upstream;
upstream.setNoDelay(true);
upstream.setTimeout(config.sessionTimeoutMs);
upstream.on("connect", () => {
if (session.closed) {
upstream.destroy();
return;
}
session.forwarded = true;
totalForwarded += 1;
socket.pipe(upstream);
upstream.pipe(socket);
socket.resume();
});
upstream.on("timeout", () => rejectSession(session));
upstream.on("error", () => rejectSession(session));
upstream.on("close", () => closeSession(session));
upstream.on("data", (chunk) => {
session.outboundBytes += chunk.length;
if (session.outboundBytes > config.maxBytesPerDirection) {
rejectSession(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;
response.end('{"ok":false,"error":"device_edge_relay_route_not_found"}\n');
return;
}
response.statusCode = 200;
response.end(`${JSON.stringify({
ok: true,
service: "nodedc-device-edge-relay",
ingress: config.ingressEnabled ? "relay-only" : "disabled",
protocolInspection: "disabled",
commandTransport: "disabled",
sourceAdmission: config.sourcePolicy,
sessions: {
active: sessions.size,
accepted: totalAccepted,
rejected: totalRejected,
forwarded: totalForwarded,
},
})}\n`);
});
return {
async start() {
await listen(healthServer, config.healthPort, config.healthHost);
if (config.ingressEnabled) {
await listen(tcpServer, config.tcpPort, config.tcpHost);
}
return {
healthAddress: healthServer.address(),
tcpAddress: config.ingressEnabled ? tcpServer.address() : null,
};
},
async stop() {
for (const session of sessions.values()) rejectSession(session);
await Promise.all([
closeServer(healthServer),
config.ingressEnabled ? closeServer(tcpServer) : Promise.resolve(),
]);
},
status() {
return {
activeSessions: sessions.size,
totalAccepted,
totalRejected,
totalForwarded,
ingress: config.ingressEnabled ? "relay-only" : "disabled",
protocolInspection: "disabled",
commandTransport: "disabled",
sourceAdmission: config.sourcePolicy,
};
},
};
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 allowsSource(remoteAddress) {
if (config.sourcePolicy === "any") return true;
return isPublicIpv4Address(remoteAddress);
}
function rejectSession(session) {
if (!session.closed) totalRejected += 1;
session.socket.destroy();
session.upstream?.destroy();
closeSession(session);
}
function closeSession(session) {
if (session.closed) return;
session.closed = true;
sessions.delete(session.socket);
decrementAddressSessions(session.remoteAddress);
}
}
function normalizeConfig(input) {
const ingressEnabled = input.ingressEnabled === true;
return {
ingressEnabled,
healthHost: normalizeHost(input.healthHost, "127.0.0.1"),
healthPort: parseInteger(
input.healthPort,
18221,
0,
65535,
"device_edge_relay_health_port_invalid",
),
tcpHost: normalizeTcpHost(input.tcpHost, ingressEnabled),
tcpPort: parseInteger(
input.tcpPort,
9921,
0,
65535,
"device_edge_relay_tcp_port_invalid",
),
upstreamHost: ingressEnabled
? normalizeUpstreamHost(input.upstreamHost)
: "disabled",
upstreamPort: ingressEnabled
? parseInteger(
input.upstreamPort,
undefined,
1,
65535,
"device_edge_relay_upstream_port_invalid",
)
: 0,
maxConcurrentSessions: parseInteger(
input.maxConcurrentSessions,
100,
1,
10000,
"device_edge_relay_session_limit_invalid",
),
maxSessionsPerAddress: parseInteger(
input.maxSessionsPerAddress,
10,
1,
1000,
"device_edge_relay_address_session_limit_invalid",
),
maxConnectionsPerMinutePerAddress: parseInteger(
input.maxConnectionsPerMinutePerAddress,
30,
1,
10000,
"device_edge_relay_connection_rate_invalid",
),
maxTrackedSourceAddresses: parseInteger(
input.maxTrackedSourceAddresses,
2048,
1,
65_536,
"device_edge_relay_source_table_limit_invalid",
),
maxBytesPerDirection: parseInteger(
input.maxBytesPerDirection,
262_144,
1_024,
16 * 1024 * 1024,
"device_edge_relay_byte_limit_invalid",
),
sourcePolicy: normalizeSourcePolicy(input.sourcePolicy, ingressEnabled),
resolveRemoteAddress: typeof input.resolveRemoteAddress === "function"
? input.resolveRemoteAddress
: normalizeRemoteAddress,
sessionTimeoutMs: parseInteger(
input.sessionTimeoutMs,
10000,
100,
60000,
"device_edge_relay_session_timeout_invalid",
),
now: typeof input.now === "function" ? input.now : () => new Date(),
};
}
function normalizeHost(value, fallback) {
const normalized = String(value || fallback).trim();
if (!["127.0.0.1", "::1", "0.0.0.0", "::"].includes(normalized)) {
throw new TypeError("device_edge_relay_health_host_invalid");
}
return normalized;
}
function normalizeTcpHost(value, ingressEnabled) {
const fallback = ingressEnabled ? "0.0.0.0" : "127.0.0.1";
const normalized = String(value || fallback).trim();
const allowed = ingressEnabled ? ["0.0.0.0", "::"] : ["127.0.0.1", "::1"];
if (!allowed.includes(normalized)) {
throw new TypeError(
ingressEnabled
? "device_edge_relay_public_ingress_host_invalid"
: "device_edge_relay_baseline_loopback_only",
);
}
return normalized;
}
function normalizeUpstreamHost(value) {
const normalized = String(value || "").trim();
if (
normalized.length === 0
|| normalized.length > 253
|| /[/:\\s]/.test(normalized)
) {
throw new TypeError("device_edge_relay_upstream_host_invalid");
}
return normalized;
}
function normalizeRemoteAddress(value) {
const normalized = String(value || "unknown").trim();
return normalized.slice(0, 64) || "unknown";
}
function normalizeSourcePolicy(value, ingressEnabled) {
const fallback = ingressEnabled ? "public-ipv4-only" : "any";
const normalized = String(value || fallback).trim().toLowerCase();
if (!["any", "public-ipv4-only"].includes(normalized)) {
throw new TypeError("device_edge_relay_source_policy_invalid");
}
if (ingressEnabled && normalized !== "public-ipv4-only") {
throw new TypeError("device_edge_relay_ingress_source_policy_invalid");
}
return normalized;
}
function isPublicIpv4Address(value) {
const normalized = String(value || "").trim().replace(/^::ffff:/i, "");
const parts = normalized.split(".");
if (parts.length !== 4) return false;
const octets = parts.map((part) => Number(part));
if (octets.some((part) => !Number.isInteger(part) || part < 0 || part > 255)) {
return false;
}
const [first, second, third] = octets;
if (
first === 0
|| first === 10
|| first === 127
|| first >= 224
|| (first === 100 && second >= 64 && second <= 127)
|| (first === 169 && second === 254)
|| (first === 172 && second >= 16 && second <= 31)
|| (first === 192 && second === 0 && third === 0)
|| (first === 192 && second === 0 && third === 2)
|| (first === 192 && second === 88 && third === 99)
|| (first === 192 && second === 168)
|| (first === 198 && (second === 18 || second === 19))
|| (first === 198 && second === 51 && third === 100)
|| (first === 203 && second === 0 && third === 113)
) {
return false;
}
return true;
}
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, () => {
server.off("error", reject);
resolve();
});
});
}
function closeServer(server) {
return new Promise((resolve, reject) => {
server.close((error) => (error ? reject(error) : resolve()));
});
}