From c9e820aed5fc3af69c4487fc7b75e80aa7aff44b Mon Sep 17 00:00:00 2001 From: DongHyeonka Date: Thu, 13 Aug 2026 23:46:03 +0900 Subject: [PATCH] fix: harden the legacy HTTP rollback path N-06: export one idempotency-key authority from mutation-intent.ts and use it in the V2 client. A caller-supplied key is validated before credentials, timers and fetch, and an invalid value is rejected as VALIDATION_REJECTED / IDEMPOTENCY_KEY_INVALID rather than trimmed, regenerated or dropped, so a keyed command can no longer replay while sending no key. N-07: bound the legacy credential wait by the existing attempt controller, which already carries the total deadline and the caller signal, so a non-cooperative owner cannot hold the request open and no extra timer is introduced. The owner receives the operation context, and the failure follows ownership: deadline to REQUEST_TIMEOUT, caller to REQUEST_ABORTED, and only a genuine rejection to AUTH_INTEGRATION_FAILURE. None of these paths fetch. N-08: readBoundedJson delegates to the common bounded reader, so cancel and releaseLock throws stay isolated inside the closed result, and the V2 content-type mismatch now cancels the response body. Co-Authored-By: Claude Opus 5 (1M context) --- docs/operations/adapter-remediation-ledger.md | 6 +- src/adapters/http/bounded-json.ts | 69 ++++----- src/adapters/http/client.ts | 135 +++++++++++++++--- src/contracts/mutation-intent.ts | 33 +++++ tests/integration/http-client.test.ts | 82 +++++++++++ .../http-execution-contract.test.ts | 3 + tests/unit/bounded-json-compatibility.test.ts | 91 ++++++++++++ 7 files changed, 356 insertions(+), 63 deletions(-) create mode 100644 tests/unit/bounded-json-compatibility.test.ts diff --git a/docs/operations/adapter-remediation-ledger.md b/docs/operations/adapter-remediation-ledger.md index 36bff33..5814003 100644 --- a/docs/operations/adapter-remediation-ledger.md +++ b/docs/operations/adapter-remediation-ledger.md @@ -77,9 +77,9 @@ Rollout state starts at `NOT_STARTED`; documented-unimplemented items start at | N-03 | Live V3 path | `corepack pnpm exec vitest run tests/unit/http-execution-v3.test.ts -t "retry-time fence"` | `fix: preserve command effect certainty across retries` | `FIXED_NOT_RELEASED` | command effect verdict regression | Red reproduced `SCOPE_FENCED` with `NOT_STARTED` after one dispatched attempt → green `MAYBE_APPLIED`; lattice table 9/9; `check:types` PASS; `lint` PASS | | N-04 | Live composition teardown | `corepack pnpm exec vitest run tests/unit/telemetry.test.ts tests/unit/runtime-adapters.test.ts` | `fix: terminate telemetry work on disposal` | `FIXED_NOT_RELEASED` | telemetry delivery loss after teardown change | Red 10 failed (5 lifecycle + 5 capacity) → green 34/34; `check:diagnostics` PASS; `check:types` PASS; `check:architecture` PASS; `lint` PASS | | N-05 | Rollout blocker (sidecar not composed) | `corepack pnpm exec vitest run tests/unit/conditional-validator-store.test.ts` | `fix: harden bounded state sidecars` | `FIXED_NOT_RELEASED` | persisted validator key incompatibility | Red collision case (two valid bindings sharing one delimiter-joined key) → green; key is now a bounded validated tuple encoded with `JSON.stringify` | -| N-06 | Legacy V2 rollback seam | `corepack pnpm exec vitest run tests/integration/http-client.test.ts` | — | `NOT_STARTED` | legacy keyed command rejection spike | — | -| N-07 | Legacy V2 rollback seam | `corepack pnpm exec vitest run tests/integration/auth-recovery.test.ts` | — | `NOT_STARTED` | legacy credential timeout regression | — | -| N-08 | Legacy V2 rollback seam | `corepack pnpm exec vitest run tests/unit/bounded-json-compatibility.test.ts` | — | `NOT_STARTED` | legacy JSON failure-code drift | — | +| N-06 | Legacy V2 rollback seam | `corepack pnpm exec vitest run tests/integration/http-client.test.ts` | `fix: harden the legacy HTTP rollback path` | `FIXED_NOT_RELEASED` | legacy keyed command rejection spike | Red 4 invalid-key cases → green; rejection happens before credentials and fetch (0 credential calls, 0 fetches) | +| N-07 | Legacy V2 rollback seam | `corepack pnpm exec vitest run tests/integration/http-client.test.ts tests/integration/auth-recovery.test.ts` | `fix: harden the legacy HTTP rollback path` | `FIXED_NOT_RELEASED` | legacy credential timeout regression | Red never-settling owner → green; credential wait races the existing attempt controller so no extra timer is added; ownership maps to REQUEST_TIMEOUT / REQUEST_ABORTED / AUTH_INTEGRATION_FAILURE with zero fetches | +| N-08 | Legacy V2 rollback seam | `corepack pnpm exec vitest run tests/unit/bounded-json-compatibility.test.ts` | `fix: harden the legacy HTTP rollback path` | `FIXED_NOT_RELEASED` | legacy JSON failure-code drift | Green 7/7 including throwing cancel/releaseLock; `readBoundedJson` now delegates to `bounded-body-reader` with the legacy codes preserved | | N-09 | Live cross-context host | `corepack pnpm exec vitest run tests/unit/cross-tab-invalidation.test.ts` | `fix: harden bounded state sidecars` | `FIXED_NOT_RELEASED` | cross-tab invalidation drop | Red foreign-area pulse accepted → green 13/13; localStorage captured once and `StorageEvent.storageArea` compared by object identity; pulse key registered as `CACHE_INVALIDATION_PULSE`; `check:registries` PASS | | N-10 | Cursor runtime `AVAILABLE_NOT_COMPOSED` | `corepack pnpm exec vitest run tests/unit/cursor-pagination-runtime.test.ts` | `fix: harden bounded state sidecars` | `FIXED_NOT_RELEASED` | pagination abort semantics change | Red never-settling loader → green `PAGINATION_ABORTED` with the late page ignored | | N-11 | Live composition | `corepack pnpm exec vitest run tests/unit/telemetry.test.ts -t capacity` | `fix: terminate telemetry work on disposal` | `FIXED_NOT_RELEASED` | capacity rejection on valid composition | Red 5/5 capacity cases → green; ceilings documented in VD-07 §7-4 | diff --git a/src/adapters/http/bounded-json.ts b/src/adapters/http/bounded-json.ts index 2d69da9..0979c5c 100644 --- a/src/adapters/http/bounded-json.ts +++ b/src/adapters/http/bounded-json.ts @@ -1,50 +1,39 @@ +import { decodeJsonBytes, readBoundedBytes } from "./bounded-body-reader.ts"; + export type BoundedJsonResult = | Readonly<{ ok: true; value: unknown }> | Readonly<{ ok: false; code: "RESPONSE_BODY_LIMIT" | "MALFORMED_JSON" }>; +/** + * N-08. The V2 compatibility reader delegates to the common bounded reader + * instead of maintaining a second stream-reading strategy. + * + * `bounded-body-reader` already isolates `cancel()` and `releaseLock()` throws + * so a cleanup failure cannot escape the closed result. Only the legacy failure + * codes are preserved here: + * + * - `RESPONSE_TOO_LARGE` → `RESPONSE_BODY_LIMIT` + * - `RESPONSE_STREAM_FAILURE` / `UTF8_INVALID` / `JSON_INVALID` → + * `MALFORMED_JSON` + */ export async function readBoundedJson( response: Response, maxBytes: number, ): Promise { - const declaredLength = Number(response.headers.get("content-length")); - if (Number.isFinite(declaredLength) && declaredLength > maxBytes) { - await response.body?.cancel(); - return { ok: false, code: "RESPONSE_BODY_LIMIT" }; - } - if (!response.body) return { ok: false, code: "MALFORMED_JSON" }; - - const reader = response.body.getReader(); - const chunks: Uint8Array[] = []; - let total = 0; - try { - while (true) { - const next = await reader.read(); - if (next.done) break; - total += next.value.byteLength; - if (total > maxBytes) { - await reader.cancel(); - return { ok: false, code: "RESPONSE_BODY_LIMIT" }; - } - chunks.push(next.value); - } - } catch { - return { ok: false, code: "MALFORMED_JSON" }; - } finally { - reader.releaseLock(); - } - - const bytes = new Uint8Array(total); - let offset = 0; - for (const chunk of chunks) { - bytes.set(chunk, offset); - offset += chunk.byteLength; - } - try { - return { - ok: true, - value: JSON.parse(new TextDecoder("utf-8", { fatal: true }).decode(bytes)), - }; - } catch { - return { ok: false, code: "MALFORMED_JSON" }; + const bytes = await readBoundedBytes(response, maxBytes); + if (!bytes.ok) { + return Object.freeze({ + ok: false as const, + code: + bytes.code === "RESPONSE_TOO_LARGE" + ? ("RESPONSE_BODY_LIMIT" as const) + : ("MALFORMED_JSON" as const), + }); } + // An absent body decodes to zero bytes, which is not valid JSON. The legacy + // contract reported that as MALFORMED_JSON, and that is preserved. + const decoded = decodeJsonBytes(bytes.bytes); + return decoded.ok + ? Object.freeze({ ok: true as const, value: decoded.value }) + : Object.freeze({ ok: false as const, code: "MALFORMED_JSON" as const }); } diff --git a/src/adapters/http/client.ts b/src/adapters/http/client.ts index 23f60bf..1f21c45 100644 --- a/src/adapters/http/client.ts +++ b/src/adapters/http/client.ts @@ -19,6 +19,10 @@ import { statusGroup, } from "../../contracts/diagnostics.ts"; import type { AuthSessionPort } from "../../application/ports/auth-session-port.ts"; +import { isValidIdempotencyKey } from "../../contracts/mutation-intent.ts"; + +/** Sentinel for a credential wait ended by the attempt lifetime. */ +const ATTEMPT_ABORTED = Symbol("ATTEMPT_ABORTED"); import type { ClockPort } from "../../application/ports/clock-port.ts"; import type { DiagnosticsPort } from "../../application/ports/diagnostics-port.ts"; import type { TelemetryPort } from "../../application/ports/telemetry-port.ts"; @@ -238,22 +242,52 @@ export function createHttpClient( } return outcome; } + // N-06. A caller-supplied key is validated before credentials, timers or + // fetch. An invalid value is rejected outright rather than trimmed or + // replaced, so a keyed command can never replay with no key at all. let logicalIdempotencyKey: string | undefined; - try { - logicalIdempotencyKey = - operation.idempotency === "keyed" - ? input.idempotencyKey ?? idempotencyKeyFactory() - : undefined; - } catch { - return finalize( - { - ok: false, - error: failure("UNKNOWN_CLIENT_FAILURE", input.operationId, 0, { - code: "IDEMPOTENCY_KEY_CREATION_FAILED", - }), - }, - "failed", - ); + if (operation.idempotency === "keyed") { + if (input.idempotencyKey !== undefined) { + if (!isValidIdempotencyKey(input.idempotencyKey)) { + return finalize( + { + ok: false, + error: failure("VALIDATION_REJECTED", input.operationId, 0, { + code: "IDEMPOTENCY_KEY_INVALID", + }), + }, + "failed", + ); + } + logicalIdempotencyKey = input.idempotencyKey; + } else { + let generated: string; + try { + generated = idempotencyKeyFactory(); + } catch { + return finalize( + { + ok: false, + error: failure("UNKNOWN_CLIENT_FAILURE", input.operationId, 0, { + code: "IDEMPOTENCY_KEY_CREATION_FAILED", + }), + }, + "failed", + ); + } + if (!isValidIdempotencyKey(generated)) { + return finalize( + { + ok: false, + error: failure("UNKNOWN_CLIENT_FAILURE", input.operationId, 0, { + code: "IDEMPOTENCY_KEY_CREATION_FAILED", + }), + }, + "failed", + ); + } + logicalIdempotencyKey = generated; + } } let retryCount = 0; let recoveryUsed = false; @@ -568,11 +602,45 @@ export function createHttpClient( }; } try { - const patch = await authSession.credentialPatch({ - origin: target.url.origin, - method: operation.method, - operationId: operation.operationId, - }); + // N-07. The credential owner is bounded by the attempt lifetime that + // already carries the total deadline and the caller signal, so a + // non-cooperative owner cannot hold the request open and no extra + // timer is introduced. A late completion is observed and discarded. + const raced = await raceAttemptSignal( + Promise.resolve( + authSession.credentialPatch( + { + origin: target.url.origin, + method: operation.method, + operationId: operation.operationId, + }, + Object.freeze({ + signal: controller.signal, + deadlineAtMonotonicMs: deadlineAt, + }), + ), + ), + controller.signal, + ); + if (raced === ATTEMPT_ABORTED) { + return { + ok: false, + error: timedOut + ? failure( + "REQUEST_TIMEOUT", + operation.operationId, + attempt, + { code: "OPERATION_DEADLINE_EXCEEDED" }, + ) + : failure( + "REQUEST_ABORTED", + operation.operationId, + attempt, + { code: "REQUEST_ABORTED" }, + ), + }; + } + const patch = raced; for (const [name, value] of Object.entries(patch.headers)) { const normalized = name.toLowerCase(); const allowedHeaders = @@ -584,6 +652,7 @@ export function createHttpClient( headers.set(normalized, value); } } catch { + // An ordinary owner rejection stays an integration failure. return { ok: false, error: failure("AUTH_INTEGRATION_FAILURE", operation.operationId, attempt, { @@ -663,6 +732,28 @@ export function createHttpClient( return Object.freeze({ execute }); + function raceAttemptSignal( + operation: Promise, + signal: AbortSignal, + ): Promise { + operation.catch(() => {}); + if (signal.aborted) return Promise.resolve(ATTEMPT_ABORTED); + return new Promise((resolve, reject) => { + const onAbort = () => resolve(ATTEMPT_ABORTED); + signal.addEventListener("abort", onAbort, { once: true }); + operation.then( + (value) => { + signal.removeEventListener("abort", onAbort); + resolve(value); + }, + (error: unknown) => { + signal.removeEventListener("abort", onAbort); + reject(error instanceof Error ? error : new Error("rejected")); + }, + ); + }); + } + function withinLogicalDeadline( promise: Promise, deadlineAt: number, @@ -714,6 +805,10 @@ async function parseResponse( const contentType = mediaType(response.headers.get("content-type")); const acceptedMedia = operation.responseMediaTypes ?? ["application/json"]; if (!contentType || !acceptedMedia.includes(contentType)) { + // N-08. A rejected response still owns an open body stream. Cancellation is + // best-effort cleanup, so it is started but not awaited: the closed result + // must not depend on stream teardown settling. + void response.body?.cancel().catch(() => {}); return { ok: false, error: failure("CONTENT_TYPE_MISMATCH", operation.operationId, attempt, { diff --git a/src/contracts/mutation-intent.ts b/src/contracts/mutation-intent.ts index 526033d..a03e7ee 100644 --- a/src/contracts/mutation-intent.ts +++ b/src/contracts/mutation-intent.ts @@ -23,6 +23,39 @@ function validBoundedString(value: unknown, maxBytes: number): value is string { ); } +/** + * N-06. The single idempotency-key authority shared by the V2 compatibility + * client and the V3 executor. + * + * A caller-supplied value is never trimmed, regenerated or silently dropped: + * an invalid key is a contract violation, because replaying a keyed command + * without its key is exactly the unsafe behaviour the key exists to prevent. + */ +export function isValidIdempotencyKey(value: unknown): value is string { + if ( + !validBoundedString( + value, + MUTATION_INTENT_BOUNDS.idempotencyKeyMaxBytes, + ) + ) { + return false; + } + for (const character of value) { + const codePoint = character.codePointAt(0) ?? 0; + if (codePoint <= 0x1f || (codePoint >= 0x7f && codePoint <= 0x9f)) { + return false; + } + } + return true; +} + +export function defineIdempotencyKey(value: unknown): string { + if (!isValidIdempotencyKey(value)) { + throw new TypeError("Idempotency key is invalid."); + } + return value; +} + export function defineMutationIntent(intent: MutationIntent): MutationIntent { if ( !validBoundedString( diff --git a/tests/integration/http-client.test.ts b/tests/integration/http-client.test.ts index b528e86..348aeb3 100644 --- a/tests/integration/http-client.test.ts +++ b/tests/integration/http-client.test.ts @@ -59,6 +59,88 @@ describe("shared HTTP client", () => { expect(attempts).toBe(3); }); + it.each(["", " ", "bad\u0000key", "x".repeat(513)])( + "rejects invalid keyed command key %j before credentials and fetch", + async (idempotencyKey) => { + let fetched = 0; + let credentialAttempts = 0; + server.use( + http.post("https://api.test/api/entities", () => { + fetched += 1; + return HttpResponse.json({ success: true, data: {} }); + }), + ); + const client = testClient({ + baseUrl: "https://api.test", + clock, + authSession: { + getState: () => "authenticated", + subscribe: () => () => {}, + beginSignIn: async () => {}, + signOut: async () => {}, + async credentialPatch() { + credentialAttempts += 1; + return { headers: {} }; + }, + recover: async () => "no-session" as const, + onUnauthenticated: () => {}, + }, + }); + + await expect( + client.execute("CREATE_ENTITY", { + body: { name: "n" }, + idempotencyKey, + }), + ).resolves.toMatchObject({ + ok: false, + error: { + kind: "VALIDATION_REJECTED", + code: "IDEMPOTENCY_KEY_INVALID", + // The repository reports attempt counts as 1-based; the invariant + // proved below is that no physical attempt happened at all. + attemptCount: 1, + }, + }); + expect(fetched).toBe(0); + expect(credentialAttempts).toBe(0); + }, + ); + + it("bounds a non-cooperative legacy credential owner by total deadline", async () => { + let fetched = 0; + let observedSignal: AbortSignal | undefined; + server.use( + http.get("https://api.test/api/entities", () => { + fetched += 1; + return HttpResponse.json({ success: true, data: [] }); + }), + ); + const client = testClient({ + baseUrl: "https://api.test", + timeoutMs: 5, + clock: { now: () => 0, sleep: async () => {} }, + authSession: { + getState: () => "authenticated", + subscribe: () => () => {}, + beginSignIn: async () => {}, + signOut: async () => {}, + credentialPatch: (_binding, context) => { + observedSignal = context?.signal; + // Never settles on its own. + return new Promise(() => {}); + }, + recover: async () => "no-session" as const, + onUnauthenticated: () => {}, + }, + }); + + const result = await client.execute("LIST_ENTITIES"); + expect(result).toMatchObject({ ok: false }); + expect(fetched).toBe(0); + expect(observedSignal).toBeDefined(); + }); + it("rejects a non-JSON response without exposing its body", async () => { server.use( http.get( diff --git a/tests/integration/http-execution-contract.test.ts b/tests/integration/http-execution-contract.test.ts index d975beb..ac15afa 100644 --- a/tests/integration/http-execution-contract.test.ts +++ b/tests/integration/http-execution-contract.test.ts @@ -307,6 +307,9 @@ describe("HTTP operation execution contract", () => { operationId: "LIST_ENTITIES", routeId: "TEST_ROUTE", }); + // The credential wait is bounded by the same attempt controller, so wait + // until the request is actually in flight before firing the deadline. + await vi.waitFor(() => expect(fetcher).toHaveBeenCalledTimes(1)); await vi.waitFor(() => expect(scheduler.callbacks).toHaveLength(1)); scheduler.callbacks[0](); await expect(timeoutResult).resolves.toMatchObject({ diff --git a/tests/unit/bounded-json-compatibility.test.ts b/tests/unit/bounded-json-compatibility.test.ts new file mode 100644 index 0000000..7c4b1b8 --- /dev/null +++ b/tests/unit/bounded-json-compatibility.test.ts @@ -0,0 +1,91 @@ +import { describe, expect, it, vi } from "vitest"; + +import { readBoundedJson } from "../../src/adapters/http/bounded-json.ts"; + +/** + * N-08. The legacy V2 reader delegates to the common bounded reader. These + * cases pin the preserved failure codes and prove that a throwing cancel or + * releaseLock cannot escape the closed result. + */ +function responseWith( + body: ReadableStream | string | null, + init: ResponseInit = {}, +): Response { + return new Response(body, init); +} + +describe("legacy bounded JSON compatibility", () => { + it("preserves RESPONSE_BODY_LIMIT for a declared oversize body", async () => { + const response = responseWith("{}", { + headers: { "content-length": "9999" }, + }); + await expect(readBoundedJson(response, 8)).resolves.toEqual({ + ok: false, + code: "RESPONSE_BODY_LIMIT", + }); + }); + + it("preserves RESPONSE_BODY_LIMIT for a headerless oversize body", async () => { + const response = responseWith("[1,2,3,4,5,6,7,8,9,10]"); + await expect(readBoundedJson(response, 4)).resolves.toEqual({ + ok: false, + code: "RESPONSE_BODY_LIMIT", + }); + }); + + it.each([ + ["invalid UTF-8", new Uint8Array([0xff, 0xfe, 0xfd])], + ["malformed JSON", new TextEncoder().encode("{")], + ["an empty body", new Uint8Array(0)], + ])("maps %s to MALFORMED_JSON", async (_label, bytes) => { + const response = responseWith( + new ReadableStream({ + start(controller) { + if (bytes.byteLength > 0) controller.enqueue(bytes); + controller.close(); + }, + }), + ); + await expect(readBoundedJson(response, 1_024)).resolves.toEqual({ + ok: false, + code: "MALFORMED_JSON", + }); + }); + + it("keeps a closed result when reader cancel or releaseLock throws", async () => { + const stream = new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode("[1,2,3,4,5,6,7,8]")); + }, + }); + const response = responseWith(stream); + const body = response.body!; + const nativeGetReader = body.getReader.bind(body); + vi.spyOn(body, "getReader").mockImplementation((() => { + const actual = nativeGetReader(); + return { + ...actual, + read: actual.read.bind(actual), + cancel: async () => { + throw new TypeError("cancel exploded"); + }, + releaseLock: () => { + throw new TypeError("releaseLock exploded"); + }, + }; + }) as never); + + await expect(readBoundedJson(response, 4)).resolves.toEqual({ + ok: false, + code: "RESPONSE_BODY_LIMIT", + }); + }); + + it("reads a bounded JSON value successfully", async () => { + const response = responseWith('{"a":1}'); + await expect(readBoundedJson(response, 1_024)).resolves.toEqual({ + ok: true, + value: { a: 1 }, + }); + }); +});