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); }