From 3d5384291e0fb56619225346a1ed05e48650a0f5 Mon Sep 17 00:00:00 2001 From: JUN Date: Wed, 23 Sep 2026 13:06:14 +0900 Subject: [PATCH 1/5] test(proxy): assert the Windows case-collapsed env for mixed-case ALL_PROXY rows --- tests/server/proxy-env.test.ts | 14 ++++++++++++++ 1 file changed, 14 insertions(+) diff --git a/tests/server/proxy-env.test.ts b/tests/server/proxy-env.test.ts index e853a7a601..881b39af8e 100644 --- a/tests/server/proxy-env.test.ts +++ b/tests/server/proxy-env.test.ts @@ -308,6 +308,10 @@ describe("applyProxyEnv", () => { // socks5h: the proxy resolves the name, so the fixture fails the request without local DNS. process.env[socksKey] = `socks5h://127.0.0.1:${address.port}`; Object.assign(process.env, inherited); + // Windows environment names are case-insensitive: ALL_PROXY and all_proxy are one variable, + // so the opposite-case HTTP value replaces the SOCKS one and no SOCKS proxy is left. + const collapsed = process.platform === "win32" + && Object.keys(inherited).some(key => key !== socksKey && key.toLowerCase() === socksKey.toLowerCase()); applyProxyEnv(configWithProxy()); // Beside an HTTP(S) proxy Bun reads NO_PROXY too, with suffix matching, so no bare localhost. expect(process.env.NO_PROXY).toBe(expectedNoProxy); @@ -321,6 +325,16 @@ describe("applyProxyEnv", () => { try { expect(await (await configuredOutboundFetch(loopbackUrl, undefined, direct)).text()).toBe("direct"); expect(directCalls).toBe(1); + if (collapsed) { + // Only the HTTP proxy remains, so nothing forces direct egress and Bun's own proxy + // environment (with the loopback NO_PROXY above) decides for both hosts. + expect(process.env[socksKey]).toBe("http://proxy.invalid:3128"); + expect(directProxy).toBeUndefined(); + expect(await (await configuredOutboundFetch("http://app.localhost:11434/v1/models", undefined, direct)).text()).toBe("direct"); + expect(directCalls).toBe(2); + expect(directProxy).toBeUndefined(); + return; + } if (forcedDirect) expect(directProxy).toBe(false); await expect(configuredOutboundFetch("http://app.localhost:11434/v1/models", undefined, direct)).rejects.toThrow(); expect(directCalls).toBe(1); From c30e968043183a8216b685fb3d7e8e7ed24d0172 Mon Sep 17 00:00:00 2001 From: JUN Date: Wed, 23 Sep 2026 13:24:08 +0900 Subject: [PATCH 2/5] test(claude): close policy history index before temp-home cleanup --- .../claude-integration/claude-native-affinity.test.ts | 10 +++++++++- 1 file changed, 9 insertions(+), 1 deletion(-) diff --git a/tests/claude-integration/claude-native-affinity.test.ts b/tests/claude-integration/claude-native-affinity.test.ts index 1bd195ae23..c52abc3743 100644 --- a/tests/claude-integration/claude-native-affinity.test.ts +++ b/tests/claude-integration/claude-native-affinity.test.ts @@ -1,5 +1,5 @@ import { afterEach, beforeEach, describe, expect, spyOn, test } from "bun:test"; -import { mkdtempSync, writeFileSync } from "node:fs"; +import { existsSync, mkdtempSync, writeFileSync } from "node:fs"; import { tmpdir } from "node:os"; import { join } from "node:path"; import { clearComboSelectionState, clearComboTargetCooldowns } from "../../src/combos"; @@ -9,6 +9,9 @@ import { handleResponsesWithPolicyFallback, rankPolicyFallbackCandidates } from import { tryAdmitTurn } from "../../src/server/lifecycle"; import { providerConfigSeed } from "../../src/providers/derive"; import { getProviderRegistryEntry } from "../../src/providers/registry"; +import { closeRequestHistoryIndex } from "../../src/routing/history/indexer"; +import { historyIndexPath } from "../../src/routing/history/schema"; +import { clearHealthHistoryCacheForTests } from "../../src/routing/health"; import type { OcxConfig } from "../../src/types"; import { fakeChatGptJwt } from "../helpers/fake-chatgpt-jwt"; import { installIsolatedCodexHome, type IsolatedCodexHome } from "../helpers/isolated-codex-home"; @@ -39,6 +42,9 @@ beforeEach(() => { releaseSpendHome = acquireOwnedSpendHome(); }); afterEach(() => { + // Policy candidate health opens a separate SQLite index under this home. + closeRequestHistoryIndex(); + clearHealthHistoryCacheForTests(); // Released before the directory is removed, so no live database sits inside it. releaseSpendHome?.(); releaseSpendHome = undefined; @@ -121,6 +127,7 @@ describe("Claude final canonical native affinity after a Go preliminary pick", ( }); test("native failure leaves policy-hop request headers free of synthesized identity", async () => { + clearHealthHistoryCacheForTests(); const cfg = config(); cfg.routingProfiles = { "native-hop": { candidates: [{ provider: "openai", model: "gpt-5.6-luna" }] } }; const trace = { version: 1, decisionId: "native-hop", createdAt: Date.now(), requestedModel: "policy/native-hop", @@ -149,6 +156,7 @@ describe("Claude final canonical native affinity after a Go preliminary pick", ( const response = await handleResponsesWithPolicyFallback(req, cfg, { model: "", provider: "" }, { claudeNativeSessionId: expectedSession }, { runCore }); await response.text(); + expect(existsSync(historyIndexPath(home))).toBe(true); expect(response.status).toBe(200); expect(requests).toHaveLength(2); expect(wires[0]?.get("session_id")).toBe(expectedSession); From d3ef4e597731dcc821a0094d61f0f8649b8066ae Mon Sep 17 00:00:00 2001 From: JUN Date: Wed, 23 Sep 2026 13:23:16 +0900 Subject: [PATCH 3/5] test: use WSL paths in service home ownership fixture --- tests/service/service-wsl-home-ownership.test.ts | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/tests/service/service-wsl-home-ownership.test.ts b/tests/service/service-wsl-home-ownership.test.ts index 20c3c19d73..b6abc4f7be 100644 --- a/tests/service/service-wsl-home-ownership.test.ts +++ b/tests/service/service-wsl-home-ownership.test.ts @@ -1,7 +1,7 @@ import { afterEach, describe, expect, test } from "bun:test"; import { mkdtempSync, writeFileSync } from "node:fs"; import { tmpdir } from "node:os"; -import { join } from "node:path"; +import { join, posix } from "node:path"; import { assertServiceEnvironmentMatchesInstall } from "../../src/service/guards"; import { repairService } from "../../src/service/repair"; import { inspectNativeCodexOwnership } from "../../src/integrations/native/ownership-preflight"; @@ -22,9 +22,9 @@ describe("WSL service ownership after Windows home discovery", () => { const root = mkdtempSync(join(tmpdir(), "ocx-wsl-service-home-")); fixtureHome = root; process.env.OPENCODEX_HOME = root; - const linuxHome = join(root, "linux", ".codex"); - const usersRoot = join(root, "windows", "Users"); - const windowsHome = join(usersRoot, "profile", ".codex"); + const linuxHome = "/home/fixture/.codex"; + const usersRoot = "/mnt/c/Users"; + const windowsHome = posix.join(usersRoot, "profile", ".codex"); const recordedHome = recorded === "linux" ? linuxHome : windowsHome; const statePath = join(root, "service-state.json"); writeFileSync(statePath, JSON.stringify({ @@ -35,9 +35,9 @@ describe("WSL service ownership after Windows home discovery", () => { const deps = { env: { WSL_DISTRO_NAME: "fixture" }, platform: "linux" as const, - homedir: () => join(root, "linux"), + homedir: () => "/home/fixture", usersRoot, - existsSync: (path: string) => path === usersRoot || path === join(windowsHome, "config.toml"), + existsSync: (path: string) => path === usersRoot || path === posix.join(windowsHome, "config.toml"), readdirSync: () => ["profile"], statSync: (() => ({ isDirectory: () => true })) as never, realpathSync: (path: string) => path, From a05cb672d4484dd91cd7fd24a934c0694a5b1941 Mon Sep 17 00:00:00 2001 From: JUN Date: Wed, 23 Sep 2026 13:54:32 +0900 Subject: [PATCH 4/5] test: isolate Windows service claim state home --- tests/service/service-claim.test.ts | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/tests/service/service-claim.test.ts b/tests/service/service-claim.test.ts index cb46b64553..d886631522 100644 --- a/tests/service/service-claim.test.ts +++ b/tests/service/service-claim.test.ts @@ -144,6 +144,8 @@ describe("runServiceClaim", () => { test("an unreadable sandbox state refuses a real claim", async () => { const home = createTempHome("ocx-claim-refusal-"); + const previousUserProfile = process.env.USERPROFILE; + if (process.platform === "win32") process.env.USERPROFILE = home.root; try { expect(serviceStatePaths().every(path => path.startsWith(home.root))).toBe(true); mkdirSync(serviceStatePath()); @@ -157,6 +159,10 @@ describe("runServiceClaim", () => { }); expect(statSync(serviceStatePath()).isDirectory()).toBe(true); } finally { + if (process.platform === "win32") { + if (previousUserProfile === undefined) delete process.env.USERPROFILE; + else process.env.USERPROFILE = previousUserProfile; + } home.remove(); } }); From a660c626340fadbe80b40b4a9ab3514400411d57 Mon Sep 17 00:00:00 2001 From: JUN Date: Wed, 23 Sep 2026 13:52:31 +0900 Subject: [PATCH 5/5] Fix Aside sync deadline and direct socket teardown --- src/cli/aside-profiles.ts | 86 ++++++----- src/server/direct-local-http.ts | 64 +++++--- .../data-plane-admission-identity.test.ts | 20 ++- .../local-aside-sync-capability.test.ts | 40 ++++- .../local-management-direct-transport.test.ts | 137 ++++++++++++++++++ 5 files changed, 283 insertions(+), 64 deletions(-) diff --git a/src/cli/aside-profiles.ts b/src/cli/aside-profiles.ts index d7817b5a87..892cc8e4e1 100644 --- a/src/cli/aside-profiles.ts +++ b/src/cli/aside-profiles.ts @@ -9,47 +9,61 @@ import { runtimeRequest, RuntimeApiError, type RuntimeApiDeps } from "./runtime- /** Aside policy and file writes share the running server's mutation owner. Never fall back locally. */ export async function refreshAsideProfilesThroughServer( // The optional transport seam observes the real direct-local exchange in listener tests. - deps: RuntimeApiDeps & { directLocalFetch?: typeof directLocalHttpFetch; exchangeDeadlineMs?: number } = {}, + deps: RuntimeApiDeps & { + directLocalFetch?: typeof directLocalHttpFetch; + exchangeDeadlineMs?: number; + scheduleExchangeDeadline?: (onTimeout: () => void, delayMs: number) => () => void; + } = {}, ): Promise { // An explicit URL is an opt-in transport used by connected callers and tests. if (!deps.baseUrl) { const localFetch = deps.directLocalFetch ?? directLocalHttpFetch; - const deadline = AbortSignal.timeout(deps.exchangeDeadlineMs ?? LOCAL_ASIDE_SYNC_CAPABILITY_TTL_MS); - const live = await (deps.findLiveProxy ?? findLiveProxy)(); - if (!live) throw new RuntimeApiError("Proxy is not running. Start it with: ocx start", 503, null); - if (live.source !== "runtime" || live.pid === null) { - throw new RuntimeApiError("Aside profile synchronization requires an attested running proxy", 503, null); + const controller = new AbortController(); + const onTimeout = () => controller.abort(new DOMException("The Aside sync exchange timed out", "TimeoutError")); + const cancelDeadline = (deps.scheduleExchangeDeadline ?? ((callback, delayMs) => { + const timer = setTimeout(callback, delayMs); + return () => clearTimeout(timer); + }))(onTimeout, deps.exchangeDeadlineMs ?? LOCAL_ASIDE_SYNC_CAPABILITY_TTL_MS); + const deadline = controller.signal; + try { + const live = await (deps.findLiveProxy ?? findLiveProxy)(); + if (!live) throw new RuntimeApiError("Proxy is not running. Start it with: ocx start", 503, null); + if (live.source !== "runtime" || live.pid === null) { + throw new RuntimeApiError("Aside profile synchronization requires an attested running proxy", 503, null); + } + const runtime = readRuntimePort(live.pid); + if (!runtime?.attestationSecret || runtime.pid !== live.pid || runtime.port !== live.port) { + throw new RuntimeApiError("Aside profile synchronization could not verify the running proxy", 503, null); + } + const nonce = createLocalAttestationChallenge(); + const baseUrl = `http://${probeHostname(live.hostname)}:${live.port}`; + const proofResponse = await localFetch(`${baseUrl}/healthz`, { headers: { [LOCAL_ATTESTATION_CHALLENGE_HEADER]: nonce }, signal: deadline }); + const health = await proofResponse.json().catch(() => null); + if (!proofResponse.ok || !isOpencodexHealthz(health) || health?.pid !== live.pid || health?.port !== live.port + || health?.asideSyncCapability !== LOCAL_ASIDE_SYNC_CAPABILITY_VERSION + || !verifyLocalAttestationProof(runtime.attestationSecret, nonce, live.pid, live.port, proofResponse.headers.get(LOCAL_ATTESTATION_PROOF_HEADER))) { + throw new RuntimeApiError("Aside profile synchronization could not attest the running proxy", 503, null); + } + const expiresAt = Date.now() + LOCAL_ASIDE_SYNC_CAPABILITY_TTL_MS; + const capability = createLocalAsideSyncCapability(runtime.attestationSecret, nonce, LOCAL_ASIDE_SYNC_METHOD, LOCAL_ASIDE_SYNC_PATH, live.pid, live.port, expiresAt); + if (!capability) throw new RuntimeApiError("Aside profile synchronization capability was unavailable", 503, null); + const response = await localFetch(`${baseUrl}${LOCAL_ASIDE_SYNC_PATH}`, { + method: LOCAL_ASIDE_SYNC_METHOD, + signal: deadline, + headers: { + [LOCAL_ASIDE_SYNC_EXPECTED_PID_HEADER]: String(live.pid), + [LOCAL_ASIDE_SYNC_NONCE_HEADER]: nonce, + [LOCAL_ASIDE_SYNC_EXPIRES_AT_HEADER]: String(expiresAt), + [LOCAL_ASIDE_SYNC_CAPABILITY_HEADER]: capability, + }, + }); + const body = await response.json().catch(() => null) as { results?: OwnedIntegrationRefreshOutcome[] } | null; + if (!response.ok) throw new RuntimeApiError("Aside profile synchronization was rejected", response.status, body); + if (!Array.isArray(body?.results)) throw new RuntimeApiError("The running proxy does not support Aside profile synchronization", 502, body); + return body.results; + } finally { + cancelDeadline(); } - const runtime = readRuntimePort(live.pid); - if (!runtime?.attestationSecret || runtime.pid !== live.pid || runtime.port !== live.port) { - throw new RuntimeApiError("Aside profile synchronization could not verify the running proxy", 503, null); - } - const nonce = createLocalAttestationChallenge(); - const baseUrl = `http://${probeHostname(live.hostname)}:${live.port}`; - const proofResponse = await localFetch(`${baseUrl}/healthz`, { headers: { [LOCAL_ATTESTATION_CHALLENGE_HEADER]: nonce }, signal: deadline }); - const health = await proofResponse.json().catch(() => null); - if (!proofResponse.ok || !isOpencodexHealthz(health) || health?.pid !== live.pid || health?.port !== live.port - || health?.asideSyncCapability !== LOCAL_ASIDE_SYNC_CAPABILITY_VERSION - || !verifyLocalAttestationProof(runtime.attestationSecret, nonce, live.pid, live.port, proofResponse.headers.get(LOCAL_ATTESTATION_PROOF_HEADER))) { - throw new RuntimeApiError("Aside profile synchronization could not attest the running proxy", 503, null); - } - const expiresAt = Date.now() + LOCAL_ASIDE_SYNC_CAPABILITY_TTL_MS; - const capability = createLocalAsideSyncCapability(runtime.attestationSecret, nonce, LOCAL_ASIDE_SYNC_METHOD, LOCAL_ASIDE_SYNC_PATH, live.pid, live.port, expiresAt); - if (!capability) throw new RuntimeApiError("Aside profile synchronization capability was unavailable", 503, null); - const response = await localFetch(`${baseUrl}${LOCAL_ASIDE_SYNC_PATH}`, { - method: LOCAL_ASIDE_SYNC_METHOD, - signal: deadline, - headers: { - [LOCAL_ASIDE_SYNC_EXPECTED_PID_HEADER]: String(live.pid), - [LOCAL_ASIDE_SYNC_NONCE_HEADER]: nonce, - [LOCAL_ASIDE_SYNC_EXPIRES_AT_HEADER]: String(expiresAt), - [LOCAL_ASIDE_SYNC_CAPABILITY_HEADER]: capability, - }, - }); - const body = await response.json().catch(() => null) as { results?: OwnedIntegrationRefreshOutcome[] } | null; - if (!response.ok) throw new RuntimeApiError("Aside profile synchronization was rejected", response.status, body); - if (!Array.isArray(body?.results)) throw new RuntimeApiError("The running proxy does not support Aside profile synchronization", 502, body); - return body.results; } const result = await runtimeRequest<{ results?: OwnedIntegrationRefreshOutcome[] }>( "/api/client-integrations/aside/sync", diff --git a/src/server/direct-local-http.ts b/src/server/direct-local-http.ts index 976c1af92b..f3891ebf10 100644 --- a/src/server/direct-local-http.ts +++ b/src/server/direct-local-http.ts @@ -6,6 +6,7 @@ const DIRECT_LOCAL_HTTP_TIMEOUT_MS = 10_000; type DirectLocalHttpIo = { timeoutMs?: number; connect?: (hostname: string, port: number) => Socket; + scheduleDeadline?: (onTimeout: () => void, timeoutMs: number) => () => void; }; function abortReason(signal: AbortSignal): Error { @@ -270,18 +271,22 @@ export async function directLocalHttpFetch( return await new Promise((resolve, reject) => { let socket: Socket | undefined; let settled = false; + let closed = false; + let completed = false; + let outcomeError: Error | undefined; + let cancelDeadline: (() => void) | undefined; let receivedBytes = 0; let responseBytes = Buffer.allocUnsafe(4 * 1024); let framing: ResponseFraming = { kind: "head", searchFrom: 0 }; - const cleanup = () => { + const disableSocketTimeout = () => { socket?.setTimeout(0); }; + const settleAfterClose = (deadlineError?: Error) => { + if (!settled || (!closed && !deadlineError) || completed) return; + completed = true; + // Framing may finish before close, so the caller's deadline owns teardown too. signal?.removeEventListener("abort", onAbort); - socket?.setTimeout(0); - }; - const finish = (error?: Error) => { - if (settled) return; - settled = true; - cleanup(); - try { socket?.destroy(); } catch { /* ignore */ } + cancelDeadline?.(); + cancelDeadline = undefined; + const error = outcomeError ?? deadlineError; if (error) { reject(error); return; @@ -292,9 +297,28 @@ export async function directLocalHttpFetch( reject(parseError instanceof Error ? parseError : new Error(String(parseError))); } }; + const finish = (error?: Error) => { + if (settled) return; + settled = true; + outcomeError = error; + disableSocketTimeout(); + try { socket?.destroy(); } catch (destroyError) { + settleAfterClose(error ?? (destroyError instanceof Error ? destroyError : new Error(String(destroyError)))); + return; + } + settleAfterClose(); + }; const onAbort = () => { const error = signal ? abortReason(signal) : new Error("direct local HTTP request aborted"); + if (settled) outcomeError ??= error; + else finish(error); + if (error.name === "TimeoutError") settleAfterClose(error); + }; + const onDeadline = () => { + const error = new Error("direct local HTTP request timed out"); + error.name = "TimeoutError"; finish(error); + settleAfterClose(error); }; socket = (io.connect ?? ((host, selectedPort) => net.createConnection({ @@ -302,17 +326,9 @@ export async function directLocalHttpFetch( port: selectedPort, autoSelectFamily: true, })))(hostname, port); - socket.setTimeout(timeoutMs, () => { - const error = new Error("direct local HTTP request timed out"); - error.name = "TimeoutError"; - finish(error); - }); - signal?.addEventListener("abort", onAbort, { once: true }); - if (signal?.aborted) { - onAbort(); - return; - } + socket.setTimeout(timeoutMs, onDeadline); socket.on("connect", () => { + if (settled) return; try { socket?.write(requestBytes); } catch (error) { finish(error instanceof Error ? error : new Error(String(error))); } @@ -342,6 +358,16 @@ export async function directLocalHttpFetch( }); socket.once("end", () => finish()); socket.once("error", error => finish(error)); - socket.once("close", () => finish()); + socket.once("close", () => { + closed = true; + finish(); + settleAfterClose(); + }); + cancelDeadline = (io.scheduleDeadline ?? ((callback, delayMs) => { + const timer = setTimeout(callback, delayMs); + return () => clearTimeout(timer); + }))(onDeadline, timeoutMs); + signal?.addEventListener("abort", onAbort, { once: true }); + if (signal?.aborted) onAbort(); }); } diff --git a/tests/server/data-plane-admission-identity.test.ts b/tests/server/data-plane-admission-identity.test.ts index 1e79431aa2..75494853bf 100644 --- a/tests/server/data-plane-admission-identity.test.ts +++ b/tests/server/data-plane-admission-identity.test.ts @@ -284,17 +284,27 @@ describe("the Responses WebSocket handshake", () => { headers: { "X-OpenCodex-API-Key": "ocx_data_secondsecret" }, } as unknown as string[]); let settled = false; + let opened = false; + let failed = false; const finish = (value: boolean) => { if (settled) return; settled = true; clearTimeout(timer); - try { socket.close(); } catch { /* already closed */ } resolve(value); }; - socket.addEventListener("open", () => finish(true)); - socket.addEventListener("error", () => finish(false)); - socket.addEventListener("close", () => finish(false)); - const timer = setTimeout(() => finish(false), 5_000); + socket.addEventListener("open", () => { + opened = true; + try { socket.close(); } catch { finish(false); } + }); + socket.addEventListener("error", () => { + failed = true; + try { socket.close(); } catch { finish(false); } + }); + socket.addEventListener("close", () => finish(opened && !failed)); + const timer = setTimeout(() => { + try { socket.close(); } catch { /* already closed */ } + finish(false); + }, 5_000); }); // The handshake now branches on the resolver rather than the boolean // wrapper, so this pins that the rewrite did not change who gets in. diff --git a/tests/server/local-aside-sync-capability.test.ts b/tests/server/local-aside-sync-capability.test.ts index 057b29b6b8..4ab0013541 100644 --- a/tests/server/local-aside-sync-capability.test.ts +++ b/tests/server/local-aside-sync-capability.test.ts @@ -17,7 +17,6 @@ import { import { createLocalAttestationChallenge, LOCAL_ATTESTATION_CHALLENGE_HEADER, LOCAL_ATTESTATION_PROOF_HEADER } from "../../src/lib/local-management-attestation"; import { directLocalHttpFetch } from "../../src/server/direct-local-http"; import { startServer } from "../../src/server"; -import { stopServerListener } from "../../src/server/lifecycle"; import { setIntegrationPathTestHooks } from "../../src/server/management/integration-routes"; import { currentServerFixtureConfig, settleServerAuthFixture } from "../helpers/server-auth-fixture"; import { serverAuthConfig } from "../helpers/server-auth-config"; @@ -54,7 +53,7 @@ beforeEach(() => { }); afterEach(async () => { - if (server) await stopServerListener(server); + if (server) await server.stop(true); server = null; setIntegrationPathTestHooks(null); await settleServerAuthFixture(home, codexHome?.path); @@ -89,13 +88,19 @@ async function send(path: string, headers: Headers, method = "POST"): Promise { + const response = await send(path, headers, method); + await response.arrayBuffer(); + return response.status; +} + test("Aside sync capability admits one exact request and refuses replay or altered bindings", async () => { if (!server) throw new Error("listener is not running"); const headers = signedHeaders(); const accepted = await send(LOCAL_ASIDE_SYNC_PATH, headers); expect(accepted.status).toBe(200); expect((await accepted.json()).results).toEqual([]); - expect((await send(LOCAL_ASIDE_SYNC_PATH, headers)).status).toBe(401); + expect(await sendStatus(LOCAL_ASIDE_SYNC_PATH, headers)).toBe(401); const wrongPort = server.port === 65_535 ? server.port - 1 : server.port + 1; const badHmac = signedHeaders(); @@ -110,7 +115,7 @@ test("Aside sync capability admits one exact request and refuses replay or alter [LOCAL_ASIDE_SYNC_PATH, "POST", signedHeaders({ expiresAt: Date.now() - 1 })], [LOCAL_ASIDE_SYNC_PATH, "POST", badHmac], ] as const) { - expect((await send(path, candidate, method)).status).toBe(401); + expect(await sendStatus(path, candidate, method)).toBe(401); } }); @@ -155,14 +160,38 @@ test("CLI default Aside sync attests the listener and sends a bodyless POST", as expect(posted).toBe(false); }); +test("CLI Aside sync clears its exchange deadline after a successful call", async () => { + if (!server) throw new Error("listener is not running"); + let active = false; + let scheduledMs: number | undefined; + await refreshAsideProfilesThroughServer({ + findLiveProxy: async () => ({ pid: process.pid, port: server!.port, hostname: "127.0.0.1", source: "runtime" }), + exchangeDeadlineMs: 12_345, + scheduleExchangeDeadline: (_onTimeout, delayMs) => { + active = true; + scheduledMs = delayMs; + return () => { active = false; }; + }, + }); + expect(scheduledMs).toBe(12_345); + expect(active).toBe(false); +}); + test("CLI Aside sync applies one absolute deadline across attestation and POST", async () => { if (!server) throw new Error("listener is not running"); const live = { pid: process.pid, port: server.port, hostname: "127.0.0.1", source: "runtime" as const }; let healthSignal: AbortSignal | undefined; let postSignal: AbortSignal | undefined; + let fireDeadline: (() => void) | undefined; + let deadlineActive = false; await expect(refreshAsideProfilesThroughServer({ findLiveProxy: async () => live, exchangeDeadlineMs: 50, + scheduleExchangeDeadline: onTimeout => { + deadlineActive = true; + fireDeadline = onTimeout; + return () => { deadlineActive = false; }; + }, directLocalFetch: async (input, init = {}) => { if (new URL(input instanceof Request ? input.url : String(input)).pathname === "/healthz") { healthSignal = init.signal ?? undefined; @@ -171,9 +200,12 @@ test("CLI Aside sync applies one absolute deadline across attestation and POST", postSignal = init.signal ?? undefined; return new Promise((_, reject) => { init.signal?.addEventListener("abort", () => reject(init.signal?.reason), { once: true }); + if (!fireDeadline) throw new Error("Aside sync did not schedule its deadline"); + fireDeadline(); }); }, })).rejects.toBeDefined(); expect(healthSignal).toBe(postSignal); expect(postSignal?.aborted).toBe(true); + expect(deadlineActive).toBe(false); }); diff --git a/tests/server/local-management-direct-transport.test.ts b/tests/server/local-management-direct-transport.test.ts index 241354faf5..5d8791df35 100644 --- a/tests/server/local-management-direct-transport.test.ts +++ b/tests/server/local-management-direct-transport.test.ts @@ -1,5 +1,6 @@ import { describe, expect, test } from "bun:test"; import { lookup } from "node:dns/promises"; +import { EventEmitter } from "node:events"; import { createServer } from "node:http"; import { createConnection, createServer as createTcpServer, type Server, type Socket } from "node:net"; import { join } from "node:path"; @@ -150,6 +151,142 @@ describe("local management direct transport", () => { } }); + test("an aborted request releases its client socket before it settles", async () => { + let accept!: () => void; + const accepted = new Promise(resolve => { accept = resolve; }); + const sockets = new Set(); + const server = createTcpServer(socket => { + sockets.add(socket); + socket.once("close", () => sockets.delete(socket)); + accept(); + }); + const controller = new AbortController(); + let client: Socket | undefined; + let clientClosed = false; + try { + const port = await listen(server); + const pending = directLocalHttpFetch(`http://127.0.0.1:${port}/healthz`, { + signal: controller.signal, + }, { + connect(hostname, selectedPort) { + client = createConnection({ host: hostname, port: selectedPort }); + client.once("close", () => { clientClosed = true; }); + return client; + }, + }); + await accepted; + controller.abort(); + await expect(pending).rejects.toMatchObject({ name: "AbortError" }); + expect(clientClosed).toBe(true); + expect(client?.destroyed).toBe(true); + } finally { + client?.destroy(); + for (const socket of sockets) socket.destroy(); + await close(server); + } + }); + + test("a never-connected socket without close still settles by its deadline", async () => { + const socket = Object.assign(new EventEmitter(), { + setTimeout() { return this; }, + destroy() { return this; }, + }) as unknown as Socket; + let deadlineActive = false; + let scheduledMs: number | undefined; + await expect(directLocalHttpFetch("http://127.0.0.1:9/healthz", {}, { + timeoutMs: 20, + connect: () => socket, + scheduleDeadline: (onTimeout, delayMs) => { + scheduledMs = delayMs; + deadlineActive = true; + const timer = setTimeout(() => { deadlineActive = false; onTimeout(); }, delayMs); + return () => { clearTimeout(timer); deadlineActive = false; }; + }, + })).rejects.toMatchObject({ name: "TimeoutError" }); + expect(scheduledMs).toBe(20); + expect(deadlineActive).toBe(false); + }); + + test("an external deadline settles a socket that never closes and clears the local timer", async () => { + const socket = Object.assign(new EventEmitter(), { + setTimeout() { return this; }, + destroy() { return this; }, + }) as unknown as Socket; + const controller = new AbortController(); + let localDeadlineActive = false; + let localDeadlineFired = false; + const pending = directLocalHttpFetch("http://127.0.0.1:9/healthz", { + signal: controller.signal, + }, { + timeoutMs: 1_000, + connect: () => socket, + scheduleDeadline: (onTimeout, delayMs) => { + localDeadlineActive = true; + const timer = setTimeout(() => { localDeadlineFired = true; onTimeout(); }, delayMs); + return () => { clearTimeout(timer); localDeadlineActive = false; }; + }, + }); + const exchangeTimer = setTimeout(() => controller.abort(new DOMException("exchange deadline", "TimeoutError")), 20); + try { + await expect(pending).rejects.toMatchObject({ name: "TimeoutError" }); + expect(localDeadlineFired).toBe(false); + expect(localDeadlineActive).toBe(false); + } finally { + clearTimeout(exchangeTimer); + } + }); + + test("an exchange deadline still cancels teardown after response framing completes", async () => { + const socket = Object.assign(new EventEmitter(), { + setTimeout() { return this; }, + destroy() { return this; }, + }) as unknown as Socket; + const controller = new AbortController(); + let localDeadlineActive = false; + let localDeadlineFired = false; + const pending = directLocalHttpFetch("http://127.0.0.1:9/healthz", { + signal: controller.signal, + }, { + timeoutMs: 1_000, + connect: () => socket, + scheduleDeadline: (onTimeout, delayMs) => { + localDeadlineActive = true; + const timer = setTimeout(() => { localDeadlineFired = true; onTimeout(); }, delayMs); + return () => { clearTimeout(timer); localDeadlineActive = false; }; + }, + }); + socket.emit("data", Buffer.from("HTTP/1.1 200 OK\r\nContent-Length: 2\r\n\r\nok")); + controller.abort(new DOMException("exchange deadline", "TimeoutError")); + await expect(pending).rejects.toMatchObject({ name: "TimeoutError" }); + expect(localDeadlineFired).toBe(false); + expect(localDeadlineActive).toBe(false); + }); + + test("destroy failure settles an aborted request without waiting for close", async () => { + const socket = Object.assign(new EventEmitter(), { + setTimeout() { return this; }, + destroy() { throw new Error("destroy failed"); }, + }) as unknown as Socket; + const controller = new AbortController(); + let deadlineActive = false; + let deadlineFired = false; + const pending = directLocalHttpFetch("http://127.0.0.1:9/healthz", { + signal: controller.signal, + }, { + timeoutMs: 200, + connect: () => socket, + scheduleDeadline: (onTimeout, delayMs) => { + deadlineActive = true; + const timer = setTimeout(() => { deadlineFired = true; onTimeout(); }, delayMs); + return () => { clearTimeout(timer); deadlineActive = false; }; + }, + }); + controller.abort(); + await expect(pending).rejects.toMatchObject({ name: "AbortError" }); + expect(deadlineFired).toBe(false); + expect(deadlineActive).toBe(false); + }); + test("times out an accepted silent socket without an AbortSignal", async () => { let accept!: () => void; const accepted = new Promise(resolve => { accept = resolve; });