NODEDC_PLATFORM/services/dc-amd-proxy/scripts/smoke-connection-pool.mjs

500 lines
22 KiB
JavaScript

import assert from "node:assert/strict";
import { execFile as execFileCallback, spawn } from "node:child_process";
import { once } from "node:events";
import { mkdtemp, readFile, rm, writeFile } from "node:fs/promises";
import { createServer as createHttpServer, request as createHttpRequest } from "node:http";
import { createServer as createHttpsServer } from "node:https";
import { connect, createServer as createNetServer } from "node:net";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { promisify } from "node:util";
const execFile = promisify(execFileCallback);
const root = await mkdtemp(join(tmpdir(), "nodedc-amd-pool-smoke-"));
const upstreamPort = await freePort();
const connectorPort = await freePort();
const proxyPort = await freePort();
const upstreamSockets = new Set();
const connectorSockets = new Set();
let upstream;
let connector;
let proxy;
let proxyOutput = "";
let tunnelCount = 0;
let upstreamRequests = 0;
let retryAttempts = 0;
let nonReusedResetAttempts = 0;
let retryLimitAttempts = 0;
let retryAbortAttempts = 0;
let queuedAbortAttempts = 0;
let nextConnectorDelayMs = 15;
const redirectCredentials = {};
let holdReleased = false;
let releaseHold;
let markHoldStarted;
let markPreHeaderSeen;
let markPreHeaderClosed;
let markIdleClosed;
let markDelayedConnectorSeen;
let markDelayedConnectorClosed;
let markRedirectSeen;
let markRedirectClosed;
let markRetryAbortSecondSeen;
let markRetryAbortSecondClosed;
const holdRelease = new Promise((resolve) => { releaseHold = resolve; });
const holdStarted = new Promise((resolve) => { markHoldStarted = resolve; });
const preHeaderSeen = new Promise((resolve) => { markPreHeaderSeen = resolve; });
const preHeaderClosed = new Promise((resolve) => { markPreHeaderClosed = resolve; });
const idleClosed = new Promise((resolve) => { markIdleClosed = resolve; });
const delayedConnectorSeen = new Promise((resolve) => { markDelayedConnectorSeen = resolve; });
const delayedConnectorClosed = new Promise((resolve) => { markDelayedConnectorClosed = resolve; });
const redirectSeen = new Promise((resolve) => { markRedirectSeen = resolve; });
const redirectClosed = new Promise((resolve) => { markRedirectClosed = resolve; });
const retryAbortSecondSeen = new Promise((resolve) => { markRetryAbortSecondSeen = resolve; });
const retryAbortSecondClosed = new Promise((resolve) => { markRetryAbortSecondClosed = resolve; });
try {
const keyPath = join(root, "key.pem");
const certificatePath = join(root, "certificate.pem");
await execFile("openssl", [
"req", "-x509", "-newkey", "rsa:2048", "-nodes",
"-keyout", keyPath,
"-out", certificatePath,
"-subj", "/CN=api.cesium.com",
"-days", "1",
]);
upstream = createHttpsServer({ key: await readFile(keyPath), cert: await readFile(certificatePath) }, (request, response) => {
upstreamRequests += 1;
if (request.url === "/retry") {
retryAttempts += 1;
if (retryAttempts === 1) {
setTimeout(() => request.socket.destroy(), 25);
return;
}
response.writeHead(200, { "content-type": "text/plain" });
response.end("retry-ok");
return;
}
if (request.url?.startsWith("/hold/")) {
markHoldStarted();
const finish = () => {
response.writeHead(200, { "content-type": "text/plain" });
response.end(`hold-${request.url.slice("/hold/".length)}`);
};
if (holdReleased) finish();
else void holdRelease.then(finish);
return;
}
if (request.url === "/preheader") {
markPreHeaderSeen();
request.socket.once("close", markPreHeaderClosed);
return;
}
if (request.url === "/redirect-hold") {
markRedirectSeen();
request.socket.once("close", markRedirectClosed);
response.writeHead(302, { location: "https://api.cesium.com/redirect-target", "content-type": "text/plain" });
response.write("redirect-partial");
return;
}
if (request.url === "/redirect-same-origin") {
redirectCredentials.sameSource = request.headers.authorization;
response.writeHead(302, { location: "/same-origin-target" });
response.end("same-origin-redirect");
return;
}
if (request.url === "/same-origin-target") {
redirectCredentials.sameTarget = request.headers.authorization;
response.writeHead(200, { "content-type": "text/plain" });
response.end("same-origin-ok");
return;
}
if (request.url === "/redirect-assets-origin") {
redirectCredentials.assetsSource = request.headers.authorization;
response.writeHead(302, { location: "https://assets.ion.cesium.com/assets-origin-target" });
response.end("assets-origin-redirect");
return;
}
if (request.url === "/assets-origin-target") {
redirectCredentials.assetsHost = request.headers.host;
redirectCredentials.assetsTarget = request.headers.authorization;
response.writeHead(200, { "content-type": "text/plain" });
response.end("assets-origin-ok");
return;
}
if (request.url === "/redirect-bing-origin") {
redirectCredentials.bingSource = request.headers.authorization;
response.writeHead(302, { location: "https://dev.virtualearth.net/bing-origin-target" });
response.end("bing-origin-redirect");
return;
}
if (request.url === "/bing-origin-target") {
redirectCredentials.bingHost = request.headers.host;
redirectCredentials.bingTarget = request.headers.authorization;
response.writeHead(200, { "content-type": "text/plain" });
response.end("bing-origin-ok");
return;
}
if (request.url === "/non-reused-reset") {
nonReusedResetAttempts += 1;
setTimeout(() => request.socket.destroy(), 25);
return;
}
if (request.url === "/queued-abort") {
queuedAbortAttempts += 1;
response.writeHead(200, { "content-type": "text/plain" });
response.end("queue-abort-reached-upstream");
return;
}
if (request.url === "/retry-limit") {
retryLimitAttempts += 1;
setTimeout(() => request.socket.destroy(), 25);
return;
}
if (request.url === "/retry-abort") {
retryAbortAttempts += 1;
if (retryAbortAttempts === 1) {
setTimeout(() => request.socket.destroy(), 25);
return;
}
markRetryAbortSecondSeen();
request.socket.once("close", markRetryAbortSecondClosed);
return;
}
if (request.url === "/idle") {
request.socket.once("close", markIdleClosed);
response.writeHead(200, { "content-type": "text/plain" });
response.write("partial");
return;
}
const responseNumber = upstreamRequests;
const finish = () => {
response.writeHead(200, { "content-type": "text/plain" });
response.end(`tile-${responseNumber}`);
};
if (request.url?.startsWith("/tile/0?")) setTimeout(finish, 25);
else finish();
});
upstream.on("connection", (socket) => {
upstreamSockets.add(socket);
socket.once("close", () => upstreamSockets.delete(socket));
});
upstream.listen(upstreamPort, "127.0.0.1");
await once(upstream, "listening");
connector = createHttpServer((_request, response) => {
response.writeHead(405);
response.end();
});
connector.on("connection", (socket) => {
connectorSockets.add(socket);
socket.once("close", () => connectorSockets.delete(socket));
});
connector.on("connect", (_request, clientSocket, head) => {
tunnelCount += 1;
// The fake connector deliberately ignores the requested hostname and
// connects to a local TLS fixture. Production still enforces the strict
// Cesium/Bing allowlist before opening this tunnel.
const connectorDelayMs = nextConnectorDelayMs;
nextConnectorDelayMs = 15;
if (connectorDelayMs > 100) {
markDelayedConnectorSeen();
clientSocket.once("close", markDelayedConnectorClosed);
}
const target = connect(upstreamPort, "127.0.0.1");
target.once("connect", () => {
setTimeout(() => {
if (clientSocket.destroyed || target.destroyed) return;
clientSocket.write("HTTP/1.1 200 Connection Established\r\n\r\n");
if (head.length) target.write(head);
clientSocket.pipe(target);
target.pipe(clientSocket);
}, connectorDelayMs);
});
target.once("error", () => clientSocket.destroy());
clientSocket.once("error", () => target.destroy());
clientSocket.once("close", () => target.destroy());
});
connector.listen(connectorPort, "127.0.0.1");
await once(connector, "listening");
const connectorToken = `connector-secret-${"a".repeat(48)}`;
const mapToken = "map-pool-smoke-secret";
const providerToken = "provider-secret-marker";
const connectorTokenFile = join(root, "connector-token");
const mapTokenFile = join(root, "map-token");
await writeFile(connectorTokenFile, `${connectorToken}\n`, { mode: 0o600 });
await writeFile(mapTokenFile, `${mapToken}\n`, { mode: 0o600 });
proxy = spawn(process.execPath, ["server.mjs"], {
cwd: new URL("..", import.meta.url),
env: {
...process.env,
NODE_ENV: "test",
PORT: String(proxyPort),
DC_AMD_PROXY_BIND_ADDRESS: "127.0.0.1",
DC_AMD_CONNECTOR_HOST: "127.0.0.1",
DC_AMD_CONNECTOR_PORT: String(connectorPort),
DC_AMD_PAIR_ALLOWED_SOURCE: "127.0.0.1",
DC_AMD_CONNECTOR_TOKEN_FILE: connectorTokenFile,
DC_AMD_MAP_EGRESS_TOKEN_FILE: mapTokenFile,
DC_AMD_PROXY_TEST_ALLOW_INSECURE_TLS: "true",
DC_AMD_POOL_MAX_SOCKETS: "1",
DC_AMD_POOL_MAX_FREE_SOCKETS: "1",
DC_AMD_SLOW_REQUEST_SECONDS: "30",
DC_AMD_UPSTREAM_TIMEOUT_SECONDS: "2",
DC_AMD_BODY_IDLE_TIMEOUT_SECONDS: "1",
},
stdio: ["ignore", "pipe", "pipe"],
});
proxy.stdout.on("data", (chunk) => { proxyOutput += String(chunk); });
proxy.stderr.on("data", (chunk) => { proxyOutput += String(chunk); });
await waitForProxy(proxyPort, proxy, () => proxyOutput);
for (let index = 0; index < 4; index += 1) {
const target = `https://api.cesium.com/tile/${index}${index === 0 ? `?access_token=${providerToken}` : ""}`;
const response = await proxyFetch(proxyPort, mapToken, target);
assert.equal(response.status, 200);
assert.equal(await response.text(), `tile-${index + 1}`);
}
assert.equal(upstreamRequests, 4);
assert.equal(tunnelCount, 1);
const retryResponse = await proxyFetch(proxyPort, mapToken, "https://api.cesium.com/retry");
assert.equal(retryResponse.status, 200);
assert.equal(await retryResponse.text(), "retry-ok");
assert.equal(retryAttempts, 2);
const redirectAuthorization = "Bearer master-api-secret-marker";
const redirectHeaders = { "x-nodedc-cesium-authorization": redirectAuthorization };
const sameOriginResponse = await proxyFetch(proxyPort, mapToken, "https://api.cesium.com/redirect-same-origin", undefined, redirectHeaders);
assert.equal(sameOriginResponse.status, 200);
assert.equal(await sameOriginResponse.text(), "same-origin-ok");
const assetsOriginResponse = await proxyFetch(proxyPort, mapToken, "https://api.cesium.com/redirect-assets-origin", undefined, redirectHeaders);
assert.equal(assetsOriginResponse.status, 200);
assert.equal(await assetsOriginResponse.text(), "assets-origin-ok");
const bingOriginResponse = await proxyFetch(proxyPort, mapToken, "https://api.cesium.com/redirect-bing-origin", undefined, redirectHeaders);
assert.equal(bingOriginResponse.status, 200);
assert.equal(await bingOriginResponse.text(), "bing-origin-ok");
assert.equal(redirectCredentials.sameSource, redirectAuthorization);
assert.equal(redirectCredentials.sameTarget, redirectAuthorization);
assert.equal(redirectCredentials.assetsSource, redirectAuthorization);
assert.equal(redirectCredentials.assetsTarget, undefined);
assert.equal(redirectCredentials.assetsHost, "assets.ion.cesium.com");
assert.equal(redirectCredentials.bingSource, redirectAuthorization);
assert.equal(redirectCredentials.bingTarget, undefined);
assert.equal(redirectCredentials.bingHost, "dev.virtualearth.net");
const held = Array.from({ length: 4 }, (_, index) => (
rawProxyRequest(proxyPort, mapToken, `https://api.cesium.com/hold/${index}`).promise
));
await withTimeout(holdStarted, 1_000, "first saturated request did not reach upstream");
const queuedAbortRequest = rawProxyRequest(proxyPort, mapToken, "https://api.cesium.com/queued-abort");
const saturatedStatus = await waitFor(async () => {
const status = await (await fetch(`http://127.0.0.1:${proxyPort}/status`)).json();
return status.transport.pool.queuedRequests >= 4 ? status : null;
}, 1_000, "agent queue never reached saturation");
assert.equal(saturatedStatus.transport.pool.maxQueuedRequests >= 4, true);
queuedAbortRequest.abort();
await assert.rejects(withTimeout(queuedAbortRequest.promise, 1_000, "queued abort did not reject client request"));
await delay(100);
holdReleased = true;
releaseHold();
const heldResults = await withTimeout(Promise.all(held), 2_000, "saturated requests did not drain");
for (const [index, result] of heldResults.entries()) {
assert.equal(result.status, 200);
assert.equal(result.body, `hold-${index}`);
}
assert.equal(queuedAbortAttempts, 0);
const redirectAbortRequest = rawProxyRequest(proxyPort, mapToken, "https://api.cesium.com/redirect-hold");
await withTimeout(redirectSeen, 1_000, "redirect request did not reach upstream");
redirectAbortRequest.abort();
await assert.rejects(withTimeout(redirectAbortRequest.promise, 1_000, "redirect abort did not reject client request"));
await withTimeout(redirectClosed, 1_000, "redirect abort did not close redirect body socket");
const preHeaderRequest = rawProxyRequest(proxyPort, mapToken, "https://api.cesium.com/preheader");
await withTimeout(preHeaderSeen, 1_000, "pre-header request did not reach upstream");
preHeaderRequest.abort();
await assert.rejects(withTimeout(preHeaderRequest.promise, 1_000, "pre-header abort did not reject client request"));
await withTimeout(preHeaderClosed, 1_000, "pre-header abort did not close upstream socket");
// With no pooled socket left, delay the next CONNECT response and abort the
// browser-side request while that tunnel is still being established.
nextConnectorDelayMs = 500;
const connectAbortRequest = rawProxyRequest(proxyPort, mapToken, "https://api.cesium.com/connect-abort");
await withTimeout(delayedConnectorSeen, 1_000, "CONNECT-abort request did not reach connector");
connectAbortRequest.abort();
await assert.rejects(withTimeout(connectAbortRequest.promise, 1_000, "CONNECT abort did not reject client request"));
await withTimeout(delayedConnectorClosed, 1_000, "CONNECT abort did not close connector socket");
await waitFor(async () => {
const status = await (await fetch(`http://127.0.0.1:${proxyPort}/status`)).json();
return status.transport.metrics.inFlightRequests === 0 ? status : null;
}, 1_000, "CONNECT abort was not terminally accounted");
// The abort above removed the pooled socket. A reset on the next freshly
// opened socket must fail directly, never trigger the safe reused-socket retry.
const nonReusedResponse = await proxyFetch(proxyPort, mapToken, "https://api.cesium.com/non-reused-reset");
assert.equal(nonReusedResponse.status, 502);
assert.equal(nonReusedResetAttempts, 1);
const warmResponse = await proxyFetch(proxyPort, mapToken, "https://api.cesium.com/warm");
assert.equal(warmResponse.status, 200);
await warmResponse.text();
const retryAbortRequest = rawProxyRequest(proxyPort, mapToken, "https://api.cesium.com/retry-abort");
await withTimeout(retryAbortSecondSeen, 1_000, "retry abort did not reach its second attempt");
retryAbortRequest.abort();
await assert.rejects(withTimeout(retryAbortRequest.promise, 1_000, "retry abort did not reject client request"));
await withTimeout(retryAbortSecondClosed, 1_000, "retry abort did not close second-attempt socket");
assert.equal(retryAbortAttempts, 2);
const retryLimitWarmResponse = await proxyFetch(proxyPort, mapToken, "https://api.cesium.com/warm-after-retry-abort");
assert.equal(retryLimitWarmResponse.status, 200);
await retryLimitWarmResponse.text();
// First failure is on a reused socket and is retried once. The second
// failure is terminal: retry remains strictly bounded to two attempts.
const retryLimitResponse = await proxyFetch(proxyPort, mapToken, "https://api.cesium.com/retry-limit");
assert.equal(retryLimitResponse.status, 502);
assert.equal(retryLimitAttempts, 2);
const idleResponse = await proxyFetch(proxyPort, mapToken, "https://api.cesium.com/idle");
assert.equal(idleResponse.status, 200);
await assert.rejects(withTimeout(idleResponse.text(), 2_000, "idle body timeout did not terminate downstream"));
await withTimeout(idleClosed, 1_000, "idle timeout did not close upstream socket");
const status = await (await fetch(`http://127.0.0.1:${proxyPort}/status`)).json();
const metrics = status.transport.metrics;
assert.equal(metrics.requests, 22);
assert.equal(metrics.terminalRequests, 22);
assert.equal(metrics.inFlightRequests, 0);
assert.equal(metrics.completedResponses, 14);
assert.equal(metrics.failedRequests, 3);
assert.equal(metrics.retries, 3);
assert.equal(metrics.streamFailures, 1);
assert.equal(metrics.clientAborts, 5);
assert.equal(metrics.responseBytes > 0, true);
assert.equal(metrics.partialBytes >= Buffer.byteLength("partial"), true);
assert.equal(metrics.reusedSocketRequests >= 8, true);
assert.equal(metrics.timing.queueMs.max >= 100, true);
assert.equal(metrics.timing.connectionMs.max > 0, true);
assert.equal(metrics.timing.ttfbMs.max >= 20, true);
assert.equal(metrics.timing.retryMs.max >= 25, true);
assert.equal(status.transport.pool.queuedRequests, 0);
assert.equal(status.transport.pool.maxQueuedRequests >= 4, true);
assert.equal(status.transport.pool.maxSockets, 1);
assert.equal(status.transport.metrics.byHost["api.cesium.com"].terminalRequests, 22);
assert.equal(proxyOutput.includes('"event":"cesium_egress_reused_socket_retry"'), true);
assert.equal(proxyOutput.includes('"event":"cesium_egress_retry_completed"'), true);
assert.equal(proxyOutput.includes('"event":"cesium_egress_retry_failed"'), true);
for (const secret of [connectorToken, mapToken, providerToken, redirectAuthorization]) assert.equal(proxyOutput.includes(secret), false);
console.log("ok: aborts propagate, queue peaks are retained, retries are bounded, timings are separated, and terminal accounting is exact");
} catch (error) {
if (proxyOutput) console.error(proxyOutput);
throw error;
} finally {
if (proxy && proxy.exitCode === null && proxy.signalCode === null) {
proxy.kill("SIGTERM");
try {
await withTimeout(once(proxy, "exit"), 2_000, "proxy did not stop after SIGTERM");
} catch {
proxy.kill("SIGKILL");
await withTimeout(once(proxy, "exit"), 2_000, "proxy did not stop after SIGKILL").catch(() => undefined);
}
}
await closeServer(connector, connectorSockets);
await closeServer(upstream, upstreamSockets);
await rm(root, { recursive: true, force: true });
}
function proxyFetch(port, token, target, signal, extraHeaders = {}) {
return fetch(`http://127.0.0.1:${port}/proxy/cesium/fetch?url=${encodeURIComponent(target)}`, {
headers: { "x-proxy-token": token, ...extraHeaders },
signal,
});
}
function rawProxyRequest(port, token, target) {
let request;
const promise = new Promise((resolve, reject) => {
request = createHttpRequest({
host: "127.0.0.1",
port,
method: "GET",
path: `/proxy/cesium/fetch?url=${encodeURIComponent(target)}`,
headers: { "x-proxy-token": token },
agent: false,
}, (response) => {
const chunks = [];
response.on("data", (chunk) => chunks.push(chunk));
response.once("end", () => resolve({ status: response.statusCode, body: Buffer.concat(chunks).toString("utf8") }));
response.once("error", reject);
});
request.once("error", reject);
request.end();
});
// Register a rejection observer immediately; the test asserts the same
// original promise after it deliberately destroys the request.
void promise.catch(() => undefined);
return {
promise,
abort() { request.destroy(new Error("fixture_client_abort")); },
};
}
async function freePort() {
const server = createNetServer();
server.listen(0, "127.0.0.1");
await once(server, "listening");
const address = server.address();
assert(address && typeof address === "object");
const closed = once(server, "close");
server.close();
await closed;
return address.port;
}
async function waitForProxy(port, child, output) {
for (let attempt = 0; attempt < 80; attempt += 1) {
if (child.exitCode !== null) throw new Error(`proxy_exited:${child.exitCode}:${output()}`);
try {
const response = await fetch(`http://127.0.0.1:${port}/healthz`);
if (response.ok) return;
} catch { /* service is still starting */ }
await delay(50);
}
throw new Error(`proxy_start_timeout:${output()}`);
}
async function waitFor(read, timeoutMs, message) {
const startedAt = Date.now();
while (Date.now() - startedAt < timeoutMs) {
const value = await read();
if (value) return value;
await delay(20);
}
throw new Error(message);
}
function withTimeout(promise, timeoutMs, message) {
let timeout;
const deadline = new Promise((_, reject) => {
timeout = setTimeout(() => reject(new Error(message)), timeoutMs);
});
return Promise.race([promise, deadline]).finally(() => clearTimeout(timeout));
}
function delay(milliseconds) {
return new Promise((resolve) => setTimeout(resolve, milliseconds));
}
async function closeServer(server, sockets) {
for (const socket of sockets) socket.destroy();
if (!server?.listening) return;
const closed = once(server, "close");
server.close();
await withTimeout(closed, 2_000, "fixture server did not close").catch(() => undefined);
}