From d3a4ca5b35788b0b71068a5c9c74ce2abf6411b2 Mon Sep 17 00:00:00 2001 From: Sebastian Otaegui Date: Mon, 13 Jul 2026 13:31:13 -0300 Subject: [PATCH 1/5] test(agent-journal): reconcile program ledger through PR 128 Record the provider-context merge and successful main CI without advancing the infrastructure acceptance gate. RED: npx vitest run packages/pi-agent-journal/__tests__/evaluation-program-ledger.test.ts Failure: ledger omitted merged PR 128. --- .../agent-work-journal-evaluation-program-state.json | 7 +++++++ .../__tests__/evaluation-program-ledger.test.ts | 7 +++++++ 2 files changed, 14 insertions(+) diff --git a/docs/evaluations/agent-work-journal-evaluation-program-state.json b/docs/evaluations/agent-work-journal-evaluation-program-state.json index 16b11f5..81667ec 100644 --- a/docs/evaluations/agent-work-journal-evaluation-program-state.json +++ b/docs/evaluations/agent-work-journal-evaluation-program-state.json @@ -44,6 +44,13 @@ "mergeCommit": "c2d0104581538d2b5d343e1a490e43b54d4fecf2", "ciRunId": 29264072142, "ciConclusion": "success" + }, + { + "prNumber": 128, + "branch": "feat/agent-journal-provider-context", + "mergeCommit": "568d0ef8b19650ce9e0a9f703cc0ffe2175db7cc", + "ciRunId": 29266604083, + "ciConclusion": "success" } ], "versions": [ diff --git a/packages/pi-agent-journal/__tests__/evaluation-program-ledger.test.ts b/packages/pi-agent-journal/__tests__/evaluation-program-ledger.test.ts index 0d48cd8..4ba2276 100644 --- a/packages/pi-agent-journal/__tests__/evaluation-program-ledger.test.ts +++ b/packages/pi-agent-journal/__tests__/evaluation-program-ledger.test.ts @@ -70,6 +70,13 @@ describe("Agent Journal evaluation program ledger", () => { ciRunId: 29264072142, ciConclusion: "success", }, + { + prNumber: 128, + branch: "feat/agent-journal-provider-context", + mergeCommit: "568d0ef8b19650ce9e0a9f703cc0ffe2175db7cc", + ciRunId: 29266604083, + ciConclusion: "success", + }, ]); }); From aa29dc132c8d831a63a2ee05aaedfcad2e62b485 Mon Sep 17 00:00:00 2001 From: Sebastian Otaegui Date: Mon, 13 Jul 2026 13:36:11 -0300 Subject: [PATCH 2/5] feat(agent-journal): witness provider fetch bytes Wrap the frozen Codex SSE fetch boundary, verify exactly one approved request against the logical filter receipt, safely bound zstd decoding, and retain only content-free logical and physical SHA-256 evidence. RED: npx vitest run packages/pi-agent-journal/__tests__/evaluation-provider-fetch-witness.test.ts Failure: provider fetch witness module was absent. Security RED: npx vitest run packages/pi-agent-journal/__tests__/evaluation-provider-fetch-witness.test.ts -t Headers Failure: a Headers subclass executed its overridden get method and was forwarded. --- .../evaluation-provider-fetch-witness.test.ts | 280 ++++++++++++++++ .../evaluation-provider-fetch-witness.ts | 302 ++++++++++++++++++ 2 files changed, 582 insertions(+) create mode 100644 packages/pi-agent-journal/__tests__/evaluation-provider-fetch-witness.test.ts create mode 100644 packages/pi-agent-journal/extensions/evaluation-provider-fetch-witness.ts diff --git a/packages/pi-agent-journal/__tests__/evaluation-provider-fetch-witness.test.ts b/packages/pi-agent-journal/__tests__/evaluation-provider-fetch-witness.test.ts new file mode 100644 index 0000000..37f5e30 --- /dev/null +++ b/packages/pi-agent-journal/__tests__/evaluation-provider-fetch-witness.test.ts @@ -0,0 +1,280 @@ +import { createHash } from "node:crypto"; +import { createServer } from "node:http"; +import { zstdCompressSync } from "node:zlib"; +import { stream as streamOpenAICodex } from "@earendil-works/pi-ai/api/openai-codex-responses"; +import { OPENAI_CODEX_MODELS } from "@earendil-works/pi-ai/providers/openai-codex.models"; +import { describe, expect, it, vi } from "vitest"; +import { + createJournalPhaseBProviderBoundary, + PHASE_B_UNTRUSTED_NOTICE, +} from "../extensions/evaluation-provider-context.js"; +import { + createProviderFetchWitness, + EvaluationProviderFetchWitnessError, +} from "../extensions/evaluation-provider-fetch-witness.js"; +import { canonicalSha256 } from "../extensions/evaluation-receipts.js"; + +const sha256 = (value: string | Uint8Array): string => createHash("sha256").update(value).digest("hex"); +const prompt = "Continue the synthetic task."; +const instructions = "Frozen synthetic instructions."; +const capsule = JSON.stringify({ notice: PHASE_B_UNTRUSTED_NOTICE, data: { checkpoint: "current" } }); +const tools = [{ type: "function", name: "read", description: "Read", parameters: { type: "object" }, strict: null }]; +const token = `e30.${Buffer.from( + JSON.stringify({ "https://api.openai.com/auth": { chatgpt_account_id: "synthetic-account" } }), +).toString("base64url")}.signature`; + +function user(content: string) { + return { role: "user", content, timestamp: 10 }; +} + +function resume() { + return { + role: "custom", + customType: "agent-journal-resume", + content: capsule, + display: false, + details: { sessionId: "journal-session", checkpointId: "checkpoint-current", fingerprint: "a".repeat(64) }, + timestamp: 11, + }; +} + +function providerBoundary() { + return createJournalPhaseBProviderBoundary({ + phaseBPrompt: prompt, + expectedSessionId: "journal-session", + expectedInstructionsDigest: sha256(instructions), + expectedToolsDigest: canonicalSha256(tools).digest, + expectedToolNames: ["read"], + expectedPromptCacheKey: "phase-b-cache", + transport: "sse", + onTerminalFailure: () => undefined, + }); +} + +function logicalBody() { + return JSON.stringify({ + model: "gpt-5.6-sol", + store: false, + stream: true, + instructions, + input: [ + { role: "user", content: [{ type: "input_text", text: prompt }] }, + { role: "user", content: [{ type: "input_text", text: capsule }] }, + ], + text: { verbosity: "low" }, + include: ["reasoning.encrypted_content"], + prompt_cache_key: "phase-b-cache", + tool_choice: "auto", + parallel_tool_calls: true, + tools, + reasoning: { effort: "high", summary: "auto" }, + }); +} + +function logicalReceipt(body = logicalBody()) { + return { + schemaVersion: 1 as const, + state: "filtered" as const, + model: "gpt-5.6-sol" as const, + transport: "sse" as const, + inputItems: 2, + previousResponseIdPresent: false as const, + byteLength: Buffer.byteLength(body, "utf8"), + payloadDigest: sha256(body), + instructionsDigest: sha256(instructions), + toolsDigest: canonicalSha256(tools).digest, + phaseBPromptDigest: sha256(prompt), + capsuleDigest: sha256(capsule), + }; +} + +function directSetup(overrides: Record = {}) { + const forwarded = vi.fn(async () => new Response("ok", { status: 200 })); + const failures: string[] = []; + const body = logicalBody(); + const witness = createProviderFetchWitness({ + attemptId: "attempt-001", + expectedUrl: "https://chatgpt.com/backend-api/codex/responses", + endpointClass: "openai-codex", + getLogicalReceipt: () => logicalReceipt(body), + fetchImpl: forwarded, + onTerminalFailure: (code) => failures.push(code), + ...overrides, + }); + return { witness, forwarded, failures, body }; +} + +describe("Agent Journal provider fetch witness", () => { + it("attests the physical zstd body sent by the real Pi 0.80.6 Codex SSE serializer", async () => { + let received = false; + const server = createServer((_request, response) => { + received = true; + response.writeHead(200, { "content-type": "text/event-stream", connection: "close" }); + response.end( + `data: ${JSON.stringify({ + type: "response.completed", + response: { + id: "response-synthetic", + status: "completed", + output: [], + usage: { input_tokens: 0, output_tokens: 0, input_tokens_details: {}, output_tokens_details: {} }, + }, + })}\n\n`, + ); + }); + await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); + const originalFetch = globalThis.fetch; + try { + const address = server.address(); + if (!address || typeof address === "string") throw new Error("loopback did not bind"); + const expectedUrl = `http://127.0.0.1:${address.port}/codex/responses`; + const boundary = providerBoundary(); + boundary.filterContext([user("PHASE_A_SECRET"), user(prompt), resume()]); + const failures: string[] = []; + const witness = createProviderFetchWitness({ + attemptId: "attempt-loopback", + expectedUrl, + endpointClass: "loopback", + getLogicalReceipt: () => boundary.getLastPayloadReceipt(), + fetchImpl: originalFetch, + onTerminalFailure: (code) => failures.push(code), + }); + vi.stubGlobal("fetch", witness.fetch); + const stream = streamOpenAICodex( + { ...OPENAI_CODEX_MODELS["gpt-5.6-sol"], baseUrl: `http://127.0.0.1:${address.port}` }, + { + systemPrompt: instructions, + messages: [ + { role: "user", content: prompt, timestamp: 10 }, + { role: "user", content: capsule, timestamp: 11 }, + ], + tools: [{ name: "read", description: "Read", parameters: { type: "object" } }], + }, + { + apiKey: token, + transport: "sse", + reasoningEffort: "high", + sessionId: "phase-b-cache", + maxRetries: 0, + onPayload: (value) => boundary.filterProviderPayload(value), + }, + ); + for await (const _event of stream) { + // Consume the synthetic terminal response. + } + const filtered = boundary.getLastPayloadReceipt(); + const receipt = witness.getReceipt(); + expect(received).toBe(true); + expect(receipt).toMatchObject({ + schemaVersion: 1, + state: "observed", + attemptId: "attempt-loopback", + endpointClass: "loopback", + model: "gpt-5.6-sol", + transport: "sse", + encoding: "zstd", + requestCount: 1, + logicalDigest: filtered?.payloadDigest, + logicalByteLength: filtered?.byteLength, + inputItems: 2, + previousResponseIdPresent: false, + }); + expect(receipt?.physicalDigest).toMatch(/^[a-f0-9]{64}$/); + expect(receipt?.physicalByteLength).toBeGreaterThan(0); + expect(JSON.stringify(receipt)).not.toContain(prompt); + expect(JSON.stringify(receipt)).not.toContain(capsule); + expect(failures).toEqual([]); + } finally { + vi.unstubAllGlobals(); + server.closeAllConnections(); + await new Promise((resolve) => server.close(() => resolve())); + } + }); + + it.each([ + ["wrong URL", { input: "https://evil.invalid/responses" }], + ["wrong method", { init: { method: "GET" } }], + ["unknown encoding", { headers: { "content-encoding": "gzip" } }], + ["malformed JSON", { body: "{" }], + ["logical digest mismatch", { receiptBody: `${logicalBody()} ` }], + ])("rejects %s without forwarding", async (_label, mutation) => { + const receiptBody = "receiptBody" in mutation ? mutation.receiptBody : logicalBody(); + const { witness, forwarded, failures } = directSetup({ getLogicalReceipt: () => logicalReceipt(receiptBody) }); + const input = "input" in mutation ? mutation.input : "https://chatgpt.com/backend-api/codex/responses"; + const headers = new Headers("headers" in mutation ? mutation.headers : undefined); + const body = "body" in mutation ? mutation.body : logicalBody(); + const init = { method: "POST", headers, body, ...("init" in mutation ? mutation.init : {}) }; + await expect(witness.fetch(input, init)).rejects.toThrow(EvaluationProviderFetchWitnessError); + expect(forwarded).not.toHaveBeenCalled(); + expect(failures).toEqual(["invalid-fetch"]); + expect(witness.getReceipt()).toBeNull(); + }); + + it("rejects a Headers subclass without invoking its overridden get method", async () => { + const trap = vi.fn((_name: string) => null); + class HostileHeaders extends Headers { + override get(name: string): string | null { + trap(name); + return null; + } + } + const { witness, forwarded, body } = directSetup(); + await expect( + witness.fetch("https://chatgpt.com/backend-api/codex/responses", { + method: "POST", + headers: new HostileHeaders(), + body, + }), + ).rejects.toThrow(EvaluationProviderFetchWitnessError); + expect(trap).not.toHaveBeenCalled(); + expect(forwarded).not.toHaveBeenCalled(); + }); + + it("rejects a second provider request without forwarding it", async () => { + const { witness, forwarded, body } = directSetup(); + const init = { method: "POST", headers: new Headers(), body }; + await witness.fetch("https://chatgpt.com/backend-api/codex/responses", init); + await expect(witness.fetch("https://chatgpt.com/backend-api/codex/responses", init)).rejects.toThrow( + EvaluationProviderFetchWitnessError, + ); + expect(forwarded).toHaveBeenCalledTimes(1); + }); + + it("rejects previous_response_id and input-count drift from logical provider JSON", async () => { + for (const mutate of [ + (value: Record) => (value.previous_response_id = "old-response"), + (value: Record) => (value.input = []), + ]) { + const value = JSON.parse(logicalBody()) as Record; + mutate(value); + const body = JSON.stringify(value); + const { witness, forwarded } = directSetup({ getLogicalReceipt: () => logicalReceipt(body) }); + await expect( + witness.fetch("https://chatgpt.com/backend-api/codex/responses", { + method: "POST", + headers: new Headers(), + body, + }), + ).rejects.toThrow(EvaluationProviderFetchWitnessError); + expect(forwarded).not.toHaveBeenCalled(); + } + }); + + it("bounds zstd decompression and normalizes callback failures", async () => { + const large = Buffer.alloc(2 * 1024 * 1024, 65); + const physical = zstdCompressSync(large); + const { witness, forwarded } = directSetup({ + onTerminalFailure: () => { + throw new Error("private callback detail"); + }, + }); + await expect( + witness.fetch("https://chatgpt.com/backend-api/codex/responses", { + method: "POST", + headers: new Headers({ "content-encoding": "zstd" }), + body: physical, + }), + ).rejects.toThrow("provider fetch witness rejected the request"); + expect(forwarded).not.toHaveBeenCalled(); + }); +}); diff --git a/packages/pi-agent-journal/extensions/evaluation-provider-fetch-witness.ts b/packages/pi-agent-journal/extensions/evaluation-provider-fetch-witness.ts new file mode 100644 index 0000000..d07480b --- /dev/null +++ b/packages/pi-agent-journal/extensions/evaluation-provider-fetch-witness.ts @@ -0,0 +1,302 @@ +import { createHash } from "node:crypto"; +import { isProxy } from "node:util/types"; +import { zstdDecompressSync } from "node:zlib"; +import type { ProviderPayloadReceipt } from "./evaluation-provider-context.js"; + +/** + * Evaluation-only physical request witness. This observes one frozen SSE fetch + * after provider filtering; it does not establish extension order, credentials, + * provider authenticity, attempt scheduling, or full B3 acceptance by itself. + */ +const SHA256 = /^[a-f0-9]{64}$/; +const OPAQUE_ID = /^[A-Za-z0-9][A-Za-z0-9._-]{0,127}$/; +const MAX_PHYSICAL_BYTES = 2 * 1024 * 1024; +const MAX_LOGICAL_BYTES = 1024 * 1024; +const PAYLOAD_KEYS = [ + "model", + "store", + "stream", + "instructions", + "input", + "text", + "include", + "prompt_cache_key", + "tool_choice", + "parallel_tool_calls", + "tools", + "reasoning", +] as const; +const LOGICAL_RECEIPT_KEYS = [ + "schemaVersion", + "state", + "model", + "transport", + "inputItems", + "previousResponseIdPresent", + "byteLength", + "payloadDigest", + "instructionsDigest", + "toolsDigest", + "phaseBPromptDigest", + "capsuleDigest", +] as const; + +export type ProviderFetchFailureCode = "invalid-fetch"; +export type ProviderEndpointClass = "openai-codex" | "loopback"; + +export interface ProviderFetchWitnessOptions { + attemptId: string; + expectedUrl: string; + endpointClass: ProviderEndpointClass; + getLogicalReceipt: () => ProviderPayloadReceipt | null; + fetchImpl: typeof fetch; + onTerminalFailure: (code: ProviderFetchFailureCode) => void; +} + +export interface ProviderFetchReceipt { + schemaVersion: 1; + state: "observed"; + attemptId: string; + endpointClass: ProviderEndpointClass; + model: "gpt-5.6-sol"; + transport: "sse"; + encoding: "identity" | "zstd"; + requestCount: 1; + logicalDigest: string; + logicalByteLength: number; + physicalDigest: string; + physicalByteLength: number; + inputItems: number; + previousResponseIdPresent: false; +} + +export interface ProviderFetchWitness { + fetch: typeof fetch; + getReceipt(): ProviderFetchReceipt | null; +} + +export class EvaluationProviderFetchWitnessError extends Error { + constructor() { + super("provider fetch witness rejected the request"); + this.name = "EvaluationProviderFetchWitnessError"; + } +} + +type DataRecord = Record; + +function reject(): never { + throw new EvaluationProviderFetchWitnessError(); +} + +function sha256(value: Uint8Array): string { + return createHash("sha256").update(value).digest("hex"); +} + +function dataRecord(value: unknown): DataRecord { + if ( + !value || + typeof value !== "object" || + isProxy(value) || + Array.isArray(value) || + Object.getPrototypeOf(value) !== Object.prototype + ) { + reject(); + } + const keys = Reflect.ownKeys(value); + if (keys.some((key) => typeof key !== "string")) reject(); + const record: DataRecord = {}; + for (const key of keys as string[]) { + const descriptor = Object.getOwnPropertyDescriptor(value, key); + if (!descriptor || !("value" in descriptor) || descriptor.enumerable !== true) reject(); + Object.defineProperty(record, key, { + value: descriptor.value, + enumerable: true, + writable: true, + configurable: true, + }); + } + return record; +} + +function exactRecord(value: unknown, keys: readonly string[]): DataRecord { + const record = dataRecord(value); + if (JSON.stringify(Object.keys(record).sort()) !== JSON.stringify([...keys].sort())) reject(); + return record; +} + +function dataArray(value: unknown): unknown[] { + if (!value || typeof value !== "object" || isProxy(value) || !Array.isArray(value)) reject(); + if (Object.getPrototypeOf(value) !== Array.prototype) reject(); + const keys = Reflect.ownKeys(value); + const expected = [...Array.from({ length: value.length }, (_, index) => String(index)), "length"].sort(); + if ( + keys.some((key) => typeof key !== "string") || + JSON.stringify((keys as string[]).sort()) !== JSON.stringify(expected) + ) { + reject(); + } + return value; +} + +function boundedString(value: unknown, maxBytes: number): string { + if ( + typeof value !== "string" || + Buffer.byteLength(value, "utf8") < 1 || + Buffer.byteLength(value, "utf8") > maxBytes + ) { + reject(); + } + return value; +} + +function positiveInteger(value: unknown): number { + if (!Number.isSafeInteger(value) || (value as number) < 1) reject(); + return value as number; +} + +function validateEndpoint(expectedUrl: string, endpointClass: ProviderEndpointClass): void { + let parsed: URL; + try { + parsed = new URL(expectedUrl); + } catch { + reject(); + } + if (parsed.username || parsed.password || parsed.search || parsed.hash) reject(); + if (endpointClass === "openai-codex") { + if (parsed.protocol !== "https:" || parsed.hostname !== "chatgpt.com") reject(); + } else if ( + parsed.protocol !== "http:" || + (parsed.hostname !== "127.0.0.1" && parsed.hostname !== "localhost" && parsed.hostname !== "[::1]") + ) { + reject(); + } +} + +function logicalReceipt(value: unknown): ProviderPayloadReceipt { + const receipt = exactRecord(value, LOGICAL_RECEIPT_KEYS); + if ( + receipt.schemaVersion !== 1 || + receipt.state !== "filtered" || + receipt.model !== "gpt-5.6-sol" || + receipt.transport !== "sse" || + receipt.previousResponseIdPresent !== false || + typeof receipt.payloadDigest !== "string" || + !SHA256.test(receipt.payloadDigest) + ) { + reject(); + } + positiveInteger(receipt.inputItems); + positiveInteger(receipt.byteLength); + return receipt as unknown as ProviderPayloadReceipt; +} + +function physicalBody(value: unknown): Buffer { + if (typeof value === "string") return Buffer.from(value, "utf8"); + if (!value || typeof value !== "object" || isProxy(value) || !(value instanceof Uint8Array)) reject(); + const prototype = Object.getPrototypeOf(value); + if (prototype !== Uint8Array.prototype && prototype !== Buffer.prototype) reject(); + return Buffer.from(value); +} + +export function createProviderFetchWitness(options: ProviderFetchWitnessOptions): ProviderFetchWitness { + const config = exactRecord(options, [ + "attemptId", + "expectedUrl", + "endpointClass", + "getLogicalReceipt", + "fetchImpl", + "onTerminalFailure", + ]); + const attemptId = boundedString(config.attemptId, 128); + const expectedUrl = boundedString(config.expectedUrl, 2048); + const endpointClass = config.endpointClass; + if (!OPAQUE_ID.test(attemptId) || (endpointClass !== "openai-codex" && endpointClass !== "loopback")) reject(); + validateEndpoint(expectedUrl, endpointClass); + if ( + typeof config.getLogicalReceipt !== "function" || + typeof config.fetchImpl !== "function" || + typeof config.onTerminalFailure !== "function" + ) { + reject(); + } + const getLogicalReceipt = config.getLogicalReceipt as () => ProviderPayloadReceipt | null; + const fetchImpl = config.fetchImpl as typeof fetch; + const onTerminalFailure = config.onTerminalFailure as (code: ProviderFetchFailureCode) => void; + + let terminal = false; + let requestCount = 0; + let receipt: ProviderFetchReceipt | null = null; + + const failClosed = (): never => { + terminal = true; + receipt = null; + try { + onTerminalFailure("invalid-fetch"); + } catch { + // The generic thrown error remains authoritative and content-free. + } + reject(); + }; + + const witnessedFetch: typeof fetch = async (input, init) => { + if (terminal || requestCount !== 0) failClosed(); + try { + const requestUrl = + typeof input === "string" ? input : input instanceof URL && !isProxy(input) ? input.href : reject(); + if (requestUrl !== expectedUrl || !init || isProxy(init)) reject(); + const request = dataRecord(init); + if (request.method !== "POST") reject(); + if (!request.headers || typeof request.headers !== "object" || isProxy(request.headers)) reject(); + if (!(request.headers instanceof Headers) || Object.getPrototypeOf(request.headers) !== Headers.prototype) + reject(); + const rawEncoding = request.headers.get("content-encoding"); + if (rawEncoding !== null && rawEncoding !== "zstd") reject(); + const encoding = rawEncoding === "zstd" ? "zstd" : "identity"; + const physical = physicalBody(request.body); + if (physical.byteLength < 1 || physical.byteLength > MAX_PHYSICAL_BYTES) reject(); + const logical = + encoding === "zstd" + ? Buffer.from(zstdDecompressSync(physical, { maxOutputLength: MAX_LOGICAL_BYTES })) + : Buffer.from(physical); + if (logical.byteLength < 1 || logical.byteLength > MAX_LOGICAL_BYTES) reject(); + const text = logical.toString("utf8"); + if (!Buffer.from(text, "utf8").equals(logical)) reject(); + + const expected = logicalReceipt(getLogicalReceipt()); + if (expected.byteLength !== logical.byteLength || expected.payloadDigest !== sha256(logical)) reject(); + const payload: unknown = JSON.parse(text); + const root = exactRecord(payload, PAYLOAD_KEYS); + if (root.model !== "gpt-5.6-sol" || root.store !== false || root.stream !== true) reject(); + if (Object.hasOwn(root, "previous_response_id")) reject(); + const items = dataArray(root.input); + if (items.length !== expected.inputItems) reject(); + + requestCount = 1; + receipt = Object.freeze({ + schemaVersion: 1, + state: "observed", + attemptId, + endpointClass, + model: "gpt-5.6-sol", + transport: "sse", + encoding, + requestCount: 1, + logicalDigest: expected.payloadDigest, + logicalByteLength: logical.byteLength, + physicalDigest: sha256(physical), + physicalByteLength: physical.byteLength, + inputItems: items.length, + previousResponseIdPresent: false, + }); + return fetchImpl(input, init); + } catch (error) { + if (requestCount === 1 && receipt && !(error instanceof EvaluationProviderFetchWitnessError)) throw error; + return failClosed(); + } + }; + + return Object.freeze({ + fetch: witnessedFetch, + getReceipt: () => receipt, + }); +} From 52edfd1ba9c1556db3f8e17c621dbc69c240df55 Mon Sep 17 00:00:00 2001 From: Sebastian Otaegui Date: Mon, 13 Jul 2026 13:37:33 -0300 Subject: [PATCH 3/5] fix(agent-journal): isolate witnessed fetch bytes Forward private copies of the validated physical body and native headers so caller-retained references cannot mutate bytes after attestation. RED: npx vitest run packages/pi-agent-journal/__tests__/evaluation-provider-fetch-witness.test.ts -t private Failure: fetch received the original mutable body and Headers references. --- .../evaluation-provider-fetch-witness.test.ts | 20 +++++++++++++++++++ .../evaluation-provider-fetch-witness.ts | 7 ++++++- 2 files changed, 26 insertions(+), 1 deletion(-) diff --git a/packages/pi-agent-journal/__tests__/evaluation-provider-fetch-witness.test.ts b/packages/pi-agent-journal/__tests__/evaluation-provider-fetch-witness.test.ts index 37f5e30..1cccade 100644 --- a/packages/pi-agent-journal/__tests__/evaluation-provider-fetch-witness.test.ts +++ b/packages/pi-agent-journal/__tests__/evaluation-provider-fetch-witness.test.ts @@ -210,6 +210,26 @@ describe("Agent Journal provider fetch witness", () => { expect(witness.getReceipt()).toBeNull(); }); + it("forwards private body and header copies rather than caller-retained references", async () => { + const body = Buffer.from(logicalBody(), "utf8"); + const original = Buffer.from(body); + const headers = new Headers({ "x-safe-metadata": "one" }); + let forwardedInit: RequestInit | undefined; + const { witness } = directSetup({ + fetchImpl: vi.fn(async (_input: URL | RequestInfo, init?: RequestInit) => { + forwardedInit = init; + return new Response("ok"); + }), + }); + await witness.fetch("https://chatgpt.com/backend-api/codex/responses", { method: "POST", headers, body }); + expect(forwardedInit?.body).not.toBe(body); + expect(forwardedInit?.headers).not.toBe(headers); + body.fill(0); + headers.set("x-safe-metadata", "mutated"); + expect(Buffer.from(forwardedInit?.body as Uint8Array)).toEqual(original); + expect(new Headers(forwardedInit?.headers).get("x-safe-metadata")).toBe("one"); + }); + it("rejects a Headers subclass without invoking its overridden get method", async () => { const trap = vi.fn((_name: string) => null); class HostileHeaders extends Headers { diff --git a/packages/pi-agent-journal/extensions/evaluation-provider-fetch-witness.ts b/packages/pi-agent-journal/extensions/evaluation-provider-fetch-witness.ts index d07480b..58dc26a 100644 --- a/packages/pi-agent-journal/extensions/evaluation-provider-fetch-witness.ts +++ b/packages/pi-agent-journal/extensions/evaluation-provider-fetch-witness.ts @@ -288,7 +288,12 @@ export function createProviderFetchWitness(options: ProviderFetchWitnessOptions) inputItems: items.length, previousResponseIdPresent: false, }); - return fetchImpl(input, init); + const forwardedInit = { + ...request, + headers: new Headers(request.headers), + body: Buffer.from(physical), + } as RequestInit; + return fetchImpl(requestUrl, forwardedInit); } catch (error) { if (requestCount === 1 && receipt && !(error instanceof EvaluationProviderFetchWitnessError)) throw error; return failClosed(); From 562209fbdaba6be5beb16e7630c89de577db0e33 Mon Sep 17 00:00:00 2001 From: Sebastian Otaegui Date: Mon, 13 Jul 2026 13:43:40 -0300 Subject: [PATCH 4/5] fix(agent-journal): fail closed across fetch lifecycle Bind exact endpoint paths, reject overridable URL and Headers objects, block callback reentry before validation, clear receipts on transport failure, and assert server-observed physical bytes against the receipt. Also state that real Pi runner proof remains deferred. RED: npx vitest run packages/pi-agent-journal/__tests__/evaluation-provider-fetch-witness.test.ts -t reentry Failures: forged Headers and URL values forwarded, callback reentry sent twice, transport failure retained a receipt, and non-Codex paths were accepted. --- .../evaluation-provider-fetch-witness.test.ts | 153 +++++++++++++++--- .../evaluation-provider-fetch-witness.ts | 50 ++++-- 2 files changed, 166 insertions(+), 37 deletions(-) diff --git a/packages/pi-agent-journal/__tests__/evaluation-provider-fetch-witness.test.ts b/packages/pi-agent-journal/__tests__/evaluation-provider-fetch-witness.test.ts index 1cccade..9842554 100644 --- a/packages/pi-agent-journal/__tests__/evaluation-provider-fetch-witness.test.ts +++ b/packages/pi-agent-journal/__tests__/evaluation-provider-fetch-witness.test.ts @@ -1,6 +1,6 @@ import { createHash } from "node:crypto"; import { createServer } from "node:http"; -import { zstdCompressSync } from "node:zlib"; +import { zstdCompressSync, zstdDecompressSync } from "node:zlib"; import { stream as streamOpenAICodex } from "@earendil-works/pi-ai/api/openai-codex-responses"; import { OPENAI_CODEX_MODELS } from "@earendil-works/pi-ai/providers/openai-codex.models"; import { describe, expect, it, vi } from "vitest"; @@ -105,22 +105,28 @@ function directSetup(overrides: Record = {}) { } describe("Agent Journal provider fetch witness", () => { - it("attests the physical zstd body sent by the real Pi 0.80.6 Codex SSE serializer", async () => { - let received = false; - const server = createServer((_request, response) => { - received = true; - response.writeHead(200, { "content-type": "text/event-stream", connection: "close" }); - response.end( - `data: ${JSON.stringify({ - type: "response.completed", - response: { - id: "response-synthetic", - status: "completed", - output: [], - usage: { input_tokens: 0, output_tokens: 0, input_tokens_details: {}, output_tokens_details: {} }, - }, - })}\n\n`, - ); + it("attests the physical zstd body sent by the Pi 0.80.6 Codex SSE serializer", async () => { + let receivedPhysical = Buffer.alloc(0); + let receivedEncoding: string | undefined; + const server = createServer((request, response) => { + const chunks: Buffer[] = []; + request.on("data", (chunk: Buffer) => chunks.push(chunk)); + request.on("end", () => { + receivedPhysical = Buffer.concat(chunks); + receivedEncoding = request.headers["content-encoding"]; + response.writeHead(200, { "content-type": "text/event-stream", connection: "close" }); + response.end( + `data: ${JSON.stringify({ + type: "response.completed", + response: { + id: "response-synthetic", + status: "completed", + output: [], + usage: { input_tokens: 0, output_tokens: 0, input_tokens_details: {}, output_tokens_details: {} }, + }, + })}\n\n`, + ); + }); }); await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); const originalFetch = globalThis.fetch; @@ -164,7 +170,8 @@ describe("Agent Journal provider fetch witness", () => { } const filtered = boundary.getLastPayloadReceipt(); const receipt = witness.getReceipt(); - expect(received).toBe(true); + expect(receivedPhysical.byteLength).toBeGreaterThan(0); + expect(receivedEncoding).toBe("zstd"); expect(receipt).toMatchObject({ schemaVersion: 1, state: "observed", @@ -179,8 +186,11 @@ describe("Agent Journal provider fetch witness", () => { inputItems: 2, previousResponseIdPresent: false, }); - expect(receipt?.physicalDigest).toMatch(/^[a-f0-9]{64}$/); - expect(receipt?.physicalByteLength).toBeGreaterThan(0); + expect(receipt?.physicalDigest).toBe(sha256(receivedPhysical)); + expect(receipt?.physicalByteLength).toBe(receivedPhysical.byteLength); + const receivedLogical = zstdDecompressSync(receivedPhysical); + expect(sha256(receivedLogical)).toBe(receipt?.logicalDigest); + expect(receivedLogical.byteLength).toBe(receipt?.logicalByteLength); expect(JSON.stringify(receipt)).not.toContain(prompt); expect(JSON.stringify(receipt)).not.toContain(capsule); expect(failures).toEqual([]); @@ -230,6 +240,37 @@ describe("Agent Journal provider fetch witness", () => { expect(new Headers(forwardedInit?.headers).get("x-safe-metadata")).toBe("one"); }); + it("rejects an own Headers.get override without invoking it", async () => { + const trap = vi.fn((_name: string) => null); + const headers = new Headers(); + Object.defineProperty(headers, "get", { value: trap, enumerable: false }); + const { witness, forwarded, body } = directSetup(); + await expect( + witness.fetch("https://chatgpt.com/backend-api/codex/responses", { method: "POST", headers, body }), + ).rejects.toThrow(EvaluationProviderFetchWitnessError); + expect(trap).not.toHaveBeenCalled(); + expect(forwarded).not.toHaveBeenCalled(); + }); + + it("rejects a URL subclass without invoking an overridden href getter", async () => { + const trap = vi.fn(() => "https://chatgpt.com/backend-api/codex/responses"); + class HostileUrl extends URL { + override get href(): string { + return trap(); + } + } + const { witness, forwarded, body } = directSetup(); + await expect( + witness.fetch(new HostileUrl("https://evil.invalid/responses"), { + method: "POST", + headers: new Headers(), + body, + }), + ).rejects.toThrow(EvaluationProviderFetchWitnessError); + expect(trap).not.toHaveBeenCalled(); + expect(forwarded).not.toHaveBeenCalled(); + }); + it("rejects a Headers subclass without invoking its overridden get method", async () => { const trap = vi.fn((_name: string) => null); class HostileHeaders extends Headers { @@ -250,6 +291,65 @@ describe("Agent Journal provider fetch witness", () => { expect(forwarded).not.toHaveBeenCalled(); }); + it("fails the whole witness on callback reentry before either request forwards", async () => { + const body = logicalBody(); + const forwarded = vi.fn(async () => new Response("ok")); + let witness: ReturnType; + let inner: Promise | undefined; + let recursing = false; + witness = createProviderFetchWitness({ + attemptId: "attempt-reentrant", + expectedUrl: "https://chatgpt.com/backend-api/codex/responses", + endpointClass: "openai-codex", + getLogicalReceipt: () => { + if (!recursing) { + recursing = true; + inner = witness.fetch("https://chatgpt.com/backend-api/codex/responses", { + method: "POST", + headers: new Headers(), + body, + }); + } + return logicalReceipt(body); + }, + fetchImpl: forwarded, + onTerminalFailure: () => undefined, + }); + const outer = witness.fetch("https://chatgpt.com/backend-api/codex/responses", { + method: "POST", + headers: new Headers(), + body, + }); + await expect(outer).rejects.toThrow(EvaluationProviderFetchWitnessError); + await expect(inner).rejects.toThrow(EvaluationProviderFetchWitnessError); + expect(forwarded).not.toHaveBeenCalled(); + expect(witness.getReceipt()).toBeNull(); + }); + + it("clears evidence and fails terminally when the underlying fetch rejects", async () => { + const failures: string[] = []; + const body = logicalBody(); + const witness = createProviderFetchWitness({ + attemptId: "attempt-fetch-failure", + expectedUrl: "https://chatgpt.com/backend-api/codex/responses", + endpointClass: "openai-codex", + getLogicalReceipt: () => logicalReceipt(body), + fetchImpl: vi.fn(async () => { + throw new Error("private transport detail"); + }), + onTerminalFailure: (code) => failures.push(code), + }); + await expect( + witness.fetch("https://chatgpt.com/backend-api/codex/responses", { + method: "POST", + headers: new Headers(), + body, + }), + ).rejects.toThrow("provider fetch witness rejected the request"); + expect(witness.getReceipt()).toBeNull(); + expect(failures).toEqual(["fetch-failed"]); + }); + it("rejects a second provider request without forwarding it", async () => { const { witness, forwarded, body } = directSetup(); const init = { method: "POST", headers: new Headers(), body }; @@ -280,6 +380,19 @@ describe("Agent Journal provider fetch witness", () => { } }); + it("requires the exact frozen OpenAI Codex endpoint path", () => { + expect(() => + createProviderFetchWitness({ + attemptId: "attempt-wrong-path", + expectedUrl: "https://chatgpt.com/other", + endpointClass: "openai-codex", + getLogicalReceipt: () => logicalReceipt(), + fetchImpl: globalThis.fetch, + onTerminalFailure: () => undefined, + }), + ).toThrow(EvaluationProviderFetchWitnessError); + }); + it("bounds zstd decompression and normalizes callback failures", async () => { const large = Buffer.alloc(2 * 1024 * 1024, 65); const physical = zstdCompressSync(large); diff --git a/packages/pi-agent-journal/extensions/evaluation-provider-fetch-witness.ts b/packages/pi-agent-journal/extensions/evaluation-provider-fetch-witness.ts index 58dc26a..0912f20 100644 --- a/packages/pi-agent-journal/extensions/evaluation-provider-fetch-witness.ts +++ b/packages/pi-agent-journal/extensions/evaluation-provider-fetch-witness.ts @@ -6,7 +6,8 @@ import type { ProviderPayloadReceipt } from "./evaluation-provider-context.js"; /** * Evaluation-only physical request witness. This observes one frozen SSE fetch * after provider filtering; it does not establish extension order, credentials, - * provider authenticity, attempt scheduling, or full B3 acceptance by itself. + * provider authenticity, attempt scheduling, a real Pi runner lifecycle, or full + * B3 acceptance by itself. */ const SHA256 = /^[a-f0-9]{64}$/; const OPAQUE_ID = /^[A-Za-z0-9][A-Za-z0-9._-]{0,127}$/; @@ -41,7 +42,7 @@ const LOGICAL_RECEIPT_KEYS = [ "capsuleDigest", ] as const; -export type ProviderFetchFailureCode = "invalid-fetch"; +export type ProviderFetchFailureCode = "invalid-fetch" | "fetch-failed"; export type ProviderEndpointClass = "openai-codex" | "loopback"; export interface ProviderFetchWitnessOptions { @@ -163,9 +164,10 @@ function validateEndpoint(expectedUrl: string, endpointClass: ProviderEndpointCl } if (parsed.username || parsed.password || parsed.search || parsed.hash) reject(); if (endpointClass === "openai-codex") { - if (parsed.protocol !== "https:" || parsed.hostname !== "chatgpt.com") reject(); + if (parsed.origin !== "https://chatgpt.com" || parsed.pathname !== "/backend-api/codex/responses") reject(); } else if ( parsed.protocol !== "http:" || + parsed.pathname !== "/codex/responses" || (parsed.hostname !== "127.0.0.1" && parsed.hostname !== "localhost" && parsed.hostname !== "[::1]") ) { reject(); @@ -224,32 +226,46 @@ export function createProviderFetchWitness(options: ProviderFetchWitnessOptions) const onTerminalFailure = config.onTerminalFailure as (code: ProviderFetchFailureCode) => void; let terminal = false; + let requestInProgress = false; let requestCount = 0; let receipt: ProviderFetchReceipt | null = null; - const failClosed = (): never => { + const failClosed = (code: ProviderFetchFailureCode): never => { + const notify = !terminal; terminal = true; receipt = null; - try { - onTerminalFailure("invalid-fetch"); - } catch { - // The generic thrown error remains authoritative and content-free. + if (notify) { + try { + onTerminalFailure(code); + } catch { + // The generic thrown error remains authoritative and content-free. + } } reject(); }; const witnessedFetch: typeof fetch = async (input, init) => { - if (terminal || requestCount !== 0) failClosed(); + if (terminal || requestInProgress || requestCount !== 0) failClosed("invalid-fetch"); + requestInProgress = true; try { - const requestUrl = - typeof input === "string" ? input : input instanceof URL && !isProxy(input) ? input.href : reject(); + let requestUrl: string; + if (typeof input === "string") requestUrl = input; + else { + if (!(input instanceof URL) || isProxy(input) || Object.getPrototypeOf(input) !== URL.prototype) reject(); + requestUrl = URL.prototype.toString.call(input); + } if (requestUrl !== expectedUrl || !init || isProxy(init)) reject(); const request = dataRecord(init); if (request.method !== "POST") reject(); if (!request.headers || typeof request.headers !== "object" || isProxy(request.headers)) reject(); - if (!(request.headers instanceof Headers) || Object.getPrototypeOf(request.headers) !== Headers.prototype) + if ( + !(request.headers instanceof Headers) || + Object.getPrototypeOf(request.headers) !== Headers.prototype || + Reflect.ownKeys(request.headers).length !== 0 + ) { reject(); - const rawEncoding = request.headers.get("content-encoding"); + } + const rawEncoding = Headers.prototype.get.call(request.headers, "content-encoding"); if (rawEncoding !== null && rawEncoding !== "zstd") reject(); const encoding = rawEncoding === "zstd" ? "zstd" : "identity"; const physical = physicalBody(request.body); @@ -263,6 +279,7 @@ export function createProviderFetchWitness(options: ProviderFetchWitnessOptions) if (!Buffer.from(text, "utf8").equals(logical)) reject(); const expected = logicalReceipt(getLogicalReceipt()); + if (terminal) reject(); if (expected.byteLength !== logical.byteLength || expected.payloadDigest !== sha256(logical)) reject(); const payload: unknown = JSON.parse(text); const root = exactRecord(payload, PAYLOAD_KEYS); @@ -293,10 +310,9 @@ export function createProviderFetchWitness(options: ProviderFetchWitnessOptions) headers: new Headers(request.headers), body: Buffer.from(physical), } as RequestInit; - return fetchImpl(requestUrl, forwardedInit); - } catch (error) { - if (requestCount === 1 && receipt && !(error instanceof EvaluationProviderFetchWitnessError)) throw error; - return failClosed(); + return await fetchImpl(requestUrl, forwardedInit); + } catch { + return failClosed(requestCount === 1 ? "fetch-failed" : "invalid-fetch"); } }; From 4106abf8dd4e65c7bb25ccfaaa197b363b1db199 Mon Sep 17 00:00:00 2001 From: Sebastian Otaegui Date: Mon, 13 Jul 2026 13:48:28 -0300 Subject: [PATCH 5/5] fix(agent-journal): freeze fetch metadata before callbacks Snapshot native headers before invoking receipt callbacks and reject an in-flight first result when a concurrent duplicate terminally invalidates the witness. RED: npx vitest run packages/pi-agent-journal/__tests__/evaluation-provider-fetch-witness.test.ts -t snapshots Failures: callback-mutated headers were forwarded and the first concurrent request resolved after duplicate invalidation. --- .../evaluation-provider-fetch-witness.test.ts | 46 +++++++++++++++++++ .../evaluation-provider-fetch-witness.ts | 9 ++-- 2 files changed, 52 insertions(+), 3 deletions(-) diff --git a/packages/pi-agent-journal/__tests__/evaluation-provider-fetch-witness.test.ts b/packages/pi-agent-journal/__tests__/evaluation-provider-fetch-witness.test.ts index 9842554..b7750cf 100644 --- a/packages/pi-agent-journal/__tests__/evaluation-provider-fetch-witness.test.ts +++ b/packages/pi-agent-journal/__tests__/evaluation-provider-fetch-witness.test.ts @@ -240,6 +240,31 @@ describe("Agent Journal provider fetch witness", () => { expect(new Headers(forwardedInit?.headers).get("x-safe-metadata")).toBe("one"); }); + it("snapshots headers before receipt callbacks can mutate caller-owned state", async () => { + const body = logicalBody(); + const headers = new Headers(); + let forwardedHeaders: Headers | undefined; + const witness = createProviderFetchWitness({ + attemptId: "attempt-header-callback", + expectedUrl: "https://chatgpt.com/backend-api/codex/responses", + endpointClass: "openai-codex", + getLogicalReceipt: () => { + headers.set("content-encoding", "zstd"); + headers.set("x-injected", "private"); + return logicalReceipt(body); + }, + fetchImpl: vi.fn(async (_input, init) => { + forwardedHeaders = new Headers(init?.headers); + return new Response("ok"); + }), + onTerminalFailure: () => undefined, + }); + await witness.fetch("https://chatgpt.com/backend-api/codex/responses", { method: "POST", headers, body }); + expect(witness.getReceipt()?.encoding).toBe("identity"); + expect(forwardedHeaders?.get("content-encoding")).toBeNull(); + expect(forwardedHeaders?.get("x-injected")).toBeNull(); + }); + it("rejects an own Headers.get override without invoking it", async () => { const trap = vi.fn((_name: string) => null); const headers = new Headers(); @@ -350,6 +375,27 @@ describe("Agent Journal provider fetch witness", () => { expect(failures).toEqual(["fetch-failed"]); }); + it("invalidates an in-flight first result when a concurrent second request is attempted", async () => { + let release: () => void = () => undefined; + const pending = new Promise((resolve) => { + release = resolve; + }); + const inFlightFetch = vi.fn(async () => { + await pending; + return new Response("ok"); + }); + const { witness, body, failures } = directSetup({ fetchImpl: inFlightFetch }); + const init = { method: "POST", headers: new Headers(), body }; + const first = witness.fetch("https://chatgpt.com/backend-api/codex/responses", init); + const second = witness.fetch("https://chatgpt.com/backend-api/codex/responses", init); + await expect(second).rejects.toThrow(EvaluationProviderFetchWitnessError); + release(); + await expect(first).rejects.toThrow(EvaluationProviderFetchWitnessError); + expect(inFlightFetch).toHaveBeenCalledTimes(1); + expect(failures).toEqual(["invalid-fetch"]); + expect(witness.getReceipt()).toBeNull(); + }); + it("rejects a second provider request without forwarding it", async () => { const { witness, forwarded, body } = directSetup(); const init = { method: "POST", headers: new Headers(), body }; diff --git a/packages/pi-agent-journal/extensions/evaluation-provider-fetch-witness.ts b/packages/pi-agent-journal/extensions/evaluation-provider-fetch-witness.ts index 0912f20..9669909 100644 --- a/packages/pi-agent-journal/extensions/evaluation-provider-fetch-witness.ts +++ b/packages/pi-agent-journal/extensions/evaluation-provider-fetch-witness.ts @@ -265,7 +265,8 @@ export function createProviderFetchWitness(options: ProviderFetchWitnessOptions) ) { reject(); } - const rawEncoding = Headers.prototype.get.call(request.headers, "content-encoding"); + const headersSnapshot = new Headers(request.headers); + const rawEncoding = Headers.prototype.get.call(headersSnapshot, "content-encoding"); if (rawEncoding !== null && rawEncoding !== "zstd") reject(); const encoding = rawEncoding === "zstd" ? "zstd" : "identity"; const physical = physicalBody(request.body); @@ -307,10 +308,12 @@ export function createProviderFetchWitness(options: ProviderFetchWitnessOptions) }); const forwardedInit = { ...request, - headers: new Headers(request.headers), + headers: headersSnapshot, body: Buffer.from(physical), } as RequestInit; - return await fetchImpl(requestUrl, forwardedInit); + const response = await fetchImpl(requestUrl, forwardedInit); + if (terminal) reject(); + return response; } catch { return failClosed(requestCount === 1 ? "fetch-failed" : "invalid-fetch"); }