From 66f791bc9473c9d11baf8c76e22f52a23003aaa3 Mon Sep 17 00:00:00 2001 From: Josh Mabry Date: Thu, 27 Aug 2026 13:32:33 -0700 Subject: [PATCH 1/4] feat(codex): capture encrypted reasoning and replay it only to its issuer MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Closes the "contained, not delivered" open item #3199 left on ADR 0097. Cross-turn reasoning continuity on `openai-codex` was off because the blob that makes it possible was never captured; this wires the capability and the guard that makes it safe on a shared thread. Capture. langchain-openai's streaming Responses path has no response.output_item.done branch for reasoning (it has one for `compaction`, which carries the same kind of blob), and the terminal response.completed event keeps only parsed/usage/response_metadata — so encrypted_content is visible in exactly one event, which the converter drops. codex_client re-emits that event as a content-block delta that merges onto the reasoning block already in flight, by index. The wrapper sits on the shared module-level converter (there is no instance seam) but is inert unless a contextvar this module's client sets is present, so every other ChatOpenAI in the process goes through the original path unchanged. output_version flipped to responses/v1. "v0" collapses a turn's reasoning into ONE additional_kwargs slot — later items overwrite earlier ones, and streamed fragments of two different items merge into each other — so it structurally cannot carry per-item blobs. The block format keeps each item separate and in order, and langchain replays it that way. The rendering half of the v0 pin was already paid off (every answer site reads AIMessage.text); text_of now skips reasoning blocks rather than writing a _[reasoning]_ placeholder into exports, session memory and chat bundles, which is what ADR 0021 asks for anyway. PROTOAGENT_CODEX_OUTPUT_VERSION=v0 is the escape hatch. Issuer stamping. encrypted_content is sealed to the endpoint AND account that minted it. Each captured item carries a truncated digest of (base_url, account_id) — so a checkpoint never stores a raw account id — and replay drops items stamped with a different issuer. Unstamped items still replay. This is what makes per-slot providers, per-tab model override and the fallback chain safe on one thread; without it the #3199 recovery middleware would fire routinely instead of never. Not verified against a live ChatGPT subscription. If the backend objects, the #3199 recovery valve strips the replay state and retries, so the thread degrades to stateless continuity rather than breaking. Co-Authored-By: Claude Opus 5 (1M context) --- ...097-native-oauth-subscription-providers.md | 47 ++- docs/reference/configuration.md | 9 + docs/reference/environment-variables.md | 7 + graph/message_blocks.py | 8 + graph/providers/codex_client.py | 245 ++++++++++---- graph/providers/openai_codex.py | 38 ++- tests/test_codex_reasoning_capture.py | 317 ++++++++++++++++++ 7 files changed, 595 insertions(+), 76 deletions(-) create mode 100644 tests/test_codex_reasoning_capture.py diff --git a/docs/adr/0097-native-oauth-subscription-providers.md b/docs/adr/0097-native-oauth-subscription-providers.md index 4e44744ba..fb2bf77fd 100644 --- a/docs/adr/0097-native-oauth-subscription-providers.md +++ b/docs/adr/0097-native-oauth-subscription-providers.md @@ -210,13 +210,50 @@ Hermes's Codex adapter (`agent/codex_responses_adapter.py`) reached the same rul independently — including the id strip and a session-wide replay kill switch — and its `_issuer_kind` stamp is the model for the cross-issuer filter listed below. +## Encrypted-reasoning replay, delivered (2026-08-27, #3199 follow-up) + +#3199 contained the damage — never send an item the backend can't verify. This wires the +capability the containment was standing in for, and closes the "contained, not delivered" +open item. + +**Capture.** langchain-openai's streaming Responses path has no `response.output_item.done` +branch for reasoning (it has one for `compaction`, which carries the same kind of blob), and +the terminal `response.completed` event keeps only `parsed`/usage/`response_metadata`. So the +blob is visible in exactly one event, which the converter drops. `codex_client` +`_install_reasoning_capture` re-emits that event as a content-block delta that merges onto the +reasoning block already in flight, by `index`. The wrapper sits on the shared module-level +converter — there is no instance seam — but is **inert unless a contextvar this module's +client sets is present**, so every other `ChatOpenAI` in the process is untouched. + +**`output_version` flipped to `responses/v1`.** `v0` collapses a turn's reasoning into ONE +`additional_kwargs` slot: later items overwrite earlier ones, and streamed fragments of two +different items merge into each other — so it structurally cannot carry per-item blobs. The +block format keeps each item separate and in order, and langchain replays it that way. The +rendering half of the v0 pin was already paid off (every answer site reads `AIMessage.text`, +which yields text blocks only); `text_of` now skips reasoning blocks outright rather than +writing a `_[reasoning]_` placeholder into exports/session memory/chat bundles, which is what +ADR 0021 asks for anyway. `PROTOAGENT_CODEX_OUTPUT_VERSION=v0` is the escape hatch. + +**Issuer stamping.** `encrypted_content` is sealed to the endpoint *and account* that minted +it. Each captured item carries `issuer_fingerprint(base_url, account_id)` — a truncated +digest, so a checkpoint never stores a raw account id — and replay drops items stamped with a +different issuer. Unstamped items (checkpointed before this) still replay. This is the guard +that makes per-slot providers, per-tab model override and the fallback chain safe on a shared +thread; without it, the recovery middleware would be firing routinely instead of never. + +**Not verified live.** The wire shape is tested end to end against the real converter, the +real merge and the real payload builder, but no turn has been driven against a real ChatGPT +subscription with this on. If the backend objects, `CodexReasoningReplayRecoveryMiddleware` +(#3199) strips the replay state and retries — the thread degrades to stateless continuity +rather than breaking, which is exactly why that half shipped first. + ## Open items -- **Encrypted-reasoning replay is contained, not delivered.** Cross-turn reasoning - continuity on `openai-codex` is OFF: capturing the blob needs an `output_item.done` - handler for reasoning items that langchain-openai does not have (worth an upstream - issue). Once captured, replayed items should carry an issuer stamp (endpoint + account) - and be filtered when the current endpoint differs — Hermes's `_classify_responses_issuer`. +- **Encrypted-reasoning replay is unverified against a live subscription.** Capture, + issuer stamping and replay are wired (above) and covered by wire-shape tests, but no turn + has been driven against a real ChatGPT account with it on. Worth an upstream issue too: + langchain-openai should handle `output_item.done` for reasoning items the way it already + does for `compaction`, which would let protoAgent drop its converter wrapper. - **Claude end-to-end still unproven on a real subscription** — the sign-in URL + PKCE + refresh are unit-tested and the flow runs, but no Pro/Max approval has been driven here yet (tool loop, streaming, `cache_control`). diff --git a/docs/reference/configuration.md b/docs/reference/configuration.md index 492e41377..00c58bbc6 100644 --- a/docs/reference/configuration.md +++ b/docs/reference/configuration.md @@ -499,6 +499,15 @@ Any slot that takes a model name — `routing.fallback_models`, `routing.aux_mod Hold a gateway key and both subscriptions and you can mix all of them at once — Claude for review, Codex for code, the gateway for cheap bulk work — whatever the main brain runs on. The qualified form is the one to reach for when two providers could plausibly serve the same model id. +::: tip Mixing providers mid-conversation costs reasoning continuity, not correctness +On `openai-codex`, the model's reasoning is threaded across turns as an encrypted blob that +is **sealed to the endpoint and account that minted it** — a blob replayed anywhere else is a +hard `400`. protoAgent stamps each captured item with its issuer and silently drops the ones +the current endpoint can't decrypt, so switching a chat's model mid-thread (or re-signing-in +under a different ChatGPT account) just restarts reasoning continuity from that point. The +conversation itself is unaffected. +::: + ```yaml model: provider: anthropic-oauth # main brain on your Claude subscription diff --git a/docs/reference/environment-variables.md b/docs/reference/environment-variables.md index a0e9cffdc..8e6862ec6 100644 --- a/docs/reference/environment-variables.md +++ b/docs/reference/environment-variables.md @@ -32,6 +32,13 @@ Every env var the template reads at runtime. | `PROTOAGENT_MODEL` | (unset) | Overrides `model.name` on every config load — used by `evals/sweep.py` to run one agent against many models without editing YAML. | | `PROTOAGENT_INSTANCE` | (unset) | Opt-in data-scoping key (ADR 0004): namespaces the knowledge/notes/tasks/checkpoint stores so several agents share a backend without colliding. Seeded from `instance.id` in config. | +## Native OAuth subscription providers (ADR 0097) + +| Variable | Default | What | +|---|---|---| +| `PROTOAGENT_CODEX_BASE_URL` | `https://chatgpt.com/backend-api/codex` | Endpoint for `model.provider: openai-codex`. Changing it changes the **issuer** a captured reasoning blob is sealed to, so items minted against the old endpoint stop being replayed (by design — the new one can't decrypt them). | +| `PROTOAGENT_CODEX_OUTPUT_VERSION` | `responses/v1` | Escape hatch for the `openai-codex` content shape. `responses/v1` keeps each reasoning item as its own content block, which is what makes cross-turn encrypted-reasoning replay possible. Set to `v0` to fall back to the legacy string-content shape — the turn still works, but reasoning continuity across turns is off. | + ## Deployment / UI tier (ADR 0010) | Variable | Default | What | diff --git a/graph/message_blocks.py b/graph/message_blocks.py index 03254d34c..cf9906244 100644 --- a/graph/message_blocks.py +++ b/graph/message_blocks.py @@ -48,6 +48,14 @@ def text_of(message) -> str: elif isinstance(block, dict): if block.get("type") == "text" and block.get("text"): parts.append(str(block["text"])) + elif block.get("type") == "reasoning": + # Skipped outright, not placeholdered. ADR 0021 says reasoning is + # never persisted, and every caller here WRITES what it returns + # (exports, session memory, chat bundles). A `_[reasoning]_` marker + # would be noise in all three — and on the Responses providers + # (openai-codex, ADR 0097) reasoning is a block on EVERY assistant + # turn, so it would be noise on every line. + continue elif block.get("type"): parts.append(f"_[{block['type']}]_") return "\n\n".join(parts) diff --git a/graph/providers/codex_client.py b/graph/providers/codex_client.py index 51516ce51..af5180bba 100644 --- a/graph/providers/codex_client.py +++ b/graph/providers/codex_client.py @@ -1,43 +1,53 @@ -"""The Codex Responses client — reasoning-item hygiene on the way out (ADR 0097). +"""The Codex Responses client — encrypted-reasoning capture + replay (ADR 0097). The ChatGPT/Codex backend runs with ``store=false``: reasoning state is not kept -server-side, so a replayed reasoning item has to carry its own -``encrypted_content`` blob (which is why the builder asks for -``include=["reasoning.encrypted_content"]``). - -langchain-openai's **streaming** Responses path never captures that blob. It reads -the reasoning item at ``response.output_item.added`` — where ``encrypted_content`` -is still null — and the terminal ``response.completed`` event rebuilds the full -message but keeps only ``parsed``/usage/``response_metadata`` from it. What DOES -survive into ``additional_kwargs["reasoning"]`` is the item's ``rs_…`` id, and -that id is replayed on the next turn. protoAgent always streams (the Codex backend -mandates it), so this is the only shape it ever produces. - -The backend rejects that half-replay: - - 400 invalid_encrypted_content — The encrypted content for item rs_… could not - be verified. Reason: Encrypted content could not be decrypted or parsed. - -and because the item rides in ``additional_kwargs`` it is checkpointed, so EVERY -later turn in the thread fails identically — the thread is bricked. Same failure -class as the dangling ``tool_call`` that ``tool_call_repair`` exists to heal. - -Two rules, applied to the outbound payload: - -- **No blob → drop the item.** An id-only reasoning item is a ghost: with - ``store=false`` the backend wrote nothing to look up and has nothing to verify. - Dropping it restores exactly the behaviour ADR 0097's live validation believed - it already had — no replay, stateless continuity. -- **Blob → keep it, drop the ``id``.** ``encrypted_content`` is self-contained; - the id only resolves against stored state that ``store=false`` never wrote. - -Hermes's Codex adapter arrives at the same two rules independently. Capturing the -blob (so replay actually works) is the follow-up: it needs an -``output_item.done`` handler for reasoning items that langchain-openai lacks. +server-side, so cross-turn reasoning continuity depends on threading each item's +own ``encrypted_content`` blob back through history. #3199 established the first +half of the contract — never send an item the backend cannot verify. This module +now also delivers the other half: capture the blob, and only replay it where it +can actually be decrypted. + +**Why capture needs code at all.** langchain-openai's streaming Responses path +reads a reasoning item at ``response.output_item.added``, where +``encrypted_content`` is still null, and has no ``response.output_item.done`` +branch for reasoning (it has one for ``compaction``, which carries the same kind +of blob). The terminal ``response.completed`` event rebuilds the full message but +keeps only ``parsed``/usage/``response_metadata``. So the blob is visible exactly +once, in an event the converter drops on the floor. `_install_reasoning_capture` +re-emits that event as a content-block delta which merges onto the reasoning block +already in flight — by ``index``, the way every other streamed block merges. + +The wrapper is installed on the module-level converter (there is no instance seam) +but is **inert unless ``_CAPTURE_ISSUER`` is set**, and only this module's client +sets it, for the duration of its own stream. Any other ``ChatOpenAI`` in the +process — a gateway client, a plain Responses user — goes through the original +code path unchanged. + +**Why the issuer stamp.** ``encrypted_content`` is sealed to the endpoint *and +account* that minted it; replaying a blob anywhere else is a hard +``400 invalid_encrypted_content`` that, once checkpointed, bricks the thread. That +is not hypothetical here: protoAgent lets every slot name its own connection +(``gateway:`` / ``anthropic-oauth:`` / ``openai-codex:``), lets each chat tab +override the model per turn, and retries a failed turn against the fallback chain. +So each captured item carries a fingerprint of its issuer, and replay drops items +minted elsewhere instead of poisoning the request. Hermes's Codex adapter reaches +the same design (`_classify_responses_issuer`); the stamp is a salted digest so a +checkpoint never stores a raw account id. + +Outbound rules, applied to every Responses ``input``: + +- **No blob → drop.** An id-only reasoning item is a ghost: with ``store=false`` + the backend wrote nothing to look up and has nothing to verify (#3199). +- **Foreign issuer → drop.** The current endpoint cannot decrypt it. Unstamped + items (written before this landed) are still replayed. +- **Otherwise keep the blob, drop the ``id``.** An item id only resolves against + stored state that ``store=false`` never wrote; the blob is self-contained. """ from __future__ import annotations +import contextvars +import hashlib import logging from typing import Any @@ -45,54 +55,175 @@ class as the dangling ``tool_call`` that ``tool_call_repair`` exists to heal. log = logging.getLogger("protoagent.providers.openai_codex") -# One warning per process: a long thread carries many ghost items and every turn -# re-sends them, so an un-throttled log would drown the turn's real output. -_GHOST_WARNED = False +# Our own key on a captured reasoning block. Never goes on the wire — the +# sanitizer below strips every key with this prefix on the way out. +_PRIVATE_PREFIX = "_protoagent_" +ISSUER_KEY = f"{_PRIVATE_PREFIX}issuer" +# Set by this module's client for the duration of ITS stream; the capture wrapper +# is a no-op for every other caller of the shared converter. +_CAPTURE_ISSUER: contextvars.ContextVar[str] = contextvars.ContextVar("protoagent_codex_issuer", default="") -def sanitize_responses_input(items: Any) -> Any: - """Drop un-verifiable reasoning items from a Responses ``input`` list. +# One warning per process per cause: a long thread carries many affected items and +# re-sends them every turn, so un-throttled logs would drown the turn's output. +_WARNED: set[str] = set() - Returns ``items`` untouched when it isn't a list (nothing to sanitize) so the - caller can apply this to any payload shape. + +def issuer_fingerprint(base_url: str, account_id: str) -> str: + """A stable, non-identifying id for the endpoint+account that mints blobs. + + Digested rather than stored raw: this value is checkpointed alongside the + conversation, and the account id has no business living in that file. + """ + raw = f"{(base_url or '').strip().rstrip('/')}\x00{(account_id or '').strip()}" + return hashlib.sha256(raw.encode()).hexdigest()[:16] + + +def _warn_once(key: str, message: str, *args: Any) -> None: + if key not in _WARNED: + _WARNED.add(key) + log.warning(message, *args) + + +def sanitize_responses_input(items: Any, *, issuer: str = "") -> Any: + """Drop reasoning items this endpoint cannot verify; strip private keys. + + ``items`` is returned untouched when it isn't a list, so a caller can apply + this to any payload shape. """ if not isinstance(items, list): return items cleaned: list = [] ghosts = 0 + foreign = 0 for item in items: if not isinstance(item, dict) or item.get("type") != "reasoning": cleaned.append(item) continue + blob = item.get("encrypted_content") if not (isinstance(blob, str) and blob): ghosts += 1 continue - cleaned.append({k: v for k, v in item.items() if k != "id"}) + + stamped = item.get(ISSUER_KEY) + if stamped and issuer and stamped != issuer: + foreign += 1 + continue + + # `id` is unresolvable under store=false; private keys are ours, not the + # API's. Everything else (summary, the blob itself) replays as-is. + cleaned.append({k: v for k, v in item.items() if k != "id" and not k.startswith(_PRIVATE_PREFIX)}) if ghosts: - global _GHOST_WARNED - if not _GHOST_WARNED: - _GHOST_WARNED = True - log.warning( - "[openai-codex] dropped %d reasoning item(s) with no encrypted_content " - "from the Responses input. The streaming path does not capture the blob, " - "and replaying the item by id alone is what the backend rejects with " - "400 invalid_encrypted_content. Cross-turn reasoning continuity is off; " - "the turn itself is unaffected.", - ghosts, - ) + _warn_once( + "ghost", + "[openai-codex] dropped %d reasoning item(s) with no encrypted_content from " + "the Responses input — replaying one by id alone is what the backend rejects " + "with 400 invalid_encrypted_content. The turn itself is unaffected.", + ghosts, + ) + if foreign: + _warn_once( + "foreign", + "[openai-codex] dropped %d reasoning item(s) minted by a different endpoint or " + "account — encrypted_content is sealed to its issuer, so this endpoint cannot " + "decrypt them. This is normal after a mid-conversation model swap or a re-login; " + "cross-turn reasoning continuity restarts from here.", + foreign, + ) return cleaned +def _install_reasoning_capture() -> None: + """Teach the shared Responses chunk converter to surface ``encrypted_content``. + + Idempotent, and gated on ``_CAPTURE_ISSUER`` so it changes nothing for any + other client in the process. Delegates to the original for every event — + it only ever ADDS a chunk where the original produced none. + """ + from langchain_openai.chat_models import base as lc_base + + original = lc_base._convert_responses_chunk_to_generation_chunk + if getattr(original, "_protoagent_reasoning_capture", False): + return + + def _capture(chunk, current_index, current_output_index, current_sub_index, *args, **kwargs): + result = original(chunk, current_index, current_output_index, current_sub_index, *args, **kwargs) + issuer = _CAPTURE_ISSUER.get() + if not issuer or result[3] is not None: + return result + if getattr(chunk, "type", "") != "response.output_item.done": + return result + item = getattr(chunk, "item", None) + if getattr(item, "type", "") != "reasoning": + return result + blob = getattr(item, "encrypted_content", None) + if not (isinstance(blob, str) and blob): + return result + + from langchain_core.messages import AIMessageChunk + from langchain_core.outputs import ChatGenerationChunk + + # `index` is what merges this onto the reasoning block opened by the + # item's `.added` event; the summary-delta events in between never + # advance it. `type` is deliberately ABSENT: merge_dicts concatenates + # two equal strings for any key but `id`, so re-sending it would yield + # "reasoningreasoning". `id` is safe (equal values are skipped) and is + # what keeps the merge from binding to a neighbouring block. + block = {"index": current_index, "id": getattr(item, "id", None), "encrypted_content": blob} + block[ISSUER_KEY] = issuer + return ( + current_index, + current_output_index, + current_sub_index, + ChatGenerationChunk(message=AIMessageChunk(content=[block])), + ) + + _capture._protoagent_reasoning_capture = True # type: ignore[attr-defined] + lc_base._convert_responses_chunk_to_generation_chunk = _capture + + class CodexChatOpenAI(_ReasoningChatOpenAI): - """``ChatOpenAI`` for the Codex backend, with reasoning items sanitized on the - way out. Every other payload is unchanged — ``input`` exists only on the - Responses path, so the guard below makes this a no-op anywhere else.""" + """``ChatOpenAI`` for the Codex backend: captures encrypted reasoning on the + way in, and replays only what this endpoint can verify on the way out.""" + + @property + def _issuer(self) -> str: + return getattr(self, "_protoagent_issuer_fp", "") or "" def _get_request_payload(self, input_, *, stop=None, **kwargs): payload = super()._get_request_payload(input_, stop=stop, **kwargs) + # `input` exists only on the Responses path, so this is a no-op elsewhere. if isinstance(payload.get("input"), list): - payload["input"] = sanitize_responses_input(payload["input"]) + payload["input"] = sanitize_responses_input(payload["input"], issuer=self._issuer) return payload + + def _stream_responses(self, *args, **kwargs): + token = _CAPTURE_ISSUER.set(self._issuer) + try: + yield from super()._stream_responses(*args, **kwargs) + finally: + _CAPTURE_ISSUER.reset(token) + + async def _astream_responses(self, *args, **kwargs): + token = _CAPTURE_ISSUER.set(self._issuer) + try: + async for chunk in super()._astream_responses(*args, **kwargs): + yield chunk + finally: + _CAPTURE_ISSUER.reset(token) + + +def build_codex_client(*, issuer: str, **kwargs: Any) -> CodexChatOpenAI: + """A ``CodexChatOpenAI`` stamped with the issuer of the endpoint it talks to. + + The stamp rides as a private attribute rather than a pydantic field — the same + way ``graph.providers.identity`` tags routing identity — so the client's + serialized shape is unchanged. + """ + _install_reasoning_capture() + client = CodexChatOpenAI(**kwargs) + object.__setattr__(client, "_protoagent_issuer_fp", issuer) + return client diff --git a/graph/providers/openai_codex.py b/graph/providers/openai_codex.py index 76a0ed90a..2b7389d01 100644 --- a/graph/providers/openai_codex.py +++ b/graph/providers/openai_codex.py @@ -20,6 +20,7 @@ from __future__ import annotations import logging +import os from typing import TYPE_CHECKING, Any from graph.providers.oauth import resolve_codex_oauth @@ -47,7 +48,7 @@ def build_codex_llm( flags. ``model_name`` overrides ``config.model_name`` for aux/subagent slots. """ # Local import — avoids a cycle at module load (the client subclasses graph.llm's). - from graph.providers.codex_client import CodexChatOpenAI + from graph.providers.codex_client import build_codex_client, issuer_fingerprint creds = resolve_codex_oauth() # raises OAuthCredentialError if none @@ -76,19 +77,24 @@ def build_codex_llm( "base_url": creds.base_url, "api_key": creds.access_token, "use_responses_api": True, - # Pin legacy string content. langchain-openai now DEFAULTS output_version to - # "responses/v1", which packs the answer into structured content blocks that - # protoAgent's answer/rendering pipeline stringifies raw (the console would show - # "[{'type':'text',...}]"). "v0" gives a plain string. The block format enables - # cross-turn encrypted-reasoning replay — that's the ADR 0097 multi-turn - # follow-up, wired once the pipeline understands the blocks. - "output_version": "v0", + # Content blocks, not the legacy "v0" string. This is the follow-up the v0 + # pin was holding open: "v0" collapses a turn's reasoning into ONE + # additional_kwargs slot (later items overwrite the first, and streamed + # fragments of two different items merge into each other), so it cannot + # carry per-item encrypted blobs. "responses/v1" keeps each reasoning item + # as its own block, in order, and langchain replays them that way — which + # is what makes cross-turn encrypted-reasoning continuity possible at all. + # The rendering half of the pin was already paid off: every answer site + # reads `AIMessage.text`, which yields text blocks only. + # PROTOAGENT_CODEX_OUTPUT_VERSION is the escape hatch back to "v0" (no + # replay, but no block-shaped content either) if a surface turns out to + # still assume the string shape. + "output_version": os.environ.get("PROTOAGENT_CODEX_OUTPUT_VERSION", "").strip() or "responses/v1", # The ChatGPT backend mandates store=false, so a replayed reasoning item must - # carry its own blob — hence `include`. langchain-openai's STREAMING path - # never captures it (only the `rs_…` id survives), so today this asks for a - # blob nothing reads; `CodexChatOpenAI` drops the id-only leftovers rather - # than replaying an item the backend cannot verify. Kept set so wiring the - # capture is the only step left. + # carry its own blob — hence `include`. langchain-openai's streaming path + # drops the event that carries it; `codex_client` re-surfaces it and stamps + # the issuer, so this now asks for a blob that is actually read, kept, and + # replayed only back to the endpoint that minted it. "store": False, "include": ["reasoning.encrypted_content"], # The Codex backend manages its own output truncation and rejects @@ -105,4 +111,8 @@ def build_codex_llm( if config.top_p is not None: kwargs["top_p"] = config.top_p - return CodexChatOpenAI(**kwargs) + # The blob is sealed to endpoint+account: stamp WHICH one, so a later turn that + # lands on another connection drops those items instead of 400-ing on them. + return build_codex_client( + issuer=issuer_fingerprint(creds.base_url, creds.account_id or ""), **kwargs + ) diff --git a/tests/test_codex_reasoning_capture.py b/tests/test_codex_reasoning_capture.py new file mode 100644 index 000000000..ff152c678 --- /dev/null +++ b/tests/test_codex_reasoning_capture.py @@ -0,0 +1,317 @@ +"""Encrypted-reasoning capture + issuer-scoped replay on the Codex path (ADR 0097). + +#3199 stopped protoAgent sending a reasoning item the backend can't verify. This is +the other direction: capture the blob that makes replay possible at all, and replay +it ONLY back to the endpoint that minted it. +""" + +from __future__ import annotations + +import pytest +from langchain_core.messages import AIMessage, AIMessageChunk, HumanMessage +from langchain_openai.chat_models import base as lc_base +from openai.types.responses import ResponseReasoningItem +from openai.types.responses.response_output_item_added_event import ResponseOutputItemAddedEvent +from openai.types.responses.response_output_item_done_event import ResponseOutputItemDoneEvent +from openai.types.responses.response_text_delta_event import ResponseTextDeltaEvent + +from graph.providers.codex_client import ( + _CAPTURE_ISSUER, + ISSUER_KEY, + CodexChatOpenAI, + _install_reasoning_capture, + build_codex_client, + issuer_fingerprint, + sanitize_responses_input, +) + +ISSUER = issuer_fingerprint("https://chatgpt.com/backend-api/codex", "acct-7") +OTHER_ISSUER = issuer_fingerprint("https://api.openai.com/v1", "acct-7") + + +@pytest.fixture(autouse=True) +def _capture_installed(): + _install_reasoning_capture() + yield + + +def _client(issuer: str = ISSUER) -> CodexChatOpenAI: + return build_codex_client( + issuer=issuer, + model="gpt-5-codex", + api_key="x", + use_responses_api=True, + output_version="responses/v1", + store=False, + include=["reasoning.encrypted_content"], + ) + + +def _reasoning_added(output_index: int, item_id: str): + return ResponseOutputItemAddedEvent( + type="response.output_item.added", + output_index=output_index, + sequence_number=output_index * 2, + item=ResponseReasoningItem(id=item_id, type="reasoning", summary=[]), + ) + + +def _reasoning_done(output_index: int, item_id: str, blob: str | None): + return ResponseOutputItemDoneEvent( + type="response.output_item.done", + output_index=output_index, + sequence_number=output_index * 2 + 1, + item=ResponseReasoningItem(id=item_id, type="reasoning", summary=[], encrypted_content=blob), + ) + + +def _text_delta(output_index: int, text: str): + return ResponseTextDeltaEvent( + type="response.output_text.delta", + output_index=output_index, + content_index=0, + item_id="msg_1", + delta=text, + sequence_number=99, + logprobs=[], + ) + + +def _accumulate(events, *, issuer: str = ISSUER): + """Drive events through the real converter the way `_stream_responses` does.""" + token = _CAPTURE_ISSUER.set(issuer) + try: + idx = out_idx = sub_idx = -1 + acc: AIMessageChunk | None = None + for event in events: + idx, out_idx, sub_idx, gen = lc_base._convert_responses_chunk_to_generation_chunk( + event, idx, out_idx, sub_idx, output_version="responses/v1" + ) + if gen is not None: + acc = gen.message if acc is None else acc + gen.message + return acc + finally: + _CAPTURE_ISSUER.reset(token) + + +# ── capture ───────────────────────────────────────────────────────────────────── + + +def test_the_blob_lands_on_the_reasoning_block_it_belongs_to(): + """`output_item.done` is the only event carrying encrypted_content, and + langchain drops it. It must merge onto the block `.added` opened — by index.""" + acc = _accumulate([_reasoning_added(0, "rs_1"), _reasoning_done(0, "rs_1", "BLOB1"), _text_delta(1, "hi")]) + + reasoning = [b for b in acc.content if b.get("type") == "reasoning"] + assert len(reasoning) == 1 + assert reasoning[0]["id"] == "rs_1" + assert reasoning[0]["encrypted_content"] == "BLOB1" + assert reasoning[0][ISSUER_KEY] == ISSUER + # `type` must survive as a plain string — merge_dicts concatenates equal + # strings for every key but `id`, so re-sending it would give "reasoningreasoning". + assert reasoning[0]["type"] == "reasoning" + assert acc.text == "hi" + + +def test_two_reasoning_items_keep_their_own_blobs_in_order(): + """The v0 slot this replaced could hold only one item per turn; a tool loop + routinely produces several.""" + acc = _accumulate( + [ + _reasoning_added(0, "rs_1"), + _reasoning_done(0, "rs_1", "BLOB1"), + _reasoning_added(1, "rs_2"), + _reasoning_done(1, "rs_2", "BLOB2"), + ] + ) + + reasoning = [b for b in acc.content if b.get("type") == "reasoning"] + assert [(b["id"], b["encrypted_content"]) for b in reasoning] == [ + ("rs_1", "BLOB1"), + ("rs_2", "BLOB2"), + ] + + +def test_capture_is_inert_for_every_other_client(): + """The wrapper sits on a module-level function shared with the gateway path, so + 'off unless we asked for it' is the whole safety argument.""" + acc = _accumulate([_reasoning_added(0, "rs_1"), _reasoning_done(0, "rs_1", "BLOB1")], issuer="") + + reasoning = [b for b in acc.content if b.get("type") == "reasoning"] + assert "encrypted_content" not in reasoning[0] + + +def test_a_done_event_with_no_blob_adds_nothing(): + acc = _accumulate([_reasoning_added(0, "rs_1"), _reasoning_done(0, "rs_1", None)]) + assert "encrypted_content" not in acc.content[0] + + +def test_install_is_idempotent(): + first = lc_base._convert_responses_chunk_to_generation_chunk + _install_reasoning_capture() + _install_reasoning_capture() + assert lc_base._convert_responses_chunk_to_generation_chunk is first + + +# ── issuer fingerprint ────────────────────────────────────────────────────────── + + +def test_fingerprint_separates_endpoints_and_accounts(): + base = "https://chatgpt.com/backend-api/codex" + assert issuer_fingerprint(base, "acct-7") != issuer_fingerprint(base, "acct-8") + assert issuer_fingerprint(base, "acct-7") != issuer_fingerprint("https://api.openai.com/v1", "acct-7") + assert issuer_fingerprint(base + "/", "acct-7") == issuer_fingerprint(base, " acct-7 ") + + +def test_fingerprint_does_not_leak_the_account_id(): + """It is checkpointed next to the conversation; a raw account id has no business + living in that file.""" + fp = issuer_fingerprint("https://chatgpt.com/backend-api/codex", "acct-secret") + assert "acct-secret" not in fp + assert len(fp) == 16 + + +# ── replay ────────────────────────────────────────────────────────────────────── + + +def test_a_blob_from_this_issuer_replays_without_its_id_or_our_private_keys(): + items = [ + { + "type": "reasoning", + "id": "rs_1", + "summary": [], + "encrypted_content": "BLOB1", + ISSUER_KEY: ISSUER, + } + ] + assert sanitize_responses_input(items, issuer=ISSUER) == [ + {"type": "reasoning", "summary": [], "encrypted_content": "BLOB1"} + ] + + +def test_a_blob_from_another_issuer_is_dropped(): + """The exact 400 this whole line of work exists to prevent: a mid-conversation + model swap replaying a Codex-minted blob at an endpoint that can't decrypt it.""" + items = [{"type": "reasoning", "encrypted_content": "BLOB1", ISSUER_KEY: OTHER_ISSUER}] + assert sanitize_responses_input(items, issuer=ISSUER) == [] + + +def test_an_unstamped_blob_still_replays(): + """Written before the stamp existed — dropping it would silently break continuity + on every thread that predates this change.""" + items = [{"type": "reasoning", "encrypted_content": "BLOB1"}] + assert sanitize_responses_input(items, issuer=ISSUER) == [{"type": "reasoning", "encrypted_content": "BLOB1"}] + + +def test_a_ghost_is_still_dropped_whatever_its_issuer(): + items = [{"type": "reasoning", "id": "rs_1", "summary": [], ISSUER_KEY: ISSUER}] + assert sanitize_responses_input(items, issuer=ISSUER) == [] + + +def test_an_unknown_issuer_on_this_client_keeps_everything_verifiable(): + """No fingerprint (an unstamped client) must not become 'drop everything'.""" + items = [{"type": "reasoning", "encrypted_content": "BLOB1", ISSUER_KEY: OTHER_ISSUER}] + assert len(sanitize_responses_input(items, issuer="")) == 1 + + +def test_end_to_end_capture_then_replay(): + """The whole contract in one pass: what the stream captured is what the next + turn sends back — blob kept, id and private keys gone.""" + acc = _accumulate([_reasoning_added(0, "rs_1"), _reasoning_done(0, "rs_1", "BLOB1"), _text_delta(1, "hi")]) + turn = AIMessage(content=acc.content, id="msg_1") + + payload = _client()._get_request_payload([HumanMessage("hi"), turn]) + reasoning = [i for i in payload["input"] if isinstance(i, dict) and i.get("type") == "reasoning"] + + assert len(reasoning) == 1 + assert reasoning[0]["encrypted_content"] == "BLOB1" + assert "id" not in reasoning[0] + assert ISSUER_KEY not in reasoning[0] + assert not any(k.startswith("_protoagent") for k in reasoning[0]) + + +def test_end_to_end_a_foreign_thread_sends_no_reasoning_at_all(): + acc = _accumulate([_reasoning_added(0, "rs_1"), _reasoning_done(0, "rs_1", "BLOB1")]) + turn = AIMessage(content=acc.content, id="msg_1") + + payload = _client(issuer=OTHER_ISSUER)._get_request_payload([HumanMessage("hi"), turn]) + assert [i for i in payload["input"] if isinstance(i, dict) and i.get("type") == "reasoning"] == [] + + +def test_a_legacy_v0_turn_in_the_same_thread_still_sends_no_ghost(): + """Threads checkpointed before the flip carry v0-shaped messages; they must keep + working alongside the new block shape.""" + v0_turn = AIMessage( + content=[{"type": "text", "text": "old"}], + id="msg_0", + additional_kwargs={"reasoning": {"type": "reasoning", "id": "rs_old", "summary": []}}, + ) + payload = _client()._get_request_payload([HumanMessage("hi"), v0_turn]) + assert [i for i in payload["input"] if isinstance(i, dict) and i.get("type") == "reasoning"] == [] + + +# ── the builder wires both halves ─────────────────────────────────────────────── + + +def test_the_builder_stamps_the_live_endpoint_and_account(monkeypatch): + import graph.providers.openai_codex as ocx + from graph.config import LangGraphConfig + from graph.llm import create_llm + from graph.providers.oauth import CodexOAuthCreds + + monkeypatch.setattr( + ocx, + "resolve_codex_oauth", + lambda *a, **k: CodexOAuthCreds( + access_token="t", + account_id="acct-7", + base_url="https://chatgpt.com/backend-api/codex", + source="instance_store", + ), + ) + llm = create_llm(LangGraphConfig(model_provider="openai-codex", model_name="gpt-5-codex")) + + assert llm.output_version == "responses/v1" + assert llm._issuer == ISSUER + + +def test_the_output_version_escape_hatch(monkeypatch): + """Reverting to v0 costs replay but restores the old content shape, if some + surface turns out to still assume the string.""" + import graph.providers.openai_codex as ocx + from graph.config import LangGraphConfig + from graph.llm import create_llm + from graph.providers.oauth import CodexOAuthCreds + + monkeypatch.setattr( + ocx, + "resolve_codex_oauth", + lambda *a, **k: CodexOAuthCreds(access_token="t", account_id="a", base_url="b", source="s"), + ) + monkeypatch.setenv("PROTOAGENT_CODEX_OUTPUT_VERSION", "v0") + llm = create_llm(LangGraphConfig(model_provider="openai-codex", model_name="gpt-5-codex")) + assert llm.output_version == "v0" + + +# ── reasoning never reaches the surfaces that WRITE what they read ────────────── + + +def test_reasoning_blocks_are_skipped_not_placeholdered(): + """`text_of` feeds exports, session memory and chat bundles — all of which + persist. ADR 0021: reasoning is never persisted.""" + from graph.message_blocks import text_of + + msg = AIMessage( + content=[ + {"type": "reasoning", "summary": [{"type": "summary_text", "text": "secret"}], "encrypted_content": "B"}, + {"type": "text", "text": "the answer"}, + ] + ) + assert text_of(msg) == "the answer" + + +def test_other_non_text_blocks_still_placeholder(): + from graph.message_blocks import text_of + + msg = AIMessage(content=[{"type": "image_url", "image_url": {"url": "x"}}, {"type": "text", "text": "hi"}]) + assert text_of(msg) == "_[image_url]_\n\nhi" From 5e318bfb23331a388ada45ba5bd4319f65e0a8ae Mon Sep 17 00:00:00 2001 From: Josh Mabry Date: Thu, 27 Aug 2026 13:33:06 -0700 Subject: [PATCH 2/4] docs(changelog): fragment for #3207 Co-Authored-By: Claude Opus 5 (1M context) --- changelog.d/3207.added.md | 8 ++++++++ 1 file changed, 8 insertions(+) create mode 100644 changelog.d/3207.added.md diff --git a/changelog.d/3207.added.md b/changelog.d/3207.added.md new file mode 100644 index 000000000..08734b6ca --- /dev/null +++ b/changelog.d/3207.added.md @@ -0,0 +1,8 @@ +- **Codex/ChatGPT-subscription chats keep their reasoning across turns again (#3207).** With + `store=false` the model's reasoning only survives a turn if its encrypted blob is threaded + back, and protoAgent never captured that blob — the streaming path drops the one event that + carries it — so every turn on `model.provider: openai-codex` started its reasoning from + scratch. It is now captured and replayed, and each item is stamped with the endpoint and + account that minted it: a blob is only ever sent back to the issuer that can decrypt it, so + switching a chat's model mid-thread (or signing in under a different ChatGPT account) costs + reasoning continuity from that point rather than failing the turn. From 87c13eac77a728c49797f5ce7def6972399e6913 Mon Sep 17 00:00:00 2001 From: Josh Mabry Date: Thu, 27 Aug 2026 13:48:17 -0700 Subject: [PATCH 3/4] fix(codex): arm the reasoning capture on the dispatched stream seam MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The contextvar that gates the capture was set in an override of _stream_responses — which ChatOpenAI._stream reaches via `super()._stream_responses(...)`, an explicit super(ChatOpenAI, self) bind that skips the subclass entirely. So the override was dead code and the capture would never have armed in production, while every unit test that set the contextvar by hand still passed. Armed on _stream/_astream instead (dispatched on the instance), and covered by a test that drives the client's OWN stream against a stubbed root_client, so nothing but the client can set the contextvar. Verified to fail against the old seam. Co-Authored-By: Claude Opus 5 (1M context) --- graph/providers/codex_client.py | 15 ++++-- tests/test_codex_reasoning_capture.py | 66 +++++++++++++++++++++++++++ 2 files changed, 77 insertions(+), 4 deletions(-) diff --git a/graph/providers/codex_client.py b/graph/providers/codex_client.py index af5180bba..f00cfb4b4 100644 --- a/graph/providers/codex_client.py +++ b/graph/providers/codex_client.py @@ -200,17 +200,24 @@ def _get_request_payload(self, input_, *, stop=None, **kwargs): payload["input"] = sanitize_responses_input(payload["input"], issuer=self._issuer) return payload - def _stream_responses(self, *args, **kwargs): + # Arming happens on `_stream`/`_astream`, NOT on `_stream_responses`. + # `ChatOpenAI._stream` routes with `super()._stream_responses(...)` — an + # explicit `super(ChatOpenAI, self)` bind that skips right past this subclass, + # so an override there is dead code and the capture would never arm in + # production while every unit test that sets the contextvar itself still + # passes. `_stream`/`_astream` are dispatched on the instance, so they are the + # seam that actually runs. + def _stream(self, *args, **kwargs): token = _CAPTURE_ISSUER.set(self._issuer) try: - yield from super()._stream_responses(*args, **kwargs) + yield from super()._stream(*args, **kwargs) finally: _CAPTURE_ISSUER.reset(token) - async def _astream_responses(self, *args, **kwargs): + async def _astream(self, *args, **kwargs): token = _CAPTURE_ISSUER.set(self._issuer) try: - async for chunk in super()._astream_responses(*args, **kwargs): + async for chunk in super()._astream(*args, **kwargs): yield chunk finally: _CAPTURE_ISSUER.reset(token) diff --git a/tests/test_codex_reasoning_capture.py b/tests/test_codex_reasoning_capture.py index ff152c678..79d505617 100644 --- a/tests/test_codex_reasoning_capture.py +++ b/tests/test_codex_reasoning_capture.py @@ -315,3 +315,69 @@ def test_other_non_text_blocks_still_placeholder(): msg = AIMessage(content=[{"type": "image_url", "image_url": {"url": "x"}}, {"type": "text", "text": "hi"}]) assert text_of(msg) == "_[image_url]_\n\nhi" + + +# ── the capture actually arms on the real stream path ─────────────────────────── + + +class _FakeStream: + """Stands in for `root_client.responses.create(...)` — a context manager over + a canned Responses event sequence.""" + + def __init__(self, events): + self._events = events + + def __enter__(self): + return iter(self._events) + + def __exit__(self, *exc): + return False + + +class _FakeResponses: + def __init__(self, events): + self._events = events + self.payloads: list[dict] = [] + + def create(self, **payload): + self.payloads.append(payload) + return _FakeStream(self._events) + + +class _FakeRootClient: + def __init__(self, events): + self.responses = _FakeResponses(events) + + +def test_capture_arms_on_the_real_stream_path(): + """Regression: arming was first written on `_stream_responses`, which + `ChatOpenAI._stream` reaches via `super()._stream_responses(...)` — an explicit + `super(ChatOpenAI, self)` bind that skips the subclass entirely. The override was + dead code, and every test that set the contextvar by hand still passed. This one + drives the client's OWN stream, so nothing sets the contextvar but the client.""" + client = _client() + events = [ + _reasoning_added(0, "rs_1"), + _reasoning_done(0, "rs_1", "BLOB1"), + _text_delta(1, "hi"), + ] + object.__setattr__(client, "root_client", _FakeRootClient(events)) + + acc = None + for gen in client._stream([HumanMessage("hi")]): + acc = gen.message if acc is None else acc + gen.message + + reasoning = [b for b in acc.content if b.get("type") == "reasoning"] + assert reasoning and reasoning[0]["encrypted_content"] == "BLOB1" + assert reasoning[0][ISSUER_KEY] == ISSUER + assert acc.text == "hi" + + +def test_the_contextvar_is_released_after_the_stream(): + """It gates a process-wide wrapper — leaking it would arm capture for every + other client in the process.""" + client = _client() + object.__setattr__(client, "root_client", _FakeRootClient([_text_delta(0, "hi")])) + + list(client._stream([HumanMessage("hi")])) + assert _CAPTURE_ISSUER.get() == "" From 86d67ac9cd20f94960ddae1eb2c5bcce2efdbdfa Mon Sep 17 00:00:00 2001 From: Josh Mabry Date: Thu, 27 Aug 2026 14:05:28 -0700 Subject: [PATCH 4/4] fix(codex): stamp the issuer on langchain-openai >= 1.6 too MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit CI installs deps UNPINNED while uv.lock holds langchain-openai 1.3.0, so CI runs 1.6.0 and local runs 1.3.0. 1.6.0 added the reasoning `response.output_item.done` branch this wrapper existed to supply — and the wrapper bailed out whenever the original produced a chunk, so on 1.6.0 the blob was captured but NEVER STAMPED. That is not a test artifact: an unstamped item is treated as legacy and always replayed, i.e. the cross-issuer guard was silently inoperative on exactly the version CI (and any fresh install) uses. CI caught it; the local suite could not. The wrapper now augments instead of bailing: it stamps the issuer whether the installed langchain surfaced the blob or we had to synthesize it. Verified against BOTH dependency sets — 1.3.0 (uv.lock) and 1.6.0 (a CI-equivalent venv): full suite green on each. Co-Authored-By: Claude Opus 5 (1M context) --- graph/providers/codex_client.py | 67 ++++++++++++++++----------- tests/test_codex_reasoning_capture.py | 7 ++- 2 files changed, 46 insertions(+), 28 deletions(-) diff --git a/graph/providers/codex_client.py b/graph/providers/codex_client.py index f00cfb4b4..3130398a1 100644 --- a/graph/providers/codex_client.py +++ b/graph/providers/codex_client.py @@ -7,21 +7,24 @@ now also delivers the other half: capture the blob, and only replay it where it can actually be decrypted. -**Why capture needs code at all.** langchain-openai's streaming Responses path -reads a reasoning item at ``response.output_item.added``, where -``encrypted_content`` is still null, and has no ``response.output_item.done`` -branch for reasoning (it has one for ``compaction``, which carries the same kind -of blob). The terminal ``response.completed`` event rebuilds the full message but -keeps only ``parsed``/usage/``response_metadata``. So the blob is visible exactly -once, in an event the converter drops on the floor. `_install_reasoning_capture` -re-emits that event as a content-block delta which merges onto the reasoning block -already in flight — by ``index``, the way every other streamed block merges. - -The wrapper is installed on the module-level converter (there is no instance seam) -but is **inert unless ``_CAPTURE_ISSUER`` is set**, and only this module's client -sets it, for the duration of its own stream. Any other ``ChatOpenAI`` in the -process — a gateway client, a plain Responses user — goes through the original -code path unchanged. +**Why capture needs code at all.** ``encrypted_content`` reaches a stream in +exactly one event — ``response.output_item.done`` for the reasoning item. The +``.added`` event that opens the item has the field still null, and the terminal +``response.completed`` event rebuilds the full message but keeps only +``parsed``/usage/``response_metadata``. langchain-openai **≥ 1.6** handles that +event; **< 1.6** drops it on the floor, so the blob was never captured at all. +``pyproject`` floors langchain-openai at 1.0, so both are live. + +`_install_reasoning_capture` therefore does two things at that one event: it +synthesizes the content-block delta when the installed langchain didn't (merging +onto the in-flight reasoning block by ``index``, the way every other streamed +block merges), and it stamps the issuer either way. On a modern langchain the +first half is a no-op — the wrapper only ever adds what is missing. + +It is installed on the module-level converter (there is no instance seam) but is +**inert unless ``_CAPTURE_ISSUER`` is set**, and only this module's client sets +it, for the duration of its own stream. Any other ``ChatOpenAI`` in the process — +a gateway client, a plain Responses user — is untouched. **Why the issuer stamp.** ``encrypted_content`` is sealed to the endpoint *and account* that minted it; replaying a blob anywhere else is a hard @@ -152,9 +155,7 @@ def _install_reasoning_capture() -> None: def _capture(chunk, current_index, current_output_index, current_sub_index, *args, **kwargs): result = original(chunk, current_index, current_output_index, current_sub_index, *args, **kwargs) issuer = _CAPTURE_ISSUER.get() - if not issuer or result[3] is not None: - return result - if getattr(chunk, "type", "") != "response.output_item.done": + if not issuer or getattr(chunk, "type", "") != "response.output_item.done": return result item = getattr(chunk, "item", None) if getattr(item, "type", "") != "reasoning": @@ -166,14 +167,28 @@ def _capture(chunk, current_index, current_output_index, current_sub_index, *arg from langchain_core.messages import AIMessageChunk from langchain_core.outputs import ChatGenerationChunk - # `index` is what merges this onto the reasoning block opened by the - # item's `.added` event; the summary-delta events in between never - # advance it. `type` is deliberately ABSENT: merge_dicts concatenates - # two equal strings for any key but `id`, so re-sending it would yield - # "reasoningreasoning". `id` is safe (equal values are skipped) and is - # what keeps the merge from binding to a neighbouring block. - block = {"index": current_index, "id": getattr(item, "id", None), "encrypted_content": blob} - block[ISSUER_KEY] = issuer + generation = result[3] + if generation is not None: + # langchain >= 1.6 already surfaces the blob; only the stamp is ours. + content = [ + {**b, ISSUER_KEY: issuer} if isinstance(b, dict) and b.get("encrypted_content") else b + for b in generation.message.content + ] + message = generation.message.model_copy(update={"content": content}) + return (result[0], result[1], result[2], ChatGenerationChunk(message=message)) + + # langchain < 1.6 drops this event, so the blob has to be re-emitted. + # `index` merges it onto the block the item's `.added` event opened; the + # summary deltas in between never advance it. `type` is deliberately + # ABSENT — merge_dicts concatenates two equal strings for any key but + # `id`, so re-sending it would yield "reasoningreasoning". `id` IS safe + # (equal values are skipped) and keeps the merge off a neighbouring block. + block = { + "index": current_index, + "id": getattr(item, "id", None), + "encrypted_content": blob, + ISSUER_KEY: issuer, + } return ( current_index, current_output_index, diff --git a/tests/test_codex_reasoning_capture.py b/tests/test_codex_reasoning_capture.py index 79d505617..aea8a931c 100644 --- a/tests/test_codex_reasoning_capture.py +++ b/tests/test_codex_reasoning_capture.py @@ -134,11 +134,14 @@ def test_two_reasoning_items_keep_their_own_blobs_in_order(): def test_capture_is_inert_for_every_other_client(): """The wrapper sits on a module-level function shared with the gateway path, so - 'off unless we asked for it' is the whole safety argument.""" + 'off unless we asked for it' is the whole safety argument. Asserted on the STAMP, + not the blob: langchain >= 1.6 surfaces the blob on its own, and suppressing that + for other clients would be a regression, not a safeguard.""" acc = _accumulate([_reasoning_added(0, "rs_1"), _reasoning_done(0, "rs_1", "BLOB1")], issuer="") reasoning = [b for b in acc.content if b.get("type") == "reasoning"] - assert "encrypted_content" not in reasoning[0] + assert ISSUER_KEY not in reasoning[0] + assert not any(k.startswith("_protoagent") for k in reasoning[0]) def test_a_done_event_with_no_blob_adds_nothing():