From b6e88a8d8443434dbe6661f5d523f327048a817f Mon Sep 17 00:00:00 2001 From: askalf <263217947+askalf@users.noreply.github.com> Date: Mon, 21 Sep 2026 22:30:26 -0400 Subject: [PATCH 1/7] feat: redact the OpenAI Responses API POST /v1/responses is what current OpenAI clients use by default and it passed through cordon verbatim, so a modern client pointed at cordon sent raw PII upstream with X-Redacted: 0. The route now normalizes to provider openai, dialect responses. The request walk covers instructions (under REDACT_SYSTEM), input as a string or item list (input_text parts, function_call arguments as parsed JSON, function_call_output text) and flat function tool definitions; image and file parts are left alone. Non-streaming replies restore output_text and refusal parts. Streaming restores output_text.delta through the existing hold-back buffer, re-emitting each delta under its original item_id, output_index and content_index, and restores the full text carried by output_text.done, content_part.done, output_item.done and response.completed. The echo stub answers /v1/responses in both shapes; the proxy suite gains 19 checks and the stream suite 8. --- CHANGELOG.md | 6 ++++ _stub-upstream.mjs | 33 +++++++++++++++++++ _test_proxy.mjs | 61 ++++++++++++++++++++++++++++++++++ _test_stream.mjs | 56 +++++++++++++++++++++++++++++++ src/index.ts | 4 +-- src/providers.ts | 71 +++++++++++++++++++++++++++++++++------- src/proxy.ts | 8 ++--- src/redact/apply.ts | 68 +++++++++++++++++++++++++++----------- src/redact/reidentify.ts | 8 ++--- src/streaming.ts | 62 +++++++++++++++++++++++++++++++++-- src/types.ts | 10 ++++++ 11 files changed, 344 insertions(+), 43 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index fcdb98e..a214fd4 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -11,6 +11,12 @@ to match, push a tag `vX.Y.Z`, and release.yml builds + pushes the GHCR image and creates the GitHub release from this file. --> +## [Unreleased] + +### Added + +- `POST /v1/responses` (the OpenAI Responses API) is redacted like `/v1/chat/completions`. Current OpenAI clients (`client.responses.create`, the Agents SDK, Codex) use it by default, and until now it passed through verbatim, so a modern OpenAI client pointed at cordon sent raw PII upstream with `X-Redacted: 0`. The request walk covers `instructions` (under `REDACT_SYSTEM`, with `system` and `developer` items), `input` as a string or an item list (`input_text` parts, `function_call` arguments as parsed JSON, `function_call_output` text) and flat function tool definitions; `input_image` and `input_file` parts are left untouched. Non-streaming replies restore `output_text` and `refusal` parts; streaming restores `response.output_text.delta` through the same hold-back buffer as chat, re-emitting each delta under its original `item_id` / `output_index` / `content_index`, and restores the full text carried by `output_text.done`, `content_part.done`, `output_item.done` and `response.completed`. Same modes, headers and audit record (provider `openai`). Sub-paths such as `/v1/responses/{id}` still pass through verbatim. + ## [0.2.1] - 2026-09-11 ### Changed diff --git a/_stub-upstream.mjs b/_stub-upstream.mjs index 8ad9450..2b5f680 100644 --- a/_stub-upstream.mjs +++ b/_stub-upstream.mjs @@ -26,9 +26,22 @@ function collectText(body, provider) { }; if (provider === "anthropic" && body.system) pushContent(body.system); for (const m of body.messages || []) if (m.role === "user") pushContent(m.content); + // Responses API: `input` is a string or a list of items (messages, function calls, outputs). + if (typeof body.input === "string") parts.push(body.input); + else for (const item of body.input || []) { + if (item?.role === "user") pushContent(item.content); + if (item?.type === "function_call_output" && typeof item.output === "string") parts.push(item.output); + } return parts.join(" "); } +const responsesBody = (n, text) => ({ + id: "resp_stub" + n, object: "response", model: "gpt-4o-mini", status: "completed", + output: [{ id: "msg_stub" + n, type: "message", role: "assistant", status: "completed", + content: [{ type: "output_text", text, annotations: [] }] }], + usage: { input_tokens: 10, output_tokens: 5, total_tokens: 15 }, _stub_call: n, +}); + const openaiBody = (n, text) => ({ id: "stub-" + n, object: "chat.completion", model: "gpt-4o-mini", choices: [{ index: 0, message: { role: "assistant", content: text }, finish_reason: "stop" }], @@ -57,6 +70,24 @@ function streamOpenAI(res, text) { res.write("data: [DONE]\n\n"); res.end(); } +function streamResponses(res, text, n) { + res.setHeader("content-type", "text/event-stream"); + let seq = 0; + const f = (type, d) => res.write(`event: ${type}\ndata: ${JSON.stringify({ type, sequence_number: seq++, ...d })}\n\n`); + const addr = { item_id: "msg_stub" + n, output_index: 0, content_index: 0 }; + const shell = (status) => ({ id: "resp_stub" + n, object: "response", model: "gpt-4o-mini", status, output: [] }); + f("response.created", { response: shell("in_progress") }); + f("response.in_progress", { response: shell("in_progress") }); + f("response.output_item.added", { output_index: 0, item: { id: addr.item_id, type: "message", role: "assistant", status: "in_progress", content: [] } }); + f("response.content_part.added", { ...addr, part: { type: "output_text", text: "", annotations: [] } }); + for (const c of chunk3(text)) f("response.output_text.delta", { ...addr, delta: c }); + f("response.output_text.done", { ...addr, text }); + f("response.content_part.done", { ...addr, part: { type: "output_text", text, annotations: [] } }); + const item = { id: addr.item_id, type: "message", role: "assistant", status: "completed", content: [{ type: "output_text", text, annotations: [] }] }; + f("response.output_item.done", { output_index: 0, item }); + f("response.completed", { response: { ...shell("completed"), output: [item], usage: { input_tokens: 10, output_tokens: 5, total_tokens: 15 } } }); + res.end(); +} function streamAnthropic(res, text, body = {}) { res.setHeader("content-type", "text/event-stream"); const f = (e, d) => res.write(`event: ${e}\ndata: ${JSON.stringify({ type: e, ...d })}\n\n`); @@ -109,6 +140,8 @@ http // Echo the received text back as the assistant reply. if (req.url.includes("/chat/completions")) return body.stream ? streamOpenAI(res, text) : json(res, openaiBody(n, text)); + if (req.url.includes("/responses")) + return body.stream ? streamResponses(res, text, n) : json(res, responsesBody(n, text)); if (req.url.includes("/messages")) return body.stream ? streamAnthropic(res, text, body) : json(res, anthropicBody(n, text)); res.statusCode = 404; diff --git a/_test_proxy.mjs b/_test_proxy.mjs index 7d88096..453612f 100644 --- a/_test_proxy.mjs +++ b/_test_proxy.mjs @@ -24,6 +24,8 @@ const post = postTo(BASE); const post2 = postTo(BASE2); const aBody = (text, extra = {}) => ({ model: "claude-haiku-4-5", messages: [{ role: "user", content: text }], ...extra }); const oBody = (text, extra = {}) => ({ model: "gpt-4o-mini", messages: [{ role: "user", content: text }], ...extra }); +const rBody = (input, extra = {}) => ({ model: "gpt-4o-mini", input, ...extra }); +const rText = (j) => j?.output?.[0]?.content?.[0]?.text || ""; const setTenantOn = (base) => (patch) => fetch(base + "/admin/tenant", { method: "POST", headers: { "content-type": "application/json", "x-admin-token": ADMIN }, body: JSON.stringify(patch) }); const setTenant = setTenantOn(BASE); @@ -82,6 +84,63 @@ const setTenant2 = setTenantOn(BASE2); sent = JSON.stringify(await calls()); ok("off: upstream saw RAW (verbatim passthrough)", sent.includes("john@acme.com")); + // ---- reversible (OpenAI Responses API): string input ---- + await reset(); + res = await post("/v1/responses", rBody(PII)); + text = rText(await res.json()); + ok("responses/string: email restored in output_text", text.includes("john@acme.com") && !text.includes("= 2", Number(res.headers.get("x-redacted")) >= 2, res.headers.get("x-redacted")); + ok("responses/string: X-Redacted-Types lists EMAIL and CREDIT_CARD", + /EMAIL:1/.test(res.headers.get("x-redacted-types") || "") && /CREDIT_CARD:1/.test(res.headers.get("x-redacted-types") || ""), + res.headers.get("x-redacted-types")); + sent = JSON.stringify(await calls()); + ok("responses/string: upstream saw placeholder", //.test(sent)); + ok("responses/string: upstream NEVER saw raw PII", !sent.includes("john@acme.com") && !sent.includes("4012888888881881")); + + // ---- reversible (Responses): item list with input_text parts, instructions, a function output ---- + await reset(); + res = await post("/v1/responses", rBody( + [ + { role: "user", content: [{ type: "input_text", text: PII }, { type: "input_image", image_url: "data:image/png;base64,AAAA" }] }, + { type: "function_call", call_id: "call_1", name: "lookup", arguments: JSON.stringify({ email: "jane@corp.io" }) }, + { type: "function_call_output", call_id: "call_1", output: "account for jane@corp.io: card 4012888888881881" }, + ], + { instructions: "You help ops@acme.com triage tickets." }, + )); + text = rText(await res.json()); + ok("responses/items: input_text restored in reply", text.includes("john@acme.com") && !text.includes(" i.type === "function_call"); + let args; + try { args = JSON.parse(fc?.arguments ?? ""); } catch {} + ok("responses/items: function_call arguments stay valid JSON with a placeholder", typeof args?.email === "string" && /^= 5, res.headers.get("x-redacted")); + + // ---- strip / off (Responses) ---- + await reset(); + res = await post("/v1/responses", rBody(PII), { "x-redact-mode": "strip" }); + text = rText(await res.json()); + ok("responses/strip: placeholders persist (not restored)", text.includes("[EMAIL]") && !text.includes("john@acme.com")); + sent = JSON.stringify(await calls()); + ok("responses/strip: upstream saw [EMAIL], not raw", sent.includes("[EMAIL]") && !sent.includes("john@acme.com")); + await reset(); + res = await post("/v1/responses", rBody(PII), { "x-redact-mode": "off" }); + text = rText(await res.json()); + ok("responses/off: reply echoes raw (nothing redacted)", text.includes("john@acme.com")); + ok("responses/off: X-Redacted 0", res.headers.get("x-redacted") === "0"); + + // ---- fail-closed (Responses) ---- + await reset(); + res = await post("/v1/responses", rBody(PII), { "x-cordon-fail": "1" }); + ok("responses/fail-closed: status 422", res.status === 422); + ok("responses/fail-closed: upstream NOT called", (await calls()).total === 0); + // ---- fail-closed ---- await reset(); res = await post("/v1/messages", aBody(PII), { "x-cordon-fail": "1" }); @@ -152,6 +211,8 @@ const setTenant2 = setTenantOn(BASE2); ok("audit: log contains NO raw email", !log.includes("john@acme.com")); ok("audit: log contains NO raw card", !log.includes("4012888888881881")); ok("audit: log contains NO raw jane", !log.includes("jane@corp.io")); + ok("audit: log contains NO raw ops address from Responses instructions", !log.includes("ops@acme.com")); + ok("audit: Responses requests recorded under provider openai", log.split("\n").some((l) => l.includes('"provider":"openai"') && l.includes('"CREDIT_CARD":1'))); console.log(`\n${pass} passed, ${fail} failed`); process.exit(fail ? 1 : 0); diff --git a/_test_stream.mjs b/_test_stream.mjs index af72b93..a9d7990 100644 --- a/_test_stream.mjs +++ b/_test_stream.mjs @@ -18,6 +18,7 @@ const post = (path, body, headers = {}) => }); const aBody = (text, extra = {}) => ({ model: "claude-haiku-4-5", messages: [{ role: "user", content: text }], stream: true, ...extra }); const oBody = (text, extra = {}) => ({ model: "gpt-4o-mini", messages: [{ role: "user", content: text }], stream: true, ...extra }); +const rBody = (text, extra = {}) => ({ model: "gpt-4o-mini", input: text, stream: true, ...extra }); /** Reconstruct assistant text from an SSE response body. */ function reconstruct(sse, provider) { @@ -35,6 +36,8 @@ function reconstruct(sse, provider) { } if (provider === "anthropic") { if (j.type === "content_block_delta") out += j.delta?.text ?? ""; + } else if (provider === "responses") { + if (j.type === "response.output_text.delta") out += j.delta ?? ""; } else { out += j.choices?.[0]?.delta?.content ?? ""; } @@ -42,6 +45,36 @@ function reconstruct(sse, provider) { return out; } +/** Every Responses frame that carries the full text, in stream order. A client SDK + * reads the final text from these, not from the deltas, so each must be restored. */ +function responsesFullTexts(sse) { + const out = []; + for (const frame of sse.split("\n\n")) { + const line = frame.split("\n").find((l) => l.startsWith("data:")); + if (!line) continue; + let j; + try { j = JSON.parse(line.slice(5).trim()); } catch { continue; } + if (j.type === "response.output_text.done") out.push(["output_text.done", j.text]); + if (j.type === "response.content_part.done") out.push(["content_part.done", j.part?.text]); + if (j.type === "response.output_item.done") out.push(["output_item.done", j.item?.content?.[0]?.text]); + if (j.type === "response.completed") out.push(["completed", j.response?.output?.[0]?.content?.[0]?.text]); + } + return out; +} + +/** Addressing on every re-emitted Responses delta must match the upstream's. */ +function responsesDeltaAddressing(sse) { + const seen = new Set(); + for (const frame of sse.split("\n\n")) { + const line = frame.split("\n").find((l) => l.startsWith("data:")); + if (!line) continue; + let j; + try { j = JSON.parse(line.slice(5).trim()); } catch { continue; } + if (j.type === "response.output_text.delta") seen.add(`${j.item_id}/${j.output_index}/${j.content_index}`); + } + return [...seen]; +} + /** Strict SSE validator mimicking the client SDK: a `text_delta` must land on a * `text` block, a `thinking_delta` on a `thinking` block. Catches the exact failure * ("Content block is not a text block") cordon hit with real CLI traffic. */ @@ -103,6 +136,29 @@ function validateSSE(sse) { sent = JSON.stringify(await calls()); ok("stream/openai: upstream saw placeholder not raw", //.test(sent) && !sent.includes("john@acme.com")); + // ---- reversible streaming (OpenAI Responses API) ---- + await reset(); + txt = await (await post("/v1/responses", rBody(PII))).text(); + text = reconstruct(txt, "responses"); + ok("stream/responses: email restored across frame split", text.includes("john@acme.com") && !text.includes(" f[0]).join(",") === "output_text.done,content_part.done,output_item.done,completed", full.map((f) => f[0]).join(",")); + ok("stream/responses: every full-text frame restored", full.every((f) => typeof f[1] === "string" && f[1].includes("john@acme.com") && !f[1].includes("/.test(sent) && !sent.includes("john@acme.com")); + + // ---- strip streaming (Responses) ---- + await reset(); + txt = await (await post("/v1/responses", rBody(PII), { "x-redact-mode": "strip" })).text(); + text = reconstruct(txt, "responses"); + ok("stream/responses/strip: placeholders persist", text.includes("[EMAIL]") && !text.includes("john@acme.com")); + // ---- strip streaming: placeholders persist, no restore, no hold-back ---- await reset(); txt = await (await post("/v1/messages", aBody(PII), { "x-redact-mode": "strip" })).text(); diff --git a/src/index.ts b/src/index.ts index c3872bd..e39ae41 100644 --- a/src/index.ts +++ b/src/index.ts @@ -138,9 +138,9 @@ async function passthroughUnknown(req: any, reply: any, method: string) { app.post("/v1/*", async (req, reply) => { const { headers, bare, path } = reqParts(req); - // Only the two generation endpoints are redacted; everything else (count_tokens, + // Only the three generation endpoints are redacted; everything else (count_tokens, // embeddings, …) forwards verbatim so it is never normalize-mangled. - if (bare !== "/v1/chat/completions" && bare !== "/v1/messages") { + if (bare !== "/v1/chat/completions" && bare !== "/v1/responses" && bare !== "/v1/messages") { return passthroughUnknown(req, reply, "POST"); } diff --git a/src/providers.ts b/src/providers.ts index 21b6f53..5805e37 100644 --- a/src/providers.ts +++ b/src/providers.ts @@ -1,19 +1,21 @@ import { config } from "./config"; import { sha256 } from "./util"; import { getPolicy } from "./policy"; -import type { CanonicalRequest, Provider, RedactMode, RedactSet } from "./types"; +import type { CanonicalRequest, Dialect, Provider, RedactMode, RedactSet } from "./types"; /** - * A ProviderAdapter knows the wire dialect of one provider: how to read a text - * delta off a streamed frame, and — the inverse cordon needs — how to synthesize a + * A ProviderAdapter knows one wire dialect: how to read a text delta off a + * streamed frame, and — the inverse cordon needs — how to synthesize a * text-carrying SSE frame from a (re-identified) text chunk. */ export interface ProviderAdapter { /** Read a streamed SSE `data:` payload. */ parseDelta(data: string): { textDelta?: string; done: boolean }; - /** Build a provider-correct SSE frame carrying one assistant-text chunk at `index` - * (the content-block index; ignored by providers without block indices). */ - frameFromText(text: string, index?: number): string; + /** Build a dialect-correct SSE frame carrying one assistant-text chunk at `index` + * (the content-block index; ignored by dialects without block indices). `ctx` is + * the addressing the dialect needs beyond an index (Responses: item_id, + * output_index, content_index), copied from the frame being re-emitted. */ + frameFromText(text: string, index?: number, ctx?: Record): string; /** Walk a non-streaming response body's assistant-text fields (for re-identify). */ responseTextSlots(body: any): Array<{ get(): string; set(v: string): void }>; } @@ -74,7 +76,41 @@ export const anthropic: ProviderAdapter = { }, }; -export const adapterFor = (p: Provider) => (p === "openai" ? openai : anthropic); +// ----------------------------- OpenAI (responses) ----------------------------- +export const openaiResponses: ProviderAdapter = { + parseDelta(data) { + try { + const j = JSON.parse(data); + if (j.type === "response.output_text.delta") return { textDelta: j.delta ?? "", done: false }; + if (j.type === "response.completed") return { done: true }; + return { done: false }; + } catch { + return { done: false }; + } + }, + frameFromText(text, _index = 0, ctx = {}) { + // The Responses stream addresses a delta by item_id / output_index / content_index, + // not by a single block index; `ctx` carries those from the frame being re-emitted. + const data = { type: "response.output_text.delta", ...ctx, delta: text }; + return `event: response.output_text.delta\ndata: ${JSON.stringify(data)}\n\n`; + }, + responseTextSlots(body) { + const slots: Array<{ get(): string; set(v: string): void }> = []; + for (const item of body?.output ?? []) { + if (item?.type !== "message" || !Array.isArray(item.content)) continue; + for (const part of item.content) { + if (part?.type === "output_text" && typeof part.text === "string") + slots.push({ get: () => part.text, set: (v) => (part.text = v) }); + else if (part?.type === "refusal" && typeof part.refusal === "string") + slots.push({ get: () => part.refusal, set: (v) => (part.refusal = v) }); + } + } + return slots; + }, +}; + +export const adapterFor = (p: Provider, dialect?: Dialect) => + p === "anthropic" ? anthropic : dialect === "responses" ? openaiResponses : openai; // ----------------------------- normalize: HTTP → CanonicalRequest ----------------------------- @@ -108,14 +144,22 @@ export function normalize( // EXACT endpoint match — sub-paths (e.g. /v1/messages/count_tokens) are NOT // generation requests and must take the transparent-passthrough route instead. const bare = path.split("?")[0]; - const provider: Provider | null = - bare === "/v1/chat/completions" ? "openai" : bare === "/v1/messages" ? "anthropic" : null; - if (!provider || !body || typeof body !== "object") return null; + const route: [Provider, Dialect] | null = + bare === "/v1/chat/completions" + ? ["openai", "chat"] + : bare === "/v1/responses" + ? ["openai", "responses"] + : bare === "/v1/messages" + ? ["anthropic", "messages"] + : null; + if (!route || !body || typeof body !== "object") return null; + const [provider, dialect] = route; const tenant = resolveTenant(headers); return { provider, + dialect, model: body.model, tenant, mode: resolveMode(headers, tenant), @@ -178,7 +222,12 @@ export function forwardUpstream(r: CanonicalRequest, bodyOverride?: any): Promis */ export function passthroughBase(path: string, headers: Record): string { if (path.startsWith("/v1/messages")) return config.upstream.anthropic; - if (path.startsWith("/v1/chat") || path.startsWith("/v1/embeddings") || path.startsWith("/v1/completions")) + if ( + path.startsWith("/v1/chat") || + path.startsWith("/v1/responses") || + path.startsWith("/v1/embeddings") || + path.startsWith("/v1/completions") + ) return config.upstream.openai; return headers["x-api-key"] ? config.upstream.anthropic : config.upstream.openai; } diff --git a/src/proxy.ts b/src/proxy.ts index 52c58c5..37353fd 100644 --- a/src/proxy.ts +++ b/src/proxy.ts @@ -18,7 +18,7 @@ import type { CanonicalRequest, HttpRes } from "./types"; */ export async function handle(r: CanonicalRequest, res: HttpRes) { const t0 = nowMs(); - const adapter = adapterFor(r.provider); + const adapter = adapterFor(r.provider, r.dialect); metrics.request(r.mode); res.setHeader("X-Redact-Mode", r.mode); @@ -51,7 +51,7 @@ export async function handle(r: CanonicalRequest, res: HttpRes) { try { if (r.testFail) throw new Error("forced detection failure (test hook)"); vault = new Vault(r.mode, { consistentPseudonyms, secret: config.tenantSecret }); - ({ deidBody, spans } = applyRedaction(r.raw, r.provider, vault, r.activeSets, detector, redactSystem)); + ({ deidBody, spans } = applyRedaction(r.raw, r.provider, vault, r.activeSets, detector, redactSystem, r.dialect)); } catch (e) { return onRedactionError(r, res, e, failMode, t0); } @@ -102,7 +102,7 @@ export async function handle(r: CanonicalRequest, res: HttpRes) { // ---- relay the response (re-identify in reversible mode) ---- if (r.stream && up.body) { - if (r.mode === "reversible") await captureAndReidentify(up.body, res, adapter, vault, r.provider); + if (r.mode === "reversible") await captureAndReidentify(up.body, res, adapter, vault, r.provider, r.dialect); else await pipeUpstream(up, res); // strip: placeholders persist, verbatim metrics.timing("stream", t0); return; @@ -110,7 +110,7 @@ export async function handle(r: CanonicalRequest, res: HttpRes) { const body = await up.json(); res.setHeader("content-type", "application/json"); - res.end(JSON.stringify(r.mode === "reversible" ? reidentifyBody(body, r.provider, vault) : body)); + res.end(JSON.stringify(r.mode === "reversible" ? reidentifyBody(body, r.provider, vault, r.dialect) : body)); metrics.timing("full", t0); } diff --git a/src/redact/apply.ts b/src/redact/apply.ts index 5d306af..a5af491 100644 --- a/src/redact/apply.ts +++ b/src/redact/apply.ts @@ -1,5 +1,5 @@ import { clone } from "../util"; -import type { Detector, Provider, RedactSet, Span } from "../types"; +import type { Detector, Dialect, Provider, RedactSet, Span } from "../types"; import type { Vault } from "./vault"; type Slot = { get(): string; set(v: string): void; numeric?: boolean }; @@ -48,18 +48,36 @@ function pushStringLeaves(node: any, slots: Slot[], depth = 0): void { } } +// A JSON-string field (tool-call arguments): parse → redact its leaves (incl. NUMERIC +// PII) → re-serialize, so a redacted number becomes a QUOTED "" and the args +// stay valid JSON (a textual replace left an unquoted placeholder). Unparseable args +// fall back to text redaction. +function pushJsonString(obj: any, key: string, slots: Slot[], finalizers: (() => void)[]): void { + let parsed: any; + try { parsed = JSON.parse(obj[key]); } catch { parsed = undefined; } + if (parsed && typeof parsed === "object") { + pushStringLeaves(parsed, slots); + finalizers.push(() => { obj[key] = JSON.stringify(parsed); }); + } else { + slots.push(slot(obj, key)); + } +} + /** * Collect every REDACTABLE text field in a provider REQUEST body. Walks message * content (string or content-part array), Anthropic system blocks + tool_result - * content, AND every model-visible structured field that can carry user data: - * OpenAI message `name` + assistant `tool_calls[].function.arguments`, tool - * definitions (descriptions + parameter schemas), and Anthropic `tool_use` inputs. - * Only image parts and raw provider-auth headers are intentionally left untouched. + * content, Responses `instructions` + `input` items, AND every model-visible + * structured field that can carry user data: OpenAI message `name` + assistant + * `tool_calls[].function.arguments`, Responses `function_call` arguments and + * `function_call_output` output, tool definitions (descriptions + parameter + * schemas), and Anthropic `tool_use` inputs. Only image / file parts and raw + * provider-auth headers are intentionally left untouched. */ function requestTextSlots( body: any, provider: Provider, redactSystem: boolean, + dialect: Dialect, ): { slots: Slot[]; finalizers: (() => void)[] } { const slots: Slot[] = []; const finalizers: (() => void)[] = []; // run after redaction (re-serialize parsed JSON-string fields) @@ -73,7 +91,7 @@ function requestTextSlots( const part = c[i]; if (typeof part === "string") { slots.push(slot(c, i)); continue; } // a bare-string content element if (!part || typeof part !== "object") continue; - if (part.type === "text" && typeof part.text === "string") { + if ((part.type === "text" || part.type === "input_text") && typeof part.text === "string") { slots.push(slot(part, "text")); } else if (part.type === "tool_result") { // Anthropic tool_result content can itself be a string or block array. @@ -100,6 +118,28 @@ function requestTextSlots( if (t && typeof t.description === "string") slots.push(slot(t, "description")); if (t && t.input_schema && typeof t.input_schema === "object") pushStringLeaves(t.input_schema, slots); } + } else if (dialect === "responses") { + // Responses: `instructions` is the system prompt; tools are flat function objects. + if (redactSystem && typeof body.instructions === "string") slots.push(slot(body, "instructions")); + if (Array.isArray(body.tools)) + for (const t of body.tools) { + if (t?.type !== "function") continue; + if (typeof t.description === "string") slots.push(slot(t, "description")); + if (t.parameters && typeof t.parameters === "object") pushStringLeaves(t.parameters, slots); + } + // `input` is a string, or a list of items: messages (string or part-array content), + // assistant function calls (JSON-string arguments) and their outputs. + if (typeof body.input === "string") slots.push(slot(body, "input")); + else if (Array.isArray(body.input)) + for (const item of body.input) { + if (!item || typeof item !== "object") continue; + if (!redactSystem && (item.role === "system" || item.role === "developer")) continue; + if (typeof item.content === "string" || Array.isArray(item.content)) pushContent(item, "content"); + if (item.type === "function_call" && typeof item.arguments === "string") + pushJsonString(item, "arguments", slots, finalizers); + if (item.type === "function_call_output" && typeof item.output === "string") slots.push(slot(item, "output")); + } + return { slots, finalizers }; } else { // OpenAI tool definitions: description + parameter schema string leaves. if (Array.isArray(body.tools)) @@ -121,18 +161,7 @@ function requestTextSlots( for (const tc of msg.tool_calls) { const fn = tc?.function; if (!fn || typeof fn.arguments !== "string") continue; - // arguments is a JSON STRING. Parse → redact its leaves (incl. NUMERIC PII) - // → re-serialize, so a redacted number becomes a QUOTED "" and the - // args stay valid JSON (a textual replace left an unquoted placeholder). - // Unparseable args fall back to text redaction. - let parsed: any; - try { parsed = JSON.parse(fn.arguments); } catch { parsed = undefined; } - if (parsed && typeof parsed === "object") { - pushStringLeaves(parsed, slots); - finalizers.push(() => { fn.arguments = JSON.stringify(parsed); }); - } else { - slots.push(slot(fn, "arguments")); - } + pushJsonString(fn, "arguments", slots, finalizers); } } @@ -169,9 +198,10 @@ export function applyRedaction( activeSets: RedactSet[], detector: Detector, redactSystem = true, + dialect: Dialect = provider === "anthropic" ? "messages" : "chat", ): RedactionResult { const deidBody = clone(rawBody); - const { slots, finalizers } = requestTextSlots(deidBody, provider, redactSystem); + const { slots, finalizers } = requestTextSlots(deidBody, provider, redactSystem, dialect); const all: Span[] = []; for (const sl of slots) { diff --git a/src/redact/reidentify.ts b/src/redact/reidentify.ts index 1bebe96..f6161d8 100644 --- a/src/redact/reidentify.ts +++ b/src/redact/reidentify.ts @@ -1,10 +1,10 @@ import { clone } from "../util"; import { adapterFor } from "../providers"; import { PLACEHOLDER_RE, isFormingPlaceholder, Vault } from "./vault"; -import type { Provider } from "../types"; +import type { Dialect, Provider } from "../types"; /** Replace every placeholder we minted with its real value; leave unknown ones as-is. */ -function restore(text: string, vault: Vault): string { +export function restore(text: string, vault: Vault): string { return text.replace(PLACEHOLDER_RE, (m) => vault.lookup(m) ?? m); } @@ -12,10 +12,10 @@ function restore(text: string, vault: Vault): string { * Full-body re-identification for a NON-streaming response: walk the provider * response's assistant-text fields and restore real values in place. */ -export function reidentifyBody(body: any, provider: Provider, vault: Vault): any { +export function reidentifyBody(body: any, provider: Provider, vault: Vault, dialect?: Dialect): any { if (!vault.hasReverse) return body; // nothing was redacted → nothing to restore const out = clone(body); - for (const sl of adapterFor(provider).responseTextSlots(out)) sl.set(restore(sl.get(), vault)); + for (const sl of adapterFor(provider, dialect).responseTextSlots(out)) sl.set(restore(sl.get(), vault)); return out; } diff --git a/src/streaming.ts b/src/streaming.ts index baa4862..8241368 100644 --- a/src/streaming.ts +++ b/src/streaming.ts @@ -1,6 +1,6 @@ -import type { HttpRes, Provider } from "./types"; +import type { Dialect, HttpRes, Provider } from "./types"; import type { ProviderAdapter } from "./providers"; -import { StreamReidentifier } from "./redact/reidentify"; +import { reidentifyBody, restore, StreamReidentifier } from "./redact/reidentify"; import type { Vault } from "./redact/vault"; /** Pipe an upstream Response straight to the client, verbatim (strip / off / errors). */ @@ -23,6 +23,11 @@ export async function pipeUpstream(up: Response, res: HttpRes): Promise { } } +/** The same SSE frame with its `data:` payload replaced by `j` (event: lines kept). */ +function reframe(frame: string, j: unknown): string { + return frame.replace(/^data:.*$/m, `data: ${JSON.stringify(j)}`); +} + /** * Reversible streaming: tee the upstream SSE stream while restoring real values in * flight. Text-carrying frames are suppressed and re-emitted (re-identified) via the @@ -37,6 +42,7 @@ export async function captureAndReidentify( adapter: ProviderAdapter, vault: Vault, provider: Provider, + dialect: Dialect = provider === "anthropic" ? "messages" : "chat", ): Promise { res.setHeader("content-type", "text/event-stream"); res.setHeader("cache-control", "no-cache"); @@ -47,8 +53,11 @@ export async function captureAndReidentify( // tool_use blocks); a re-emitted text_delta MUST carry the matching block index or the // client SDK throws "Content block is not a text block". OpenAI has no block index → 0. let emitIndex = 0; + // Responses addresses a delta by item_id / output_index / content_index; a re-emitted + // delta carries the addressing of the frame it stands in for. + let emitCtx: Record = {}; const emitText = (chunk: string, index = emitIndex) => { - if (chunk) res.write(adapter.frameFromText(chunk, index)); + if (chunk) res.write(adapter.frameFromText(chunk, index, emitCtx)); }; const flushTail = () => emitText(reider.end(), emitIndex); @@ -92,6 +101,53 @@ export async function captureAndReidentify( return; } + if (dialect === "responses") { + let j: any; + try { + j = JSON.parse(data); + } catch {} + const type = j?.type; + if (type === "response.output_text.delta" || type === "response.refusal.delta") { + // Suppress the original delta, re-emit the restorable prefix under the same + // addressing; the possibly-forming tail stays held in the re-identifier. + const { type: _t, delta: _d, ...ctx } = j; + emitCtx = ctx; + emitText(reider.push(j.delta ?? "")); + return; + } + // Every frame below carries the text in full, so the held tail is flushed first + // (a placeholder split across deltas is then already restored on the client) and + // the frame's own text is restored before it passes. + if (type === "response.output_text.done" || type === "response.refusal.done") { + flushTail(); + if (typeof j.text === "string") j.text = restore(j.text, vault); + if (typeof j.refusal === "string") j.refusal = restore(j.refusal, vault); + res.write(reframe(frame, j)); + return; + } + if (type === "response.content_part.done" && j.part && typeof j.part === "object") { + flushTail(); + if (typeof j.part.text === "string") j.part.text = restore(j.part.text, vault); + if (typeof j.part.refusal === "string") j.part.refusal = restore(j.part.refusal, vault); + res.write(reframe(frame, j)); + return; + } + if (type === "response.output_item.done" && j.item && typeof j.item === "object") { + flushTail(); + res.write(reframe(frame, { ...j, item: reidentifyBody({ output: [j.item] }, provider, vault, dialect).output[0] })); + return; + } + if (type === "response.completed" || type === "response.incomplete" || type === "response.failed") { + flushTail(); + if (j.response && typeof j.response === "object") + res.write(reframe(frame, { ...j, response: reidentifyBody(j.response, provider, vault, dialect) })); + else res.write(frame); + return; + } + res.write(frame); // response.created / in_progress / output_item.added / content_part.added / function-call frames / ping + return; + } + // openai let ch: any; try { diff --git a/src/types.ts b/src/types.ts index 092d68d..2a3b68c 100644 --- a/src/types.ts +++ b/src/types.ts @@ -1,5 +1,13 @@ export type Provider = "openai" | "anthropic"; +/** + * The wire dialect of a generation request. A provider can speak more than one: + * OpenAI serves both chat.completions (`messages`) and the Responses API + * (`input`/`output`), which differ in request shape, response shape and SSE + * event grammar while sharing the upstream base, auth headers and audit provider. + */ +export type Dialect = "chat" | "responses" | "messages"; + export interface Msg { role: string; content: any; @@ -58,6 +66,8 @@ export interface Detector { export interface CanonicalRequest { provider: Provider; + /** Which of the provider's wire shapes this request uses (decided by path). */ + dialect: Dialect; model: string; tenant: string; /** Resolved: X-Redact-Mode header > tenant policy > config.defaultMode. */ From 7b6987d93bd23dce68f47ee9b8f1eac7c9b8c502 Mon Sep 17 00:00:00 2001 From: askalf <263217947+askalf@users.noreply.github.com> Date: Tue, 22 Sep 2026 03:17:32 +0000 Subject: [PATCH 2/7] fix: preserve the Responses event kind when re-emitting a restored delta A streamed response.refusal.delta shared the re-emit path with response.output_text.delta, but frameFromText hardcoded the output-text type on both the SSE event line and the data payload. A refusal reached the client as ordinary output text, so a client consuming refusal events never saw its incremental refusal. streaming.ts now keeps the upstream type in the emission context, and the Responses adapter emits that type (falling back to output_text for any other caller). Refs FIX-e5e0127b. --- src/providers.ts | 21 ++++++++++++++++----- src/streaming.ts | 6 ++++-- 2 files changed, 20 insertions(+), 7 deletions(-) diff --git a/src/providers.ts b/src/providers.ts index 5805e37..eea8b7a 100644 --- a/src/providers.ts +++ b/src/providers.ts @@ -14,7 +14,8 @@ export interface ProviderAdapter { /** Build a dialect-correct SSE frame carrying one assistant-text chunk at `index` * (the content-block index; ignored by dialects without block indices). `ctx` is * the addressing the dialect needs beyond an index (Responses: item_id, - * output_index, content_index), copied from the frame being re-emitted. */ + * output_index, content_index) plus the upstream event `type`, copied from the + * frame being re-emitted. */ frameFromText(text: string, index?: number, ctx?: Record): string; /** Walk a non-streaming response body's assistant-text fields (for re-identify). */ responseTextSlots(body: any): Array<{ get(): string; set(v: string): void }>; @@ -77,11 +78,16 @@ export const anthropic: ProviderAdapter = { }; // ----------------------------- OpenAI (responses) ----------------------------- +/** Streamed Responses events carrying restorable assistant text. Output text and a + * refusal are DISTINCT event types and a client consuming refusals reads only its + * own, so a refusal delta must be re-emitted as a refusal delta. */ +const RESPONSES_TEXT_DELTAS = new Set(["response.output_text.delta", "response.refusal.delta"]); + export const openaiResponses: ProviderAdapter = { parseDelta(data) { try { const j = JSON.parse(data); - if (j.type === "response.output_text.delta") return { textDelta: j.delta ?? "", done: false }; + if (RESPONSES_TEXT_DELTAS.has(j.type)) return { textDelta: j.delta ?? "", done: false }; if (j.type === "response.completed") return { done: true }; return { done: false }; } catch { @@ -90,9 +96,14 @@ export const openaiResponses: ProviderAdapter = { }, frameFromText(text, _index = 0, ctx = {}) { // The Responses stream addresses a delta by item_id / output_index / content_index, - // not by a single block index; `ctx` carries those from the frame being re-emitted. - const data = { type: "response.output_text.delta", ...ctx, delta: text }; - return `event: response.output_text.delta\ndata: ${JSON.stringify(data)}\n\n`; + // not by a single block index; `ctx` carries those from the frame being re-emitted, + // along with that frame's event type so a refusal delta stays a refusal delta. + const type = + typeof ctx.type === "string" && RESPONSES_TEXT_DELTAS.has(ctx.type) + ? ctx.type + : "response.output_text.delta"; + const data = { ...ctx, type, delta: text }; + return `event: ${type}\ndata: ${JSON.stringify(data)}\n\n`; }, responseTextSlots(body) { const slots: Array<{ get(): string; set(v: string): void }> = []; diff --git a/src/streaming.ts b/src/streaming.ts index 8241368..b8e8a1d 100644 --- a/src/streaming.ts +++ b/src/streaming.ts @@ -109,8 +109,10 @@ export async function captureAndReidentify( const type = j?.type; if (type === "response.output_text.delta" || type === "response.refusal.delta") { // Suppress the original delta, re-emit the restorable prefix under the same - // addressing; the possibly-forming tail stays held in the re-identifier. - const { type: _t, delta: _d, ...ctx } = j; + // addressing AND the same event type — a refusal delta re-emitted as output + // text would break a client that reads refusal events; the possibly-forming + // tail stays held in the re-identifier. + const { delta: _d, ...ctx } = j; emitCtx = ctx; emitText(reider.push(j.delta ?? "")); return; From 597803642eed99c441371e546d1d1c31f9b4ed39 Mon Sep 17 00:00:00 2001 From: askalf <263217947+askalf@users.noreply.github.com> Date: Tue, 22 Sep 2026 03:19:40 +0000 Subject: [PATCH 3/7] test: streaming Responses refusal keeps its event type and restores its text Adds a refusal-flavoured Responses stream to the stub (FORCE_REFUSAL in the input) and asserts the re-emitted deltas stay response.refusal.delta on both the event: line and the data payload, that no output_text delta leaks, and that the split placeholder is restored in the deltas and in refusal.done. Refs FIX-e5e0127b. --- _stub-upstream.mjs | 25 ++++++++++++++++++++++++- _test_stream.mjs | 40 ++++++++++++++++++++++++++++++++++++++++ 2 files changed, 64 insertions(+), 1 deletion(-) diff --git a/_stub-upstream.mjs b/_stub-upstream.mjs index 2b5f680..b1813b3 100644 --- a/_stub-upstream.mjs +++ b/_stub-upstream.mjs @@ -88,6 +88,27 @@ function streamResponses(res, text, n) { f("response.completed", { response: { ...shell("completed"), output: [item], usage: { input_tokens: 10, output_tokens: 5, total_tokens: 15 } } }); res.end(); } +// A REFUSED Responses turn streams the same frame shape under refusal-flavoured +// event types (`response.refusal.delta` / `.done`, a `refusal` content part). A +// client consuming refusals reads only those, so cordon must re-emit a restored +// refusal delta as a refusal delta — not as output text. +function streamResponsesRefusal(res, text, n) { + res.setHeader("content-type", "text/event-stream"); + let seq = 0; + const f = (type, d) => res.write(`event: ${type}\ndata: ${JSON.stringify({ type, sequence_number: seq++, ...d })}\n\n`); + const addr = { item_id: "msg_stub" + n, output_index: 0, content_index: 0 }; + const shell = (status) => ({ id: "resp_stub" + n, object: "response", model: "gpt-4o-mini", status, output: [] }); + f("response.created", { response: shell("in_progress") }); + f("response.output_item.added", { output_index: 0, item: { id: addr.item_id, type: "message", role: "assistant", status: "in_progress", content: [] } }); + f("response.content_part.added", { ...addr, part: { type: "refusal", refusal: "" } }); + for (const c of chunk3(text)) f("response.refusal.delta", { ...addr, delta: c }); + f("response.refusal.done", { ...addr, refusal: text }); + f("response.content_part.done", { ...addr, part: { type: "refusal", refusal: text } }); + const item = { id: addr.item_id, type: "message", role: "assistant", status: "completed", content: [{ type: "refusal", refusal: text }] }; + f("response.output_item.done", { output_index: 0, item }); + f("response.completed", { response: { ...shell("completed"), output: [item], usage: { input_tokens: 10, output_tokens: 5, total_tokens: 15 } } }); + res.end(); +} function streamAnthropic(res, text, body = {}) { res.setHeader("content-type", "text/event-stream"); const f = (e, d) => res.write(`event: ${e}\ndata: ${JSON.stringify({ type: e, ...d })}\n\n`); @@ -141,7 +162,9 @@ http if (req.url.includes("/chat/completions")) return body.stream ? streamOpenAI(res, text) : json(res, openaiBody(n, text)); if (req.url.includes("/responses")) - return body.stream ? streamResponses(res, text, n) : json(res, responsesBody(n, text)); + return body.stream + ? text.includes("FORCE_REFUSAL") ? streamResponsesRefusal(res, text, n) : streamResponses(res, text, n) + : json(res, responsesBody(n, text)); if (req.url.includes("/messages")) return body.stream ? streamAnthropic(res, text, body) : json(res, anthropicBody(n, text)); res.statusCode = 404; diff --git a/_test_stream.mjs b/_test_stream.mjs index a9d7990..32cbaa3 100644 --- a/_test_stream.mjs +++ b/_test_stream.mjs @@ -62,6 +62,24 @@ function responsesFullTexts(sse) { return out; } +/** Every re-emitted Responses delta as [eventLine, dataType, delta], in stream order. + * A refusal must stay a refusal on BOTH the `event:` line and the data payload — a + * client subscribing to refusal events reads nothing if either says output_text. */ +function responsesDeltaKinds(sse) { + const out = []; + for (const frame of sse.split("\n\n")) { + const lines = frame.split("\n"); + const ev = lines.find((l) => l.startsWith("event:")); + const line = lines.find((l) => l.startsWith("data:")); + if (!line) continue; + let j; + try { j = JSON.parse(line.slice(5).trim()); } catch { continue; } + if (j.type === "response.output_text.delta" || j.type === "response.refusal.delta") + out.push([ev?.slice(6).trim(), j.type, j.delta]); + } + return out; +} + /** Addressing on every re-emitted Responses delta must match the upstream's. */ function responsesDeltaAddressing(sse) { const seen = new Set(); @@ -153,6 +171,28 @@ function validateSSE(sse) { sent = JSON.stringify(await calls()); ok("stream/responses: upstream saw placeholder not raw", //.test(sent) && !sent.includes("john@acme.com")); + // ---- reversible streaming (Responses REFUSAL) ---- + // A refusal carries restorable text too, but under its own event type. Regression: + // the re-emit path shared frameFromText with output text and hardcoded the + // output-text type, so a refusal reached the client as ordinary output text. + await reset(); + txt = await (await post("/v1/responses", rBody("FORCE_REFUSAL cannot help with john@acme.com"))).text(); + { + const kinds = responsesDeltaKinds(txt); + ok("stream/responses/refusal: deltas present", kinds.length > 0, String(kinds.length)); + ok("stream/responses/refusal: every re-emitted delta is a refusal delta", + kinds.every((k) => k[1] === "response.refusal.delta"), JSON.stringify(kinds.map((k) => k[1]))); + ok("stream/responses/refusal: event: line matches the data type", + kinds.every((k) => k[0] === k[1]), JSON.stringify(kinds.map((k) => [k[0], k[1]]))); + ok("stream/responses/refusal: no output_text delta leaked", + !txt.includes("event: response.output_text.delta") && !txt.includes('"response.output_text.delta"')); + const refusal = kinds.map((k) => k[2] ?? "").join(""); + ok("stream/responses/refusal: refusal text restored across frame split", + refusal.includes("john@acme.com") && !refusal.includes(" Date: Tue, 22 Sep 2026 07:17:53 -0400 Subject: [PATCH 4/7] fix: redact prior assistant turns and stop replacement metacharacters reframing A stateless Responses client appends the previous reply's output to the next request's input. Those parts are output_text and refusal, which the request walk did not recognise, so a restored email went back upstream raw on the second turn. Both part kinds are now redacted, with a multi-turn regression that checks the captured upstream body. reframe replaced the data: line with a string, so $&, $` and $' in model text were replacement metacharacters and spliced the original line into the payload. It takes a function now. --- _test_proxy.mjs | 22 ++++++++++++++++++++++ src/redact/apply.ts | 10 +++++++++- src/streaming.ts | 5 ++++- 3 files changed, 35 insertions(+), 2 deletions(-) diff --git a/_test_proxy.mjs b/_test_proxy.mjs index 453612f..835395c 100644 --- a/_test_proxy.mjs +++ b/_test_proxy.mjs @@ -122,6 +122,28 @@ const setTenant2 = setTenantOn(BASE2); } ok("responses/items: X-Redacted counts every field", Number(res.headers.get("x-redacted")) >= 5, res.headers.get("x-redacted")); + // ---- reversible (Responses): a prior assistant turn fed back as input ---- + // The reply cordon restores carries real values, and a stateless client appends it to + // the next request's input. Those parts are output_text/refusal, not input_text. + await reset(); + res = await post("/v1/responses", rBody([ + { role: "user", content: [{ type: "input_text", text: "who do I contact" }] }, + { + type: "message", role: "assistant", + content: [ + { type: "output_text", text: "Contact john@acme.com about card 4012888888881881", annotations: [] }, + { type: "refusal", refusal: "I cannot share ops@acme.com" }, + ], + }, + { role: "user", content: [{ type: "input_text", text: "thanks" }] }, + ])); + sent = JSON.stringify(await calls()); + ok("responses/history: assistant output_text never reaches upstream raw", + !sent.includes("john@acme.com") && !sent.includes("4012888888881881"), sent.slice(0, 200)); + ok("responses/history: assistant refusal never reaches upstream raw", !sent.includes("ops@acme.com")); + ok("responses/history: the assistant turn is placeholdered", //.test(sent)); + ok("responses/history: X-Redacted counts the assistant turn", Number(res.headers.get("x-redacted")) >= 3, res.headers.get("x-redacted")); + // ---- strip / off (Responses) ---- await reset(); res = await post("/v1/responses", rBody(PII), { "x-redact-mode": "strip" }); diff --git a/src/redact/apply.ts b/src/redact/apply.ts index a5af491..38bcd49 100644 --- a/src/redact/apply.ts +++ b/src/redact/apply.ts @@ -91,8 +91,16 @@ function requestTextSlots( const part = c[i]; if (typeof part === "string") { slots.push(slot(c, i)); continue; } // a bare-string content element if (!part || typeof part !== "object") continue; - if ((part.type === "text" || part.type === "input_text") && typeof part.text === "string") { + if ( + (part.type === "text" || part.type === "input_text" || part.type === "output_text") && + typeof part.text === "string" + ) { + // output_text is an ASSISTANT part: a stateless Responses conversation appends + // the previous reply's output to the next request's input, and that reply left + // here with its real values restored. slots.push(slot(part, "text")); + } else if (part.type === "refusal" && typeof part.refusal === "string") { + slots.push(slot(part, "refusal")); } else if (part.type === "tool_result") { // Anthropic tool_result content can itself be a string or block array. if (typeof part.content === "string") slots.push(slot(part, "content")); diff --git a/src/streaming.ts b/src/streaming.ts index b8e8a1d..58a6711 100644 --- a/src/streaming.ts +++ b/src/streaming.ts @@ -25,7 +25,10 @@ export async function pipeUpstream(up: Response, res: HttpRes): Promise { /** The same SSE frame with its `data:` payload replaced by `j` (event: lines kept). */ function reframe(frame: string, j: unknown): string { - return frame.replace(/^data:.*$/m, `data: ${JSON.stringify(j)}`); + // A function replacement, never a string: `$&`, `$\`` and `$'` inside model text are + // replacement metacharacters to String.replace and would splice the original line + // into the payload unescaped. + return frame.replace(/^data:.*$/m, () => `data: ${JSON.stringify(j)}`); } /** From 43a62bfcf35d25c955ba664ec0dd5ff90e25214e Mon Sep 17 00:00:00 2001 From: askalf <263217947+askalf@users.noreply.github.com> Date: Tue, 22 Sep 2026 07:33:01 -0400 Subject: [PATCH 5/7] fix: give each Responses content part its own re-identifier A Responses turn with more than one output item interleaves the deltas of its parts on the wire, addressed by item_id/output_index/content_index. One buffer for the whole stream spliced the second part's text into a placeholder the first had half-written and re-emitted it under the wrong address. Each part now keeps its own re-identifier and addressing, flushed on that part's own done frame; the frames that close the whole response flush whatever is still open. The stub streams two parts at once and the regression reassembles by address: without the fix the first part reads 'email AIL_101073_1> about card REDIT_CARD_101073_1>' and the second carries the lost ' res.write(`event: ${type}\ndata: ${JSON.stringify({ type, sequence_number: seq++, ...d })}\n\n`); + const a = { item_id: "msg_stub" + n + "a", output_index: 0, content_index: 0 }; + const b = { item_id: "msg_stub" + n + "b", output_index: 1, content_index: 0 }; + const aText = text.replace("INTERLEAVE ", ""); + const bText = "the second part says nothing private"; + const shell = (status) => ({ id: "resp_stub" + n, object: "response", model: "gpt-4o-mini", status, output: [] }); + const item = (addr, t, status) => ({ id: addr.item_id, type: "message", role: "assistant", status, content: status === "completed" ? [{ type: "output_text", text: t, annotations: [] }] : [] }); + f("response.created", { response: shell("in_progress") }); + for (const [addr, t] of [[a, aText], [b, bText]]) { + f("response.output_item.added", { output_index: addr.output_index, item: item(addr, t, "in_progress") }); + f("response.content_part.added", { ...addr, part: { type: "output_text", text: "", annotations: [] } }); + } + const ac = chunk3(aText), bc = chunk3(bText); + for (let i = 0; i < Math.max(ac.length, bc.length); i++) { + if (i < ac.length) f("response.output_text.delta", { ...a, delta: ac[i] }); + if (i < bc.length) f("response.output_text.delta", { ...b, delta: bc[i] }); + } + for (const [addr, t] of [[a, aText], [b, bText]]) { + f("response.output_text.done", { ...addr, text: t }); + f("response.content_part.done", { ...addr, part: { type: "output_text", text: t, annotations: [] } }); + f("response.output_item.done", { output_index: addr.output_index, item: item(addr, t, "completed") }); + } + f("response.completed", { response: { ...shell("completed"), output: [item(a, aText, "completed"), item(b, bText, "completed")], usage: { input_tokens: 10, output_tokens: 5, total_tokens: 15 } } }); + res.end(); +} + // A REFUSED Responses turn streams the same frame shape under refusal-flavoured // event types (`response.refusal.delta` / `.done`, a `refusal` content part). A // client consuming refusals reads only those, so cordon must re-emit a restored @@ -163,7 +195,9 @@ http return body.stream ? streamOpenAI(res, text) : json(res, openaiBody(n, text)); if (req.url.includes("/responses")) return body.stream - ? text.includes("FORCE_REFUSAL") ? streamResponsesRefusal(res, text, n) : streamResponses(res, text, n) + ? text.includes("FORCE_REFUSAL") ? streamResponsesRefusal(res, text, n) + : text.includes("INTERLEAVE") ? streamResponsesInterleaved(res, text, n) + : streamResponses(res, text, n) : json(res, responsesBody(n, text)); if (req.url.includes("/messages")) return body.stream ? streamAnthropic(res, text, body) : json(res, anthropicBody(n, text)); diff --git a/_test_stream.mjs b/_test_stream.mjs index 32cbaa3..1dbc4df 100644 --- a/_test_stream.mjs +++ b/_test_stream.mjs @@ -171,6 +171,30 @@ function validateSSE(sse) { sent = JSON.stringify(await calls()); ok("stream/responses: upstream saw placeholder not raw", //.test(sent) && !sent.includes("john@acme.com")); + // ---- reversible streaming (Responses, two interleaved content parts) ---- + // Responses addresses a delta by item_id/output_index/content_index and parts can + // interleave. With one buffer for the whole stream, part B's text lands inside a + // placeholder part A had half-written and is re-emitted under B's address. + await reset(); + txt = await (await post("/v1/responses", rBody("INTERLEAVE " + PII))).text(); + { + const byAddr = new Map(); + for (const f of txt.split("\n\n")) { + const d = f.split("\n").find((l) => l.startsWith("data:")); + if (!d) continue; + let j; try { j = JSON.parse(d.slice(5).trim()); } catch { continue; } + if (j?.type !== "response.output_text.delta") continue; + const k = `${j.item_id}/${j.output_index}/${j.content_index}`; + byAddr.set(k, (byAddr.get(k) ?? "") + (j.delta ?? "")); + } + ok("stream/responses/interleave: both parts present", byAddr.size === 2, JSON.stringify([...byAddr.keys()])); + const joined = [...byAddr.values()]; + ok("stream/responses/interleave: the PII part is restored whole", + joined.some((t) => t.includes("john@acme.com")) && !joined.join("").includes(" t.includes("second part") && !t.includes("john@acme.com")), JSON.stringify(joined)); + } + // ---- reversible streaming (Responses REFUSAL) ---- // A refusal carries restorable text too, but under its own event type. Regression: // the re-emit path shared frameFromText with output text and hardcoded the diff --git a/src/streaming.ts b/src/streaming.ts index 58a6711..59ad13a 100644 --- a/src/streaming.ts +++ b/src/streaming.ts @@ -59,11 +59,42 @@ export async function captureAndReidentify( // Responses addresses a delta by item_id / output_index / content_index; a re-emitted // delta carries the addressing of the frame it stands in for. let emitCtx: Record = {}; - const emitText = (chunk: string, index = emitIndex) => { - if (chunk) res.write(adapter.frameFromText(chunk, index, emitCtx)); + const emitText = (chunk: string, index = emitIndex, ctx = emitCtx) => { + if (chunk) res.write(adapter.frameFromText(chunk, index, ctx)); }; const flushTail = () => emitText(reider.end(), emitIndex); + // A Responses stream can carry several content parts at once (two output items, or two + // content indices of one item) and their deltas interleave. One buffer for all of them + // would splice part B's text into a placeholder part A had half-written, and re-emit it + // under B's address. Each part gets its own re-identifier and its own addressing, and + // is flushed and dropped when that part's own done frame arrives. + const parts = new Map }>(); + const partKey = (j: any) => `${j?.item_id ?? ""}/${j?.output_index ?? 0}/${j?.content_index ?? 0}`; + const partFor = (j: any, ctx: Record) => { + const key = partKey(j); + const found = parts.get(key); + if (found) { found.ctx = ctx; return found; } + const made = { reider: new StreamReidentifier(vault), ctx }; + parts.set(key, made); + return made; + }; + /** Flush one part's held tail under its own addressing and forget it. */ + const flushPart = (j: any) => { + const key = partKey(j); + const part = parts.get(key); + if (!part) return; + emitText(part.reider.end(), 0, part.ctx); + parts.delete(key); + }; + /** Flush every part still open, for the frames that close the whole response. */ + const flushAllParts = () => { + for (const [key, part] of parts) { + emitText(part.reider.end(), 0, part.ctx); + parts.delete(key); + } + }; + const handleFrame = (frame: string) => { const dataLine = frame.split("\n").find((l) => l.startsWith("data:")); const data = dataLine ? dataLine.slice(5).trim() : ""; @@ -72,7 +103,8 @@ export async function captureAndReidentify( return; } if (data === "[DONE]") { - flushTail(); + if (dialect === "responses") flushAllParts(); + else flushTail(); res.write(frame); return; } @@ -116,34 +148,34 @@ export async function captureAndReidentify( // text would break a client that reads refusal events; the possibly-forming // tail stays held in the re-identifier. const { delta: _d, ...ctx } = j; - emitCtx = ctx; - emitText(reider.push(j.delta ?? "")); + const part = partFor(j, ctx); + emitText(part.reider.push(j.delta ?? ""), 0, part.ctx); return; } // Every frame below carries the text in full, so the held tail is flushed first // (a placeholder split across deltas is then already restored on the client) and // the frame's own text is restored before it passes. if (type === "response.output_text.done" || type === "response.refusal.done") { - flushTail(); + flushPart(j); if (typeof j.text === "string") j.text = restore(j.text, vault); if (typeof j.refusal === "string") j.refusal = restore(j.refusal, vault); res.write(reframe(frame, j)); return; } if (type === "response.content_part.done" && j.part && typeof j.part === "object") { - flushTail(); + flushPart(j); if (typeof j.part.text === "string") j.part.text = restore(j.part.text, vault); if (typeof j.part.refusal === "string") j.part.refusal = restore(j.part.refusal, vault); res.write(reframe(frame, j)); return; } if (type === "response.output_item.done" && j.item && typeof j.item === "object") { - flushTail(); + flushAllParts(); res.write(reframe(frame, { ...j, item: reidentifyBody({ output: [j.item] }, provider, vault, dialect).output[0] })); return; } if (type === "response.completed" || type === "response.incomplete" || type === "response.failed") { - flushTail(); + flushAllParts(); if (j.response && typeof j.response === "object") res.write(reframe(frame, { ...j, response: reidentifyBody(j.response, provider, vault, dialect) })); else res.write(frame); From 29195f94ce458952f2d7ad22753cdcdd0ef2e5a9 Mon Sep 17 00:00:00 2001 From: askalf <263217947+askalf@users.noreply.github.com> Date: Tue, 22 Sep 2026 07:38:20 -0400 Subject: [PATCH 6/7] fix: closing one output item flushes only that item's parts response.output_item.done ended the re-identifier for every open part, not just the item named by the frame. An item still mid-placeholder had its fragment emitted as-is, and the suffix that followed arrived on a fresh buffer, so the client saw the placeholder's tail in its text. Only the parts whose addressing matches the closing item are flushed; response.completed, incomplete, failed and [DONE] still flush the rest. The regression now cuts the second part's deltas inside its own placeholder and closes the first item at that moment. On the previous head the second part reads 'john@acme.comIL_F9D5A1_1> about card'. --- _stub-upstream.mjs | 31 ++++++++++++++++++++----------- _test_stream.mjs | 7 ++++++- src/streaming.ts | 19 ++++++++++++++++--- 3 files changed, 42 insertions(+), 15 deletions(-) diff --git a/_stub-upstream.mjs b/_stub-upstream.mjs index d00cdc5..7c6fbbd 100644 --- a/_stub-upstream.mjs +++ b/_stub-upstream.mjs @@ -98,7 +98,9 @@ function streamResponsesInterleaved(res, text, n) { const a = { item_id: "msg_stub" + n + "a", output_index: 0, content_index: 0 }; const b = { item_id: "msg_stub" + n + "b", output_index: 1, content_index: 0 }; const aText = text.replace("INTERLEAVE ", ""); - const bText = "the second part says nothing private"; + // B's text is the SAME de-identified text the stub received, so it carries real + // placeholders and the cut below can land inside one. + const bText = "second part repeats " + aText; const shell = (status) => ({ id: "resp_stub" + n, object: "response", model: "gpt-4o-mini", status, output: [] }); const item = (addr, t, status) => ({ id: addr.item_id, type: "message", role: "assistant", status, content: status === "completed" ? [{ type: "output_text", text: t, annotations: [] }] : [] }); f("response.created", { response: shell("in_progress") }); @@ -106,16 +108,23 @@ function streamResponsesInterleaved(res, text, n) { f("response.output_item.added", { output_index: addr.output_index, item: item(addr, t, "in_progress") }); f("response.content_part.added", { ...addr, part: { type: "output_text", text: "", annotations: [] } }); } - const ac = chunk3(aText), bc = chunk3(bText); - for (let i = 0; i < Math.max(ac.length, bc.length); i++) { - if (i < ac.length) f("response.output_text.delta", { ...a, delta: ac[i] }); - if (i < bc.length) f("response.output_text.delta", { ...b, delta: bc[i] }); - } - for (const [addr, t] of [[a, aText], [b, bText]]) { - f("response.output_text.done", { ...addr, text: t }); - f("response.content_part.done", { ...addr, part: { type: "output_text", text: t, annotations: [] } }); - f("response.output_item.done", { output_index: addr.output_index, item: item(addr, t, "completed") }); - } + // B goes first and stops PART-WAY THROUGH ITS PLACEHOLDER, then A runs to + // completion INCLUDING its output_item.done, and only then does B finish. Closing A + // must not end B's re-identifier while B still holds a half-written placeholder. + // The cut lands INSIDE B's placeholder: the text the stub echoes is already + // de-identified, so `<` is the placeholder's first character and B stops four + // characters in, holding a fragment no re-identifier can resolve yet. + const ac = chunk3(aText); + const cut = bText.indexOf("<") >= 0 ? bText.indexOf("<") + 4 : Math.ceil(bText.length / 2); + for (const c of chunk3(bText.slice(0, cut))) f("response.output_text.delta", { ...b, delta: c }); + for (const c of ac) f("response.output_text.delta", { ...a, delta: c }); + f("response.output_text.done", { ...a, text: aText }); + f("response.content_part.done", { ...a, part: { type: "output_text", text: aText, annotations: [] } }); + f("response.output_item.done", { output_index: a.output_index, item: item(a, aText, "completed") }); + for (const c of chunk3(bText.slice(cut))) f("response.output_text.delta", { ...b, delta: c }); + f("response.output_text.done", { ...b, text: bText }); + f("response.content_part.done", { ...b, part: { type: "output_text", text: bText, annotations: [] } }); + f("response.output_item.done", { output_index: b.output_index, item: item(b, bText, "completed") }); f("response.completed", { response: { ...shell("completed"), output: [item(a, aText, "completed"), item(b, bText, "completed")], usage: { input_tokens: 10, output_tokens: 5, total_tokens: 15 } } }); res.end(); } diff --git a/_test_stream.mjs b/_test_stream.mjs index 1dbc4df..7ec2dfb 100644 --- a/_test_stream.mjs +++ b/_test_stream.mjs @@ -192,7 +192,12 @@ function validateSSE(sse) { ok("stream/responses/interleave: the PII part is restored whole", joined.some((t) => t.includes("john@acme.com")) && !joined.join("").includes(" t.includes("second part") && !t.includes("john@acme.com")), JSON.stringify(joined)); + joined.some((t) => t.startsWith("second part repeats")), JSON.stringify(joined)); + // The stub closes item A (output_text.done, content_part.done, output_item.done) + // while item B is holding a half-written placeholder. Closing A must not end B's + // re-identifier, or B's email resolves early and its suffix arrives on its own. + ok("stream/responses/interleave: a part still open survives another item closing", + joined.every((t) => t.includes("john@acme.com") && !//.test(t)), JSON.stringify(joined)); } // ---- reversible streaming (Responses REFUSAL) ---- diff --git a/src/streaming.ts b/src/streaming.ts index 59ad13a..3797d0f 100644 --- a/src/streaming.ts +++ b/src/streaming.ts @@ -87,13 +87,26 @@ export async function captureAndReidentify( emitText(part.reider.end(), 0, part.ctx); parts.delete(key); }; - /** Flush every part still open, for the frames that close the whole response. */ - const flushAllParts = () => { + /** Flush the parts matching `pick`, each under its own addressing, and forget them. */ + const flushParts = (pick: (ctx: Record) => boolean) => { for (const [key, part] of parts) { + if (!pick(part.ctx)) continue; emitText(part.reider.end(), 0, part.ctx); parts.delete(key); } }; + /** Every part still open: only for the frames that close the whole response. */ + const flushAllParts = () => flushParts(() => true); + /** + * The parts of ONE output item, on its `output_item.done`. Closing item A must not + * end item B's re-identifier: B may be holding a half-written placeholder, and + * ending it early emits the resolved value and then B's own suffix separately. + */ + const flushOutputItem = (j: any) => { + const id = j?.item?.id; + const idx = j?.output_index; + flushParts((ctx) => (id !== undefined && ctx["item_id"] === id) || (idx !== undefined && ctx["output_index"] === idx)); + }; const handleFrame = (frame: string) => { const dataLine = frame.split("\n").find((l) => l.startsWith("data:")); @@ -170,7 +183,7 @@ export async function captureAndReidentify( return; } if (type === "response.output_item.done" && j.item && typeof j.item === "object") { - flushAllParts(); + flushOutputItem(j); res.write(reframe(frame, { ...j, item: reidentifyBody({ output: [j.item] }, provider, vault, dialect).output[0] })); return; } From a1e3e480fad1d621d1afe7b03e8e53525464adce Mon Sep 17 00:00:00 2001 From: askalf <263217947+askalf@users.noreply.github.com> Date: Tue, 22 Sep 2026 07:45:33 -0400 Subject: [PATCH 7/7] fix: flush the held Responses text on every path that ends a stream The per-part buffers were flushed for the terminal frames only. When the upstream stopped without one, the reader's end path and the catch path flushed the chat-dialect buffer, which for a Responses stream has received nothing, so the client's text simply stopped where the last resolved delta ended. One flushHeld now serves all three. The regression stops a Responses stream with a placeholder still held: on the previous head the client receives 'email john@acme.com about card ' and the card never arrives. --- _stub-upstream.mjs | 21 ++++++++++++++++++++- _test_stream.mjs | 22 ++++++++++++++++++++++ src/streaming.ts | 14 ++++++++++---- 3 files changed, 52 insertions(+), 5 deletions(-) diff --git a/_stub-upstream.mjs b/_stub-upstream.mjs index 7c6fbbd..985410d 100644 --- a/_stub-upstream.mjs +++ b/_stub-upstream.mjs @@ -129,6 +129,24 @@ function streamResponsesInterleaved(res, text, n) { res.end(); } +// A Responses stream that just STOPS: deltas up to a point four characters inside a +// placeholder, then the socket ends with no output_text.done, no response.completed +// and no [DONE]. Whatever the re-identifier is holding has to reach the client anyway. +function streamResponsesTruncated(res, text, n) { + res.setHeader("content-type", "text/event-stream"); + let seq = 0; + const f = (type, d) => res.write(`event: ${type}\ndata: ${JSON.stringify({ type, sequence_number: seq++, ...d })}\n\n`); + const addr = { item_id: "msg_stub" + n, output_index: 0, content_index: 0 }; + const body = text.replace("TRUNCATE ", ""); + const marks = [...body.matchAll(/ m.index); + const cut = marks.length > 1 ? marks[1] + 4 : Math.ceil(body.length / 2); + f("response.created", { response: { id: "resp_stub" + n, object: "response", model: "gpt-4o-mini", status: "in_progress", output: [] } }); + f("response.output_item.added", { output_index: 0, item: { id: addr.item_id, type: "message", role: "assistant", status: "in_progress", content: [] } }); + f("response.content_part.added", { ...addr, part: { type: "output_text", text: "", annotations: [] } }); + for (const c of chunk3(body.slice(0, cut))) f("response.output_text.delta", { ...addr, delta: c }); + res.end(); +} + // A REFUSED Responses turn streams the same frame shape under refusal-flavoured // event types (`response.refusal.delta` / `.done`, a `refusal` content part). A // client consuming refusals reads only those, so cordon must re-emit a restored @@ -204,7 +222,8 @@ http return body.stream ? streamOpenAI(res, text) : json(res, openaiBody(n, text)); if (req.url.includes("/responses")) return body.stream - ? text.includes("FORCE_REFUSAL") ? streamResponsesRefusal(res, text, n) + ? text.includes("TRUNCATE") ? streamResponsesTruncated(res, text, n) + : text.includes("FORCE_REFUSAL") ? streamResponsesRefusal(res, text, n) : text.includes("INTERLEAVE") ? streamResponsesInterleaved(res, text, n) : streamResponses(res, text, n) : json(res, responsesBody(n, text)); diff --git a/_test_stream.mjs b/_test_stream.mjs index 7ec2dfb..3918f35 100644 --- a/_test_stream.mjs +++ b/_test_stream.mjs @@ -200,6 +200,28 @@ function validateSSE(sse) { joined.every((t) => t.includes("john@acme.com") && !//.test(t)), JSON.stringify(joined)); } + // ---- Responses stream that ends with no terminal frame ---- + // The upstream stops after a delta that leaves a placeholder half-written: the + // per-part buffers are the only place that text lives, so the end-of-reader path + // has to flush THEM. Flushing the chat-dialect buffer there emits nothing and the + // client silently loses the tail. + await reset(); + txt = await (await post("/v1/responses", rBody("TRUNCATE " + PII))).text(); + { + let got = ""; + for (const f of txt.split("\n\n")) { + const d = f.split("\n").find((l) => l.startsWith("data:")); + if (!d) continue; + let j; try { j = JSON.parse(d.slice(5).trim()); } catch { continue; } + if (j?.type === "response.output_text.delta") got += j.delta ?? ""; + } + ok("stream/responses/truncated: the resolved text arrives", got.includes("john@acme.com"), JSON.stringify(got)); + // The card's placeholder was still held when the socket closed. Flushing the + // per-part buffer restores it; flushing the chat buffer instead emits nothing and + // the client's text stops at "about card ". + ok("stream/responses/truncated: the held tail is not dropped", got.includes("4012888888881881"), JSON.stringify(got)); + } + // ---- reversible streaming (Responses REFUSAL) ---- // A refusal carries restorable text too, but under its own event type. Regression: // the re-emit path shared frameFromText with output text and hardcoded the diff --git a/src/streaming.ts b/src/streaming.ts index 3797d0f..b57d371 100644 --- a/src/streaming.ts +++ b/src/streaming.ts @@ -97,6 +97,13 @@ export async function captureAndReidentify( }; /** Every part still open: only for the frames that close the whole response. */ const flushAllParts = () => flushParts(() => true); + /** + * Whatever this dialect is holding, for the paths that end a stream rather than + * close it: the terminal frame, the reader throwing, and an upstream that simply + * stops. Responses keeps its held text in the per-part map, so calling the single + * `flushTail` buffer there emits nothing and the client loses the tail. + */ + const flushHeld = () => (dialect === "responses" ? flushAllParts() : flushTail()); /** * The parts of ONE output item, on its `output_item.done`. Closing item A must not * end item B's re-identifier: B may be holding a half-written placeholder, and @@ -116,8 +123,7 @@ export async function captureAndReidentify( return; } if (data === "[DONE]") { - if (dialect === "responses") flushAllParts(); - else flushTail(); + flushHeld(); res.write(frame); return; } @@ -233,11 +239,11 @@ export async function captureAndReidentify( } if (buf.length) handleFrame(buf); // trailing frame without a terminating blank line } catch { - flushTail(); // best-effort restore of whatever was held + flushHeld(); // best-effort restore of whatever was held res.end(); return; } - flushTail(); // safety: flush if the stream ended without an explicit close frame + flushHeld(); // safety: flush if the stream ended without an explicit close frame res.end(); }