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..985410d 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,104 @@ 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(); +} +// Two content parts of ONE response, streaming at the same time: the deltas of +// output_index 0 and 1 alternate on the wire, which is what a Responses turn with +// more than one output item looks like. Only part 0 carries PII. +function streamResponsesInterleaved(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 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 ", ""); + // 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") }); + 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: [] } }); + } + // 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(); +} + +// 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 +// 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`); @@ -109,6 +220,13 @@ 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 + ? 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)); 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..835395c 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,85 @@ 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")); + + // ---- 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" }); + 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 +233,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..3918f35 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,54 @@ 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; +} + +/** 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(); + 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 +154,102 @@ 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")); + + // ---- 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.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)); + } + + // ---- 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 + // 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(" { 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..eea8b7a 100644 --- a/src/providers.ts +++ b/src/providers.ts @@ -1,19 +1,22 @@ 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) 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 }>; } @@ -74,7 +77,51 @@ export const anthropic: ProviderAdapter = { }, }; -export const adapterFor = (p: Provider) => (p === "openai" ? openai : anthropic); +// ----------------------------- 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 (RESPONSES_TEXT_DELTAS.has(j.type)) 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, + // 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 }> = []; + 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 +155,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 +233,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..38bcd49 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,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" && 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")); @@ -100,6 +126,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 +169,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 +206,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..b57d371 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,14 @@ 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 { + // 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)}`); +} + /** * 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 +45,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,11 +56,65 @@ 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; - const emitText = (chunk: string, index = emitIndex) => { - if (chunk) res.write(adapter.frameFromText(chunk, index)); + // 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, 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 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); + /** + * 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 + * 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:")); const data = dataLine ? dataLine.slice(5).trim() : ""; @@ -60,7 +123,7 @@ export async function captureAndReidentify( return; } if (data === "[DONE]") { - flushTail(); + flushHeld(); res.write(frame); return; } @@ -92,6 +155,55 @@ 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 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; + 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") { + 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") { + 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") { + flushOutputItem(j); + 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") { + flushAllParts(); + 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 { @@ -127,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(); } 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. */