988 lines
39 KiB
JavaScript
988 lines
39 KiB
JavaScript
import { timingSafeEqual } from "node:crypto";
|
|
import { link, mkdir, open, readFile, unlink } from "node:fs/promises";
|
|
import { createServer } from "node:http";
|
|
import https from "node:https";
|
|
import { connect as connectNet } from "node:net";
|
|
import { dirname } from "node:path";
|
|
import { connect as connectTls } from "node:tls";
|
|
|
|
const allowedHosts = new Set([
|
|
"api.cesium.com",
|
|
"assets.ion.cesium.com",
|
|
"dev.virtualearth.net",
|
|
"ecn.t0.tiles.virtualearth.net",
|
|
"ecn.t1.tiles.virtualearth.net",
|
|
"ecn.t2.tiles.virtualearth.net",
|
|
"ecn.t3.tiles.virtualearth.net",
|
|
]);
|
|
const connectionTimingSymbol = Symbol("nodedc.connection-timing");
|
|
const config = await readConfig();
|
|
const transport = createTransport(config);
|
|
|
|
const server = createServer(async (request, response) => {
|
|
const url = new URL(request.url || "/", "http://localhost");
|
|
try {
|
|
if (request.method === "GET" && ["/healthz", "/status"].includes(url.pathname)) return writeJson(response, 200, await statusBody());
|
|
if (request.method === "POST" && url.pathname === "/api/pair") return await pairConnector(request, response);
|
|
if (["GET", "HEAD"].includes(request.method || "") && url.pathname === "/proxy/cesium/fetch") return await forwardCesiumRequest(request, response, url);
|
|
return writeJson(response, 404, { ok: false, error: "not_found" });
|
|
} catch (error) {
|
|
if (response.destroyed || response.headersSent || response.writableEnded) {
|
|
if (!response.writableEnded) response.destroy();
|
|
return;
|
|
}
|
|
const errorCode = safeError(error);
|
|
const status = Number(error?.statusCode || 502);
|
|
console.warn(JSON.stringify({ event: "cesium_egress_request_failed", error: errorCode, status, ...telemetryLogFields(error?.telemetry) }));
|
|
return writeJson(response, status, { ok: false, error: errorCode }, { "x-nodedc-egress-error": errorCode });
|
|
}
|
|
});
|
|
|
|
server.listen(config.port, config.bindAddress, () => console.log(JSON.stringify({ event: "dc_amd_proxy_started", bindAddress: config.bindAddress, port: config.port, connector: `${config.connectorHost}:${config.connectorPort}`, directEgress: false })));
|
|
for (const signal of ["SIGINT", "SIGTERM"]) process.on(signal, () => server.close(() => transport.close(() => process.exit(0))));
|
|
|
|
async function readConfig() {
|
|
const nodeEnv = String(process.env.NODE_ENV || "production").trim();
|
|
return {
|
|
port: parsePort(process.env.PORT, 8790),
|
|
bindAddress: parseIpv4(process.env.DC_AMD_PROXY_BIND_ADDRESS || "172.22.0.222", "bind_address"),
|
|
connectorHost: parseIpv4(process.env.DC_AMD_CONNECTOR_HOST || "172.22.0.183", "connector_host"),
|
|
connectorPort: parsePort(process.env.DC_AMD_CONNECTOR_PORT, 8791),
|
|
pairAllowedSource: parseIpv4(process.env.DC_AMD_PAIR_ALLOWED_SOURCE || "172.22.0.183", "pair_allowed_source"),
|
|
connectorTokenFile: String(process.env.DC_AMD_CONNECTOR_TOKEN_FILE || "/var/lib/dc-amd-proxy/connector-access").trim(),
|
|
mapToken: Buffer.from(await readSecretFile(process.env.DC_AMD_MAP_EGRESS_TOKEN_FILE || "/run/nodedc-secrets/map-egress-proxy-token", "map_egress_token"), "utf8"),
|
|
connectTimeoutMs: parseDuration(process.env.DC_AMD_CONNECT_TIMEOUT_SECONDS, 20),
|
|
upstreamTimeoutMs: parseDuration(process.env.DC_AMD_UPSTREAM_TIMEOUT_SECONDS, 30),
|
|
bodyIdleTimeoutMs: parseDuration(process.env.DC_AMD_BODY_IDLE_TIMEOUT_SECONDS, 30),
|
|
poolMaxSockets: parsePoolSize(process.env.DC_AMD_POOL_MAX_SOCKETS, 8),
|
|
poolMaxFreeSockets: parsePoolSize(process.env.DC_AMD_POOL_MAX_FREE_SOCKETS, 4),
|
|
slowRequestMs: parseDuration(process.env.DC_AMD_SLOW_REQUEST_SECONDS, 2),
|
|
allowInsecureTls: nodeEnv === "test" && parseBoolean(process.env.DC_AMD_PROXY_TEST_ALLOW_INSECURE_TLS, false),
|
|
};
|
|
}
|
|
|
|
async function statusBody() {
|
|
const paired = Boolean(await readConnectorToken());
|
|
return {
|
|
ok: true,
|
|
service: "dc-amd-proxy",
|
|
state: paired ? "paired" : "awaiting_pair",
|
|
forwarding: paired ? "amd_connector_only" : "disabled",
|
|
directEgress: false,
|
|
connector: `${config.connectorHost}:${config.connectorPort}`,
|
|
transport: transport.status(),
|
|
};
|
|
}
|
|
|
|
async function pairConnector(request, response) {
|
|
if (normalizeAddress(request.socket.remoteAddress) !== config.pairAllowedSource) return writeJson(response, 403, { ok: false, error: "pair_source_not_allowed" });
|
|
if (await readConnectorToken()) return writeJson(response, 409, { ok: false, error: "already_paired" });
|
|
const body = await readJsonBody(request, 1024);
|
|
const token = String(body?.connectorAccessToken || "").trim();
|
|
if (!isSecret(token)) return writeJson(response, 400, { ok: false, error: "connector_access_token_invalid" });
|
|
await persistConnectorToken(token);
|
|
console.log(JSON.stringify({ event: "amd_connector_paired" }));
|
|
return writeJson(response, 201, { ok: true, state: "paired" });
|
|
}
|
|
|
|
async function forwardCesiumRequest(request, response, requestUrl) {
|
|
if (!isMapAuthorized(request.headers["x-proxy-token"])) return writeJson(response, 401, { ok: false, error: "map_egress_unauthorized" });
|
|
const target = parseTarget(requestUrl.searchParams.get("url") || "");
|
|
const startedAt = Date.now();
|
|
const terminal = transport.beginRequest(target.hostname);
|
|
const controller = new AbortController();
|
|
let telemetry = emptyTelemetry();
|
|
let upstream;
|
|
let responseBytes = 0;
|
|
let upstreamEnded = false;
|
|
let downstreamFinished = false;
|
|
let clientAborted = false;
|
|
let terminalFinished = false;
|
|
let bodyIdleTimer;
|
|
let abortUpstream;
|
|
const clearBodyIdleTimer = () => {
|
|
if (bodyIdleTimer) clearTimeout(bodyIdleTimer);
|
|
bodyIdleTimer = undefined;
|
|
};
|
|
const cleanup = () => {
|
|
clearBodyIdleTimer();
|
|
request.off("aborted", onClientAbort);
|
|
response.off("close", onClientAbort);
|
|
if (abortUpstream) controller.signal.removeEventListener("abort", abortUpstream);
|
|
};
|
|
const finish = (outcome) => {
|
|
if (terminalFinished) return false;
|
|
terminalFinished = terminal.complete({
|
|
...outcome,
|
|
durationMs: Date.now() - startedAt,
|
|
bytes: responseBytes,
|
|
telemetry,
|
|
});
|
|
if (terminalFinished) cleanup();
|
|
return terminalFinished;
|
|
};
|
|
function onClientAbort() {
|
|
if (downstreamFinished || terminalFinished) return;
|
|
clientAborted = true;
|
|
if (!controller.signal.aborted) controller.abort(abortError());
|
|
// Before headers, requestOnce owns the active/queued ClientRequest and its
|
|
// signal. The await/catch below records its complete attempt telemetry.
|
|
if (upstream) finish({ type: "client-abort" });
|
|
}
|
|
const armBodyIdleTimer = () => {
|
|
clearBodyIdleTimer();
|
|
bodyIdleTimer = setTimeout(() => {
|
|
const error = typedError("amd_upstream_body_idle_timeout", 502);
|
|
error.code = "ETIMEDOUT";
|
|
upstream.destroy(error);
|
|
}, config.bodyIdleTimeoutMs);
|
|
bodyIdleTimer.unref?.();
|
|
};
|
|
request.once("aborted", onClientAbort);
|
|
response.once("close", onClientAbort);
|
|
|
|
try {
|
|
const connectorToken = await readConnectorToken();
|
|
if (clientAborted || controller.signal.aborted) {
|
|
finish({ type: "client-abort" });
|
|
return;
|
|
}
|
|
if (!connectorToken) {
|
|
finish({ type: "request-failure", errorCode: "amd_connector_not_paired" });
|
|
return writeJson(response, 503, { ok: false, error: "amd_connector_not_paired" }, { "x-nodedc-egress-error": "amd_connector_not_paired" });
|
|
}
|
|
const result = await requestThroughAmd(target, request.method || "GET", upstreamHeaders(request), controller.signal);
|
|
upstream = result.upstream;
|
|
telemetry = result.telemetry;
|
|
if (clientAborted || controller.signal.aborted) {
|
|
upstream.destroy();
|
|
finish({ type: "client-abort" });
|
|
return;
|
|
}
|
|
abortUpstream = () => {
|
|
if (!upstreamEnded && !upstream.destroyed) upstream.destroy(abortError());
|
|
};
|
|
controller.signal.addEventListener("abort", abortUpstream, { once: true });
|
|
if (controller.signal.aborted) abortUpstream();
|
|
|
|
upstream.on("data", (chunk) => {
|
|
responseBytes += chunk.length;
|
|
armBodyIdleTimer();
|
|
});
|
|
upstream.once("end", () => {
|
|
upstreamEnded = true;
|
|
clearBodyIdleTimer();
|
|
});
|
|
upstream.once("error", (error) => {
|
|
upstreamEnded = true;
|
|
clearBodyIdleTimer();
|
|
if (!finish({ type: clientAborted ? "client-abort" : "stream-failure", errorCode: safeTransportError(error) })) return;
|
|
if (!clientAborted) {
|
|
console.warn(JSON.stringify({ event: "cesium_egress_stream_failed", host: target.hostname, error: safeTransportError(error), bytes: responseBytes }));
|
|
}
|
|
if (!response.destroyed) response.destroy();
|
|
});
|
|
armBodyIdleTimer();
|
|
response.writeHead(upstream.statusCode || 502, safeResponseHeaders(upstream.headers));
|
|
response.once("finish", () => {
|
|
downstreamFinished = true;
|
|
clearBodyIdleTimer();
|
|
const durationMs = Date.now() - startedAt;
|
|
const status = Number(upstream.statusCode || 502);
|
|
if (!finish({ type: "response", status })) return;
|
|
// Normal tile traffic is intentionally not logged per object. Slow or
|
|
// failing responses are sufficient to diagnose VPN/transport regressions
|
|
// without creating a noisy, credential-bearing request log.
|
|
if (durationMs >= config.slowRequestMs || status >= 400) {
|
|
console.warn(JSON.stringify({
|
|
event: "cesium_egress_slow_or_failed",
|
|
host: target.hostname,
|
|
status,
|
|
durationMs,
|
|
queueMs: telemetry.queueMs,
|
|
connectionMs: telemetry.connectionMs,
|
|
ttfbMs: telemetry.ttfbMs,
|
|
retryMs: telemetry.retryMs,
|
|
attempts: telemetry.attempts,
|
|
bytes: responseBytes,
|
|
reusedSocket: telemetry.reusedSocket,
|
|
retries: telemetry.retries,
|
|
}));
|
|
}
|
|
});
|
|
upstream.pipe(response);
|
|
} catch (error) {
|
|
telemetry = mergeTelemetry(telemetry, error?.telemetry);
|
|
if (clientAborted || controller.signal.aborted) {
|
|
finish({ type: "client-abort" });
|
|
return;
|
|
}
|
|
if (upstream && !upstream.destroyed) upstream.destroy();
|
|
finish({ type: "request-failure", errorCode: safeTransportError(error) });
|
|
error.telemetry = telemetry;
|
|
throw error;
|
|
}
|
|
}
|
|
|
|
async function requestThroughAmd(initialTarget, method, headers, signal) {
|
|
let target = initialTarget;
|
|
let requestHeaders = { ...headers };
|
|
let aggregate = emptyTelemetry();
|
|
for (let redirects = 0; redirects <= 5; redirects += 1) {
|
|
let result;
|
|
try {
|
|
result = await requestWithSafeRetry(target, method, requestHeaders, signal);
|
|
} catch (error) {
|
|
error.telemetry = mergeTelemetry(aggregate, error?.telemetry);
|
|
throw error;
|
|
}
|
|
const { upstream } = result;
|
|
aggregate = mergeTelemetry(aggregate, result.telemetry);
|
|
const statusCode = Number(upstream.statusCode || 502);
|
|
const location = Array.isArray(upstream.headers.location) ? upstream.headers.location[0] : upstream.headers.location;
|
|
if ([301, 302, 303, 307, 308].includes(statusCode) && location) {
|
|
try {
|
|
aggregate.attemptMs += await drainRedirect(upstream, signal);
|
|
const redirectedTarget = parseTarget(new URL(location, target).toString());
|
|
if (redirectedTarget.origin !== target.origin) requestHeaders = withoutAuthorization(requestHeaders);
|
|
target = redirectedTarget;
|
|
} catch (error) {
|
|
aggregate.attemptMs += Math.max(0, Number(error?.redirectDrainMs || 0));
|
|
error.telemetry = aggregate;
|
|
throw error;
|
|
}
|
|
continue;
|
|
}
|
|
return { upstream, telemetry: aggregate };
|
|
}
|
|
const error = typedError("amd_upstream_redirect_limit", 502);
|
|
error.telemetry = aggregate;
|
|
throw error;
|
|
}
|
|
|
|
function withoutAuthorization(headers) {
|
|
const output = { ...headers };
|
|
delete output.authorization;
|
|
return output;
|
|
}
|
|
|
|
function drainRedirect(upstream, signal) {
|
|
return new Promise((resolve, reject) => {
|
|
const startedAt = Date.now();
|
|
let finished = false;
|
|
let idleTimer;
|
|
const cleanup = () => {
|
|
if (idleTimer) clearTimeout(idleTimer);
|
|
signal.removeEventListener("abort", onAbort);
|
|
upstream.off("data", armIdleTimer);
|
|
upstream.off("end", onEnd);
|
|
upstream.off("aborted", onUpstreamAbort);
|
|
upstream.off("error", onError);
|
|
};
|
|
const complete = () => {
|
|
if (finished) return;
|
|
finished = true;
|
|
cleanup();
|
|
resolve(Date.now() - startedAt);
|
|
};
|
|
const fail = (error) => {
|
|
if (finished) return;
|
|
finished = true;
|
|
cleanup();
|
|
error.redirectDrainMs = Date.now() - startedAt;
|
|
upstream.destroy();
|
|
reject(error);
|
|
};
|
|
const armIdleTimer = () => {
|
|
if (idleTimer) clearTimeout(idleTimer);
|
|
idleTimer = setTimeout(() => {
|
|
const error = typedError("amd_upstream_redirect_body_idle_timeout", 502);
|
|
error.code = "ETIMEDOUT";
|
|
fail(error);
|
|
}, config.bodyIdleTimeoutMs);
|
|
idleTimer.unref?.();
|
|
};
|
|
const onAbort = () => fail(abortReason(signal));
|
|
const onEnd = () => complete();
|
|
const onUpstreamAbort = () => {
|
|
const error = typedError("amd_upstream_redirect_aborted", 502);
|
|
error.code = "ECONNRESET";
|
|
fail(error);
|
|
};
|
|
const onError = (error) => fail(error);
|
|
signal.addEventListener("abort", onAbort, { once: true });
|
|
upstream.on("data", armIdleTimer);
|
|
upstream.once("end", onEnd);
|
|
upstream.once("aborted", onUpstreamAbort);
|
|
upstream.once("error", onError);
|
|
if (signal.aborted) onAbort();
|
|
else {
|
|
armIdleTimer();
|
|
upstream.resume();
|
|
}
|
|
});
|
|
}
|
|
|
|
async function requestWithSafeRetry(target, method, headers, signal) {
|
|
let retries = 0;
|
|
let aggregate = emptyTelemetry();
|
|
while (true) {
|
|
try {
|
|
const result = await requestOnce(target, method, headers, signal);
|
|
aggregate = mergeTelemetry(aggregate, result.telemetry);
|
|
aggregate.retries = retries;
|
|
if (retries > 0) {
|
|
console.warn(JSON.stringify({
|
|
event: "cesium_egress_retry_completed",
|
|
host: target.hostname,
|
|
attempt: retries + 1,
|
|
status: Number(result.upstream.statusCode || 502),
|
|
...telemetryLogFields(result.telemetry),
|
|
totalAttemptMs: aggregate.attemptMs,
|
|
}));
|
|
}
|
|
return { upstream: result.upstream, telemetry: aggregate };
|
|
} catch (error) {
|
|
const failedAttemptTelemetry = error?.telemetry;
|
|
aggregate = mergeTelemetry(aggregate, failedAttemptTelemetry);
|
|
if (signal.aborted || retries >= 1 || !["GET", "HEAD"].includes(method) || error?.reusedSocket !== true || !isRetryableSocketError(error)) {
|
|
if (retries > 0) {
|
|
console.warn(JSON.stringify({
|
|
event: "cesium_egress_retry_failed",
|
|
host: target.hostname,
|
|
attempt: retries + 1,
|
|
error: safeTransportError(error),
|
|
...telemetryLogFields(failedAttemptTelemetry),
|
|
totalAttemptMs: aggregate.attemptMs,
|
|
}));
|
|
}
|
|
error.telemetry = aggregate;
|
|
throw error;
|
|
}
|
|
retries += 1;
|
|
aggregate.retryMs += Number(error?.telemetry?.attemptMs || 0);
|
|
aggregate.retries = retries;
|
|
const errorCode = safeTransportError(error);
|
|
transport.recordRetry(target.hostname, errorCode);
|
|
console.warn(JSON.stringify({
|
|
event: "cesium_egress_reused_socket_retry",
|
|
host: target.hostname,
|
|
attempt: retries,
|
|
error: errorCode,
|
|
...telemetryLogFields(failedAttemptTelemetry),
|
|
}));
|
|
}
|
|
}
|
|
}
|
|
|
|
function requestOnce(target, method, headers, externalSignal) {
|
|
return new Promise((resolve, reject) => {
|
|
const startedAt = Date.now();
|
|
let socketAssignedAt = null;
|
|
let socketTiming = null;
|
|
let reusedSocket = false;
|
|
let settled = false;
|
|
let timedOut = false;
|
|
const pendingConnectionTiming = { startedAt: null, readyAt: null };
|
|
const attemptController = new AbortController();
|
|
const abortAttempt = () => {
|
|
if (!attemptController.signal.aborted) attemptController.abort(abortReason(externalSignal));
|
|
};
|
|
if (externalSignal.aborted) abortAttempt();
|
|
else externalSignal.addEventListener("abort", abortAttempt, { once: true });
|
|
const client = https.request({
|
|
protocol: "https:",
|
|
hostname: target.hostname,
|
|
port: 443,
|
|
method,
|
|
path: `${target.pathname || "/"}${target.search || ""}`,
|
|
headers,
|
|
agent: transport.agent,
|
|
signal: attemptController.signal,
|
|
nodedcAbortSignal: attemptController.signal,
|
|
nodedcConnectionTiming: pendingConnectionTiming,
|
|
});
|
|
transport.observeQueuePeak();
|
|
queueMicrotask(() => transport.observeQueuePeak());
|
|
const timeoutError = typedError("amd_upstream_timeout", 502);
|
|
timeoutError.code = "ETIMEDOUT";
|
|
const timeout = setTimeout(() => {
|
|
timedOut = true;
|
|
if (!attemptController.signal.aborted) attemptController.abort(timeoutError);
|
|
}, config.upstreamTimeoutMs);
|
|
timeout.unref?.();
|
|
const cleanup = () => {
|
|
clearTimeout(timeout);
|
|
externalSignal.removeEventListener("abort", abortAttempt);
|
|
};
|
|
const attemptTelemetry = (endedAt, receivedHeaders, error) => {
|
|
const failedConnectionTiming = error?.[connectionTimingSymbol];
|
|
const timing = socketTiming || failedConnectionTiming || (pendingConnectionTiming.startedAt ? pendingConnectionTiming : null);
|
|
const newConnection = !reusedSocket && timing;
|
|
const queueEnd = newConnection?.startedAt || socketAssignedAt || endedAt;
|
|
const connectionMs = newConnection ? Math.max(0, Number(newConnection.readyAt || endedAt) - Number(newConnection.startedAt || endedAt)) : 0;
|
|
return {
|
|
queueMs: Math.max(0, Number(queueEnd) - startedAt),
|
|
connectionMs,
|
|
ttfbMs: receivedHeaders && socketAssignedAt ? Math.max(0, endedAt - socketAssignedAt) : 0,
|
|
attemptMs: Math.max(0, endedAt - startedAt),
|
|
retryMs: 0,
|
|
attempts: 1,
|
|
retries: 0,
|
|
reusedSocket,
|
|
};
|
|
};
|
|
client.once("socket", (socket) => {
|
|
socketAssignedAt = Date.now();
|
|
socketTiming = socket[connectionTimingSymbol] || null;
|
|
const priorAssignments = Number(socketTiming?.assignments || 0);
|
|
reusedSocket = client.reusedSocket === true || priorAssignments > 0;
|
|
if (socketTiming) socketTiming.assignments = priorAssignments + 1;
|
|
});
|
|
client.once("error", (rawError) => {
|
|
if (settled) return;
|
|
settled = true;
|
|
cleanup();
|
|
const error = timedOut && !externalSignal.aborted ? timeoutError : externalSignal.aborted ? abortReason(externalSignal) : rawError;
|
|
if (rawError?.[connectionTimingSymbol] && !error[connectionTimingSymbol]) error[connectionTimingSymbol] = rawError[connectionTimingSymbol];
|
|
error.reusedSocket = reusedSocket || client.reusedSocket === true;
|
|
error.telemetry = attemptTelemetry(Date.now(), false, error);
|
|
reject(error);
|
|
});
|
|
client.once("response", (upstream) => {
|
|
if (settled) {
|
|
upstream.destroy();
|
|
return;
|
|
}
|
|
settled = true;
|
|
cleanup();
|
|
const endedAt = Date.now();
|
|
resolve({
|
|
upstream,
|
|
telemetry: attemptTelemetry(endedAt, true),
|
|
});
|
|
});
|
|
client.end();
|
|
transport.observeQueuePeak();
|
|
});
|
|
}
|
|
|
|
function emptyTelemetry() {
|
|
return { queueMs: 0, connectionMs: 0, ttfbMs: 0, attemptMs: 0, retryMs: 0, attempts: 0, retries: 0, reusedSocket: false };
|
|
}
|
|
|
|
function telemetryLogFields(telemetry) {
|
|
if (!telemetry) return {};
|
|
return {
|
|
queueMs: Math.max(0, Number(telemetry.queueMs || 0)),
|
|
connectionMs: Math.max(0, Number(telemetry.connectionMs || 0)),
|
|
ttfbMs: Math.max(0, Number(telemetry.ttfbMs || 0)),
|
|
attemptMs: Math.max(0, Number(telemetry.attemptMs || 0)),
|
|
retryMs: Math.max(0, Number(telemetry.retryMs || 0)),
|
|
attempts: Math.max(0, Number(telemetry.attempts || 0)),
|
|
retries: Math.max(0, Number(telemetry.retries || 0)),
|
|
reusedSocket: telemetry.reusedSocket === true,
|
|
};
|
|
}
|
|
|
|
function mergeTelemetry(left, right) {
|
|
const merged = { ...emptyTelemetry(), ...(left || {}) };
|
|
if (!right) return merged;
|
|
for (const name of ["queueMs", "connectionMs", "ttfbMs", "attemptMs", "retryMs", "attempts", "retries"]) {
|
|
if (name === "retries") merged[name] = Math.max(Number(merged[name] || 0), Number(right[name] || 0));
|
|
else merged[name] += Math.max(0, Number(right[name] || 0));
|
|
}
|
|
merged.reusedSocket ||= right.reusedSocket === true;
|
|
return merged;
|
|
}
|
|
|
|
function createTransport(runtimeConfig) {
|
|
const metrics = {
|
|
openedTunnels: 0,
|
|
tunnelFailures: 0,
|
|
requests: 0,
|
|
terminalRequests: 0,
|
|
inFlightRequests: 0,
|
|
completedResponses: 0,
|
|
failedRequests: 0,
|
|
slowRequests: 0,
|
|
responseBytes: 0,
|
|
partialBytes: 0,
|
|
reusedSocketRequests: 0,
|
|
retries: 0,
|
|
streamFailures: 0,
|
|
clientAborts: 0,
|
|
maxQueuedRequests: 0,
|
|
timing: {
|
|
queueMs: { total: 0, max: 0 },
|
|
connectionMs: { total: 0, max: 0 },
|
|
ttfbMs: { total: 0, max: 0 },
|
|
attemptMs: { total: 0, max: 0 },
|
|
retryMs: { total: 0, max: 0 },
|
|
durationMs: { total: 0, max: 0 },
|
|
},
|
|
byHost: {},
|
|
lastFailure: null,
|
|
lastFailureAt: null,
|
|
};
|
|
const agent = new https.Agent({
|
|
keepAlive: true,
|
|
keepAliveMsecs: 30_000,
|
|
maxSockets: runtimeConfig.poolMaxSockets,
|
|
maxFreeSockets: runtimeConfig.poolMaxFreeSockets,
|
|
scheduling: "lifo",
|
|
});
|
|
// https.Agent maintains a separate pool per origin. A Cesium/Bing tile
|
|
// therefore reuses an already authenticated NAS -> AMD -> VPN -> TLS path
|
|
// for its host instead of paying a fresh CONNECT and TLS handshake.
|
|
agent.createConnection = (options, callback) => {
|
|
const hostname = String(options.servername || options.hostname || options.host || "").toLowerCase();
|
|
const signal = options.nodedcAbortSignal;
|
|
const pendingConnectionTiming = options.nodedcConnectionTiming;
|
|
void (async () => {
|
|
if (!allowedHosts.has(hostname)) throw typedError("cesium_target_not_allowed", 403);
|
|
const connectorToken = await readConnectorToken();
|
|
if (!connectorToken) throw typedError("amd_connector_not_paired", 503);
|
|
const startedAt = Date.now();
|
|
if (pendingConnectionTiming) pendingConnectionTiming.startedAt = startedAt;
|
|
let connectedAt = startedAt;
|
|
try {
|
|
const tunnel = await openConnectorTunnel({ hostname }, connectorToken, signal);
|
|
connectedAt = Date.now();
|
|
const secureSocket = await openTlsTunnel(tunnel, hostname, signal);
|
|
const readyAt = Date.now();
|
|
if (pendingConnectionTiming) pendingConnectionTiming.readyAt = readyAt;
|
|
secureSocket[connectionTimingSymbol] = { startedAt, readyAt, assignments: 0 };
|
|
metrics.openedTunnels += 1;
|
|
console.log(JSON.stringify({
|
|
event: "cesium_egress_tunnel_opened",
|
|
host: hostname,
|
|
connectorMs: connectedAt - startedAt,
|
|
tlsMs: readyAt - connectedAt,
|
|
}));
|
|
callback(null, secureSocket);
|
|
} catch (error) {
|
|
const readyAt = Date.now();
|
|
if (pendingConnectionTiming) pendingConnectionTiming.readyAt = readyAt;
|
|
error[connectionTimingSymbol] = { startedAt, readyAt, assignments: 0 };
|
|
throw error;
|
|
}
|
|
})().catch((error) => {
|
|
const errorCode = safeError(error);
|
|
if (error?.code !== "ABORT_ERR") {
|
|
metrics.tunnelFailures += 1;
|
|
metrics.lastFailure = errorCode;
|
|
metrics.lastFailureAt = new Date().toISOString();
|
|
}
|
|
callback(error);
|
|
});
|
|
return undefined;
|
|
};
|
|
|
|
return {
|
|
agent,
|
|
close(callback) { agent.destroy(); callback(); },
|
|
beginRequest(host) {
|
|
const hostMetrics = perHostMetrics(metrics, host);
|
|
metrics.requests += 1;
|
|
hostMetrics.requests += 1;
|
|
metrics.inFlightRequests += 1;
|
|
hostMetrics.inFlightRequests += 1;
|
|
let finished = false;
|
|
return {
|
|
complete({ type, status = 0, durationMs = 0, bytes = 0, telemetry = emptyTelemetry(), errorCode = "" }) {
|
|
if (finished) return false;
|
|
finished = true;
|
|
metrics.terminalRequests += 1;
|
|
metrics.inFlightRequests = Math.max(0, metrics.inFlightRequests - 1);
|
|
hostMetrics.terminalRequests += 1;
|
|
hostMetrics.inFlightRequests = Math.max(0, hostMetrics.inFlightRequests - 1);
|
|
const safeBytes = Math.max(0, Number(bytes || 0));
|
|
if (telemetry.reusedSocket) {
|
|
metrics.reusedSocketRequests += 1;
|
|
hostMetrics.reusedSocketRequests += 1;
|
|
}
|
|
observe(metrics.timing.queueMs, telemetry.queueMs);
|
|
observe(metrics.timing.connectionMs, telemetry.connectionMs);
|
|
observe(metrics.timing.ttfbMs, telemetry.ttfbMs);
|
|
observe(metrics.timing.attemptMs, telemetry.attemptMs);
|
|
observe(metrics.timing.retryMs, telemetry.retryMs);
|
|
observe(metrics.timing.durationMs, durationMs);
|
|
if (type === "response") {
|
|
metrics.completedResponses += 1;
|
|
hostMetrics.completedResponses += 1;
|
|
metrics.responseBytes += safeBytes;
|
|
hostMetrics.bytes += safeBytes;
|
|
} else {
|
|
metrics.partialBytes += safeBytes;
|
|
hostMetrics.partialBytes += safeBytes;
|
|
}
|
|
if (type === "client-abort") {
|
|
metrics.clientAborts += 1;
|
|
hostMetrics.clientAborts += 1;
|
|
}
|
|
const failure = type === "stream-failure" || type === "request-failure" || type === "response" && Number(status) >= 400;
|
|
if (failure) {
|
|
metrics.failedRequests += 1;
|
|
hostMetrics.failedRequests += 1;
|
|
if (type === "stream-failure") {
|
|
metrics.streamFailures += 1;
|
|
hostMetrics.streamFailures += 1;
|
|
}
|
|
metrics.lastFailure = errorCode || `upstream_http_${status}`;
|
|
metrics.lastFailureAt = new Date().toISOString();
|
|
hostMetrics.lastTransportError = errorCode || `upstream_http_${status}`;
|
|
}
|
|
if (type !== "client-abort" && durationMs >= runtimeConfig.slowRequestMs) metrics.slowRequests += 1;
|
|
return true;
|
|
},
|
|
};
|
|
},
|
|
recordRetry(host, errorCode) {
|
|
metrics.retries += 1;
|
|
const hostMetrics = perHostMetrics(metrics, host);
|
|
hostMetrics.retries += 1;
|
|
hostMetrics.lastTransportError = errorCode;
|
|
},
|
|
observeQueuePeak() {
|
|
const queuedRequests = socketCount(agent.requests);
|
|
metrics.maxQueuedRequests = Math.max(metrics.maxQueuedRequests, queuedRequests);
|
|
},
|
|
status() {
|
|
const queuedRequests = socketCount(agent.requests);
|
|
metrics.maxQueuedRequests = Math.max(metrics.maxQueuedRequests, queuedRequests);
|
|
return {
|
|
pool: {
|
|
activeSockets: socketCount(agent.sockets),
|
|
idleSockets: socketCount(agent.freeSockets),
|
|
queuedRequests,
|
|
maxQueuedRequests: metrics.maxQueuedRequests,
|
|
maxSockets: runtimeConfig.poolMaxSockets,
|
|
maxFreeSockets: runtimeConfig.poolMaxFreeSockets,
|
|
byOrigin: poolByOrigin(agent),
|
|
},
|
|
metrics: publicTransportMetrics(metrics),
|
|
};
|
|
},
|
|
};
|
|
}
|
|
|
|
function perHostMetrics(metrics, host) {
|
|
const safeHost = allowedHosts.has(String(host || "").toLowerCase()) ? String(host).toLowerCase() : "unknown";
|
|
metrics.byHost[safeHost] ||= {
|
|
requests: 0,
|
|
terminalRequests: 0,
|
|
inFlightRequests: 0,
|
|
completedResponses: 0,
|
|
failedRequests: 0,
|
|
reusedSocketRequests: 0,
|
|
retries: 0,
|
|
streamFailures: 0,
|
|
clientAborts: 0,
|
|
bytes: 0,
|
|
partialBytes: 0,
|
|
lastTransportError: null,
|
|
};
|
|
return metrics.byHost[safeHost];
|
|
}
|
|
|
|
function observe(bucket, rawValue) {
|
|
const value = Math.max(0, Number(rawValue || 0));
|
|
bucket.total += value;
|
|
bucket.max = Math.max(bucket.max, value);
|
|
}
|
|
|
|
function publicTransportMetrics(metrics) {
|
|
const average = (bucket) => metrics.terminalRequests ? Math.round(bucket.total / metrics.terminalRequests) : 0;
|
|
return {
|
|
openedTunnels: metrics.openedTunnels,
|
|
tunnelFailures: metrics.tunnelFailures,
|
|
requests: metrics.requests,
|
|
terminalRequests: metrics.terminalRequests,
|
|
inFlightRequests: metrics.inFlightRequests,
|
|
completedResponses: metrics.completedResponses,
|
|
failedRequests: metrics.failedRequests,
|
|
slowRequests: metrics.slowRequests,
|
|
responseBytes: metrics.responseBytes,
|
|
partialBytes: metrics.partialBytes,
|
|
reusedSocketRequests: metrics.reusedSocketRequests,
|
|
retries: metrics.retries,
|
|
streamFailures: metrics.streamFailures,
|
|
clientAborts: metrics.clientAborts,
|
|
timing: {
|
|
queueMs: { average: average(metrics.timing.queueMs), max: metrics.timing.queueMs.max },
|
|
connectionMs: { average: average(metrics.timing.connectionMs), max: metrics.timing.connectionMs.max },
|
|
ttfbMs: { average: average(metrics.timing.ttfbMs), max: metrics.timing.ttfbMs.max },
|
|
attemptMs: { average: average(metrics.timing.attemptMs), max: metrics.timing.attemptMs.max },
|
|
retryMs: { average: average(metrics.timing.retryMs), max: metrics.timing.retryMs.max },
|
|
durationMs: { average: average(metrics.timing.durationMs), max: metrics.timing.durationMs.max },
|
|
},
|
|
byHost: metrics.byHost,
|
|
lastFailure: metrics.lastFailure,
|
|
lastFailureAt: metrics.lastFailureAt,
|
|
};
|
|
}
|
|
|
|
function poolByOrigin(agent) {
|
|
const origins = {};
|
|
for (const [key, sockets] of Object.entries(agent.sockets)) poolOrigin(origins, key).active += sockets.length;
|
|
for (const [key, sockets] of Object.entries(agent.freeSockets)) poolOrigin(origins, key).idle += sockets.length;
|
|
for (const [key, requests] of Object.entries(agent.requests)) poolOrigin(origins, key).queued += requests.length;
|
|
return origins;
|
|
}
|
|
|
|
function poolOrigin(origins, agentKey) {
|
|
const host = String(agentKey || "").split(":")[0].toLowerCase();
|
|
const safeHost = allowedHosts.has(host) ? host : "unknown";
|
|
origins[safeHost] ||= { active: 0, idle: 0, queued: 0 };
|
|
return origins[safeHost];
|
|
}
|
|
|
|
function openConnectorTunnel(target, connectorToken, signal) {
|
|
return new Promise((resolve, reject) => {
|
|
if (signal?.aborted) {
|
|
reject(abortReason(signal));
|
|
return;
|
|
}
|
|
const socket = connectNet({ host: config.connectorHost, port: config.connectorPort });
|
|
let buffer = Buffer.alloc(0);
|
|
let finished = false;
|
|
const timeoutError = typedError("amd_connector_timeout", 502);
|
|
timeoutError.code = "ETIMEDOUT";
|
|
const timeout = setTimeout(() => fail(timeoutError), config.connectTimeoutMs);
|
|
timeout.unref?.();
|
|
const cleanup = () => {
|
|
clearTimeout(timeout);
|
|
signal?.removeEventListener("abort", onAbort);
|
|
socket.off("connect", onConnect);
|
|
socket.off("error", fail);
|
|
socket.off("data", onData);
|
|
};
|
|
const finish = (value) => {
|
|
if (finished) return;
|
|
finished = true;
|
|
cleanup();
|
|
resolve(value);
|
|
};
|
|
const fail = (error) => {
|
|
if (finished) return;
|
|
finished = true;
|
|
cleanup();
|
|
socket.destroy();
|
|
reject(error);
|
|
};
|
|
const onAbort = () => fail(abortReason(signal));
|
|
const onConnect = () => socket.write(`CONNECT ${target.hostname}:443 HTTP/1.1\r\nHost: ${target.hostname}:443\r\nProxy-Authorization: Bearer ${connectorToken}\r\n\r\n`);
|
|
const onData = (chunk) => {
|
|
buffer = Buffer.concat([buffer, chunk]);
|
|
if (buffer.length > 16 * 1024) return fail(typedError("amd_connector_response_invalid", 502));
|
|
const end = buffer.indexOf("\r\n\r\n");
|
|
if (end < 0) return;
|
|
if (!/^HTTP\/1\.[01] 200\b/.test(buffer.subarray(0, end).toString("ascii"))) return fail(typedError("amd_connector_rejected", 502));
|
|
const remainder = buffer.subarray(end + 4);
|
|
if (remainder.length) socket.unshift(remainder);
|
|
finish(socket);
|
|
};
|
|
signal?.addEventListener("abort", onAbort, { once: true });
|
|
socket.once("connect", onConnect);
|
|
socket.on("error", fail);
|
|
socket.on("data", onData);
|
|
});
|
|
}
|
|
|
|
function openTlsTunnel(socket, hostname, signal) {
|
|
return new Promise((resolve, reject) => {
|
|
if (signal?.aborted) {
|
|
socket.destroy();
|
|
reject(abortReason(signal));
|
|
return;
|
|
}
|
|
const secureSocket = connectTls({ socket, servername: hostname, ALPNProtocols: ["http/1.1"], rejectUnauthorized: !config.allowInsecureTls });
|
|
let finished = false;
|
|
const timeoutError = typedError("amd_upstream_tls_timeout", 502);
|
|
timeoutError.code = "ETIMEDOUT";
|
|
const timeout = setTimeout(() => fail(timeoutError), config.connectTimeoutMs);
|
|
timeout.unref?.();
|
|
const cleanup = () => {
|
|
clearTimeout(timeout);
|
|
signal?.removeEventListener("abort", onAbort);
|
|
secureSocket.off("secureConnect", finish);
|
|
secureSocket.off("error", fail);
|
|
};
|
|
const finish = () => {
|
|
if (finished) return;
|
|
finished = true;
|
|
cleanup();
|
|
resolve(secureSocket);
|
|
};
|
|
const fail = (error) => {
|
|
if (finished) return;
|
|
finished = true;
|
|
cleanup();
|
|
secureSocket.destroy();
|
|
reject(error);
|
|
};
|
|
const onAbort = () => fail(abortReason(signal));
|
|
signal?.addEventListener("abort", onAbort, { once: true });
|
|
secureSocket.once("secureConnect", finish);
|
|
secureSocket.once("error", fail);
|
|
});
|
|
}
|
|
|
|
function upstreamHeaders(request) {
|
|
// Do not force `Connection: close`: the shared Agent owns a deliberately
|
|
// bounded keep-alive pool for the approved upstream hosts.
|
|
const headers = { accept: safeHeader(request.headers.accept, "*/*"), "accept-encoding": "identity" };
|
|
for (const name of ["range", "if-none-match", "if-modified-since"]) {
|
|
const value = safeHeader(request.headers[name], "");
|
|
if (value) headers[name] = value;
|
|
}
|
|
const authorization = safeHeader(request.headers["x-nodedc-cesium-authorization"], "");
|
|
if (authorization) headers.authorization = authorization;
|
|
const referer = safeHeader(request.headers["x-nodedc-map-referer"], "");
|
|
if (referer) headers.referer = referer;
|
|
return headers;
|
|
}
|
|
|
|
function safeResponseHeaders(headers) {
|
|
const output = {};
|
|
const skipped = new Set(["connection", "keep-alive", "proxy-authenticate", "proxy-authorization", "te", "trailer", "transfer-encoding", "upgrade"]);
|
|
for (const [name, raw] of Object.entries(headers)) {
|
|
if (skipped.has(name.toLowerCase()) || raw === undefined) continue;
|
|
const value = Array.isArray(raw) ? raw.join(", ") : String(raw);
|
|
if (value.length <= 8192 && !/[\r\n\u0000]/.test(value)) output[name] = value;
|
|
}
|
|
return output;
|
|
}
|
|
|
|
function parseTarget(value) {
|
|
let target;
|
|
try { target = new URL(String(value)); } catch { throw typedError("cesium_target_invalid", 400); }
|
|
if (target.protocol !== "https:" || target.username || target.password || target.port && target.port !== "443" || !allowedHosts.has(target.hostname.toLowerCase())) throw typedError("cesium_target_not_allowed", 403);
|
|
target.hostname = target.hostname.toLowerCase();
|
|
return target;
|
|
}
|
|
|
|
async function readConnectorToken() {
|
|
try {
|
|
const token = (await readFile(config.connectorTokenFile, "utf8")).trim();
|
|
return isSecret(token) ? token : null;
|
|
} catch (error) {
|
|
if (error?.code === "ENOENT") return null;
|
|
throw typedError("connector_pair_state_unreadable", 500);
|
|
}
|
|
}
|
|
|
|
async function persistConnectorToken(token) {
|
|
await mkdir(dirname(config.connectorTokenFile), { recursive: true, mode: 0o700 });
|
|
const temporary = `${config.connectorTokenFile}.${process.pid}.${Date.now()}.tmp`;
|
|
let handle;
|
|
try {
|
|
handle = await open(temporary, "wx", 0o600);
|
|
await handle.writeFile(`${token}\n`, "ascii");
|
|
await handle.sync();
|
|
await handle.close();
|
|
handle = null;
|
|
// link(2) is the compare-and-set: unlike rename(), it cannot overwrite an
|
|
// existing pair state. Two concurrent pair attempts therefore produce one
|
|
// durable secret and one harmless 409 response.
|
|
try {
|
|
await link(temporary, config.connectorTokenFile);
|
|
} catch (error) {
|
|
if (error?.code === "EEXIST") throw typedError("already_paired", 409);
|
|
throw error;
|
|
}
|
|
} finally {
|
|
if (handle) await handle.close();
|
|
await unlink(temporary).catch(() => {});
|
|
}
|
|
}
|
|
|
|
async function readSecretFile(path, label) {
|
|
let value;
|
|
try { value = (await readFile(String(path), "utf8")).trim(); } catch { throw new Error(`${label}_unreadable`); }
|
|
if (!value || value.length > 4096 || /[\u0000-\u001f\u007f\s]/.test(value)) throw new Error(`${label}_invalid`);
|
|
return value;
|
|
}
|
|
|
|
function isMapAuthorized(raw) {
|
|
if (Array.isArray(raw)) return false;
|
|
const candidate = Buffer.from(String(raw || ""), "utf8");
|
|
return candidate.length === config.mapToken.length && timingSafeEqual(candidate, config.mapToken);
|
|
}
|
|
|
|
function readJsonBody(request, maxBytes) {
|
|
return new Promise((resolve, reject) => {
|
|
const chunks = [];
|
|
let bytes = 0;
|
|
request.on("data", (chunk) => {
|
|
bytes += chunk.length;
|
|
if (bytes > maxBytes) return reject(typedError("pair_payload_too_large", 413));
|
|
chunks.push(chunk);
|
|
});
|
|
request.once("error", reject);
|
|
request.once("end", () => {
|
|
try { resolve(JSON.parse(Buffer.concat(chunks).toString("utf8"))); } catch { reject(typedError("pair_payload_invalid", 400)); }
|
|
});
|
|
});
|
|
}
|
|
|
|
function writeJson(response, status, body, extraHeaders = {}) {
|
|
const payload = JSON.stringify(body);
|
|
response.writeHead(status, {
|
|
"content-type": "application/json; charset=utf-8",
|
|
"content-length": Buffer.byteLength(payload),
|
|
"cache-control": "no-store",
|
|
...extraHeaders,
|
|
});
|
|
response.end(payload);
|
|
}
|
|
|
|
function parsePort(raw, fallback) {
|
|
const value = Number(String(raw || fallback).trim());
|
|
if (!Number.isInteger(value) || value < 1024 || value > 65535) throw new Error("invalid_port");
|
|
return value;
|
|
}
|
|
|
|
function parseDuration(raw, fallbackSeconds) {
|
|
const value = Number(String(raw || fallbackSeconds).trim());
|
|
if (!Number.isInteger(value) || value < 1 || value > 120) throw new Error("invalid_duration");
|
|
return value * 1000;
|
|
}
|
|
|
|
function parsePoolSize(raw, fallback) {
|
|
const value = Number(String(raw || fallback).trim());
|
|
if (!Number.isInteger(value) || value < 1 || value > 32) throw new Error("invalid_pool_size");
|
|
return value;
|
|
}
|
|
|
|
function socketCount(table) {
|
|
return Object.values(table).reduce((total, sockets) => total + (Array.isArray(sockets) ? sockets.length : 0), 0);
|
|
}
|
|
|
|
function parseIpv4(raw, label) {
|
|
const value = String(raw || "").trim();
|
|
const octets = value.split(".").map(Number);
|
|
if (octets.length !== 4 || octets.some((part) => !Number.isInteger(part) || part < 0 || part > 255)) throw new Error(`${label}_invalid`);
|
|
return value;
|
|
}
|
|
|
|
function normalizeAddress(value) { return String(value || "").replace(/^::ffff:/, ""); }
|
|
function safeHeader(raw, fallback) {
|
|
if (Array.isArray(raw)) return fallback;
|
|
const value = String(raw || "").trim();
|
|
return value.length <= 4096 && !/[\r\n\u0000]/.test(value) ? value : fallback;
|
|
}
|
|
function parseBoolean(value, fallback) { return value === undefined || value === "" ? fallback : ["1", "true", "yes", "on"].includes(String(value).trim().toLowerCase()); }
|
|
function isSecret(value) { return /^[A-Za-z0-9_-]{48,256}$/.test(value); }
|
|
function typedError(message, statusCode) { const error = new Error(message); error.statusCode = statusCode; return error; }
|
|
function abortError() { const error = typedError("amd_client_aborted", 499); error.code = "ABORT_ERR"; return error; }
|
|
function abortReason(signal) { return signal?.reason instanceof Error ? signal.reason : abortError(); }
|
|
function safeError(error) { return String(error?.message || "dc_amd_proxy_error").replace(/[^A-Za-z0-9_.:-]/g, "_").slice(0, 120); }
|
|
function isRetryableSocketError(error) { return new Set(["ECONNRESET", "EPIPE", "ETIMEDOUT", "ECONNABORTED"]).has(String(error?.code || "").toUpperCase()); }
|
|
function safeTransportError(error) {
|
|
const code = String(error?.code || "").trim().toLowerCase();
|
|
if (/^[a-z0-9_]{1,48}$/.test(code)) return `amd_upstream_${code}`;
|
|
return safeError(error);
|
|
}
|