fix(codex): scope encrypted-reasoning replay to the issuing model

Encrypted reasoning blobs are sealed to the model that minted them, not
only to the endpoint. Switching models on the same custom Responses
endpoint therefore replayed blobs the new model cannot decrypt and the
turn failed with HTTP 400.

Stamp captured reasoning items with `_issuer_model` (the canonical wire
model) alongside `_issuer_kind`, and replay an item only when both the
issuer kind and the model match the current request. Endpoint-stamped
legacy items without model provenance are dropped once the current
model is known (fail closed); ordinary assistant text stays replayable.
The transport threads the effective wire model (request_overrides win)
into conversion and normalization; the auxiliary Codex adapter stamps
and filters against its own model rather than the main agent's. The
400 classifier also recognises the custom-endpoint wording
"encrypted content could not be decrypted or parsed" so recovery strips
the replay state instead of aborting.

Hand-grafted from #95849 (final head d9cf6bcc08) onto current main; the
middleware-model-rewrite half is intentionally left out.

Closes #95834
This commit is contained in:
Fangliquan
2026-09-14 23:27:12 +05:30
committed by kshitij
parent f7b6a2b59f
commit 51ebdff570
6 changed files with 131 additions and 35 deletions
+15 -7
View File
@@ -107,10 +107,7 @@ def aux_probe_mode():
from agent.credential_pool import load_pool
from agent.model_metadata import (
MINIMUM_CONTEXT_LENGTH, get_model_context_length,
strip_codex_context_variant_suffix as _strip_codex_ctx_variant,
)
from agent.model_metadata import MINIMUM_CONTEXT_LENGTH, get_model_context_length
from hermes_cli.config import get_hermes_home
from agent.auxiliary_health import _custom_health_base_url, _unhealthy_cache_key
from hermes_constants import OPENROUTER_BASE_URL, hermes_home_key
@@ -1340,8 +1337,13 @@ class _CodexCompletionsAdapter:
# includes assistant tool_calls + role="tool" results). The shared converter encodes assistant tool
# calls as `function_call` items and tool results as `function_call_output` items with a valid
# call_id, so every Responses path normalizes tool history identically and cannot drift.
from agent.codex_responses_adapter import _chat_messages_to_responses_input
from agent.codex_responses_adapter import (
_chat_messages_to_responses_input,
_classify_responses_issuer,
_wire_model_identity,
)
model = kwargs.get("model", self._model)
wire_model = _wire_model_identity(model)
host = str(getattr(self._client, "base_url", "") or "")
is_xai = base_url_host_matches(host, "x.ai") or base_url_host_matches(host, "api.x.ai")
is_copilot = base_url_host_matches(host, "githubcopilot.com")
@@ -1363,12 +1365,18 @@ class _CodexCompletionsAdapter:
# Auxiliary calls (context compression, flush_memories, MoA aggregation) go through this adapter
# instead of agent/transports/codex.py's build_kwargs, so they need the same guard applied
# independently. See #32716.
# Aux requests run their own model; stamp/filter reasoning provenance against it, not the main agent's.
input_items = _chat_messages_to_responses_input(
replay_messages, is_github_responses=is_copilot, native_compaction_eligible=False
replay_messages, is_github_responses=is_copilot,
current_issuer_kind=_classify_responses_issuer(
is_xai_responses=is_xai, is_github_responses=is_github,
is_codex_backend=base_url_host_matches(host, "chatgpt.com"), base_url=host,
),
current_issuer_model=wire_model, native_compaction_eligible=False,
)
resp_kwargs: Dict[str, Any] = {
# Codex only knows the base slug; strip the Hermes ``-900k`` picker suffix.
"model": _strip_codex_ctx_variant(model), "instructions": instructions,
"model": wire_model, "instructions": instructions,
"input": input_items or [{"role": "user", "content": ""}], "store": False,
}
# Forward the chat.completions timeout; otherwise a Codex stream can sit behind a
+45 -18
View File
@@ -33,6 +33,14 @@ def _classify_responses_issuer(
# Per-process throttle for the cross-issuer skip warning.
_CROSS_ISSUER_WARN_EMITTED = False
def _wire_model_identity(model: Any) -> Optional[str]:
"""Canonical Responses wire model stamped on encrypted reasoning: blobs are sealed to the issuing
model too, so a same-endpoint model switch must not replay them (HTTP 400)."""
from agent.model_metadata import strip_codex_context_variant_suffix
return str(strip_codex_context_variant_suffix(model or "")).strip() or None
# Codex/Harmony tool-call serialization leaked into assistant text (no structured function_call).
_TOOL_CALL_LEAK_PATTERN = re.compile(r"(?:^|[\s>|])to=functions\.[A-Za-z_][\w.]*", re.IGNORECASE)
@@ -336,13 +344,15 @@ def _assistant_message_item(
def _replay_reasoning_items(
msg: Dict[str, Any], *, seen_item_ids: set, current_issuer_kind: Optional[str], native_compaction_eligible: bool,
msg: Dict[str, Any], *, seen_item_ids: set, current_issuer_kind: Optional[str],
current_issuer_model: Optional[str] = None, native_compaction_eligible: bool,
) -> List[Dict[str, Any]]:
"""Replay persisted encrypted reasoning/compaction items for one assistant turn. Skips duplicate
ids, ``compaction`` checkpoints unless THIS request carries ``context_management`` (else a persisted
checkpoint erases pre-checkpoint history on a model that cannot decrypt it), and items stamped by
another issuer (HTTP 400); unstamped legacy items pass. ``id`` (store=False lookups 404) and
``_issuer_kind`` are stripped."""
another issuer or model (HTTP 400). Endpoint-stamped legacy items without model provenance drop
(fail closed) when the current model is known; fully unstamped items pass. ``id`` (store=False
lookups 404) and the Hermes provenance fields are stripped."""
global _CROSS_ISSUER_WARN_EMITTED
replayed: List[Dict[str, Any]] = []
for ri in _as_list(msg.get("codex_reasoning_items")):
@@ -352,16 +362,22 @@ def _replay_reasoning_items(
if (item_id and item_id in seen_item_ids) or (ri.get("type") == "compaction" and not native_compaction_eligible):
continue
item_issuer = ri.get("_issuer_kind")
if current_issuer_kind is not None and item_issuer is not None and item_issuer != current_issuer_kind:
item_model = ri.get("_issuer_model")
foreign_issuer = current_issuer_kind is not None and item_issuer is not None and item_issuer != current_issuer_kind
foreign_model = current_issuer_model is not None and (
(item_model is not None and item_model != current_issuer_model)
or (item_model is None and item_issuer is not None)
)
if foreign_issuer or foreign_model:
if not _CROSS_ISSUER_WARN_EMITTED:
logger.warning(
"Dropping reasoning item minted by %s while calling %s — encrypted_content is sealed to "
"its issuer. This happens when a session switches model providers mid-conversation.",
item_issuer, current_issuer_kind,
"Dropping reasoning item minted by %s/%s while calling %s/%s — encrypted_content is "
"sealed to its issuer and model. This happens when a session switches model mid-conversation.",
item_issuer, item_model, current_issuer_kind, current_issuer_model,
)
_CROSS_ISSUER_WARN_EMITTED = True
continue
replayed.append({k: v for k, v in ri.items() if k not in ("id", "_issuer_kind")})
replayed.append({k: v for k, v in ri.items() if k not in ("id", "_issuer_kind", "_issuer_model")})
if item_id:
seen_item_ids.add(item_id)
return replayed
@@ -428,7 +444,7 @@ def _tool_output_items(msg: Dict[str, Any]) -> List[Dict[str, Any]]:
def _chat_messages_to_responses_input(
messages: List[Dict[str, Any]], *, is_xai_responses: bool = False, is_github_responses: bool = False,
replay_encrypted_reasoning: bool = True, current_issuer_kind: Optional[str] = None,
native_compaction_eligible: bool = False,
current_issuer_model: Optional[str] = None, native_compaction_eligible: bool = False,
) -> List[Dict[str, Any]]:
"""Convert internal chat-style messages to Responses input items.
@@ -436,7 +452,8 @@ def _chat_messages_to_responses_input(
``replay_encrypted_reasoning``: per-session kill switch, threaded False by
``AIAgent._disable_codex_reasoning_replay`` after an ``invalid_encrypted_content`` 400.
``is_github_responses``: drops ``id`` from replayed message items (Copilot 401s on stale ids).
``current_issuer_kind``: cross-issuer guard; foreign-stamped items drop, legacy items replay.
``current_issuer_kind`` / ``current_issuer_model``: provenance guard; foreign-stamped items drop, as do
endpoint-stamped legacy items without model provenance when the current model is known.
``native_compaction_eligible``: THIS request carries ``context_management``; gates both replaying ``compaction``
checkpoints and ``prune_pre_checkpoint_items``. Checkpoints persist across model swaps / compression flips / resume,
so without the gate one checkpoint would erase pre-checkpoint history on a model that cannot decrypt it (lossless:
@@ -498,7 +515,7 @@ def _chat_messages_to_responses_input(
continue
reasoning_items = [] if not replay_encrypted_reasoning else _replay_reasoning_items(
msg, seen_item_ids=seen_item_ids, current_issuer_kind=current_issuer_kind,
native_compaction_eligible=native_compaction_eligible,
current_issuer_model=current_issuer_model, native_compaction_eligible=native_compaction_eligible,
)
emit(reasoning_items, msg)
message_items = _replay_message_items(
@@ -569,13 +586,17 @@ def _native_responses_replay_items(
return None
route = classify_responses_route(agent)._asdict()
from agent.native_compaction import native_compaction_context_management
from agent.fast_mode import effective_request_overrides
if not native_compaction_context_management(agent, **route):
return None
# The wire model may be rewritten per request (fast mode); provenance must match what the transport stamps.
effective_model = effective_request_overrides(agent).get("model", getattr(agent, "model", None))
try:
items = _chat_messages_to_responses_input(
messages, is_xai_responses=route["is_xai_responses"], is_github_responses=route["is_github_responses"],
replay_encrypted_reasoning=bool(getattr(agent, "_codex_reasoning_replay_enabled", True)),
current_issuer_kind=_classify_responses_issuer(base_url=getattr(agent, "base_url", None), **route),
current_issuer_model=_wire_model_identity(effective_model),
native_compaction_eligible=True,
)
except Exception:
@@ -916,8 +937,10 @@ def _response_tool_call(item: Any, item_type: str, index: int) -> SimpleNamespac
)
def _capture_encrypted_item(item: Any, item_type: str, issuer_kind: Optional[str]) -> Optional[Dict[str, Any]]:
"""``{type, encrypted_content[, _issuer_kind]}`` for replay, or None without a blob. Reasoning
def _capture_encrypted_item(
item: Any, item_type: str, issuer_kind: Optional[str], issuer_model: Optional[str] = None,
) -> Optional[Dict[str, Any]]:
"""``{type, encrypted_content[, _issuer_kind, _issuer_model]}`` for replay, or None without a blob. Reasoning
items also carry ``id`` + ``summary`` (required by the API on replay); transient ``rs_tmp_`` skip."""
encrypted = getattr(item, "encrypted_content", None)
if not _nonempty_str(encrypted):
@@ -925,6 +948,8 @@ def _capture_encrypted_item(item: Any, item_type: str, issuer_kind: Optional[str
raw_item: Dict[str, Any] = {"type": item_type, "encrypted_content": encrypted}
if issuer_kind:
raw_item["_issuer_kind"] = issuer_kind
if issuer_model:
raw_item["_issuer_model"] = issuer_model
if item_type != "reasoning":
return raw_item
item_id = getattr(item, "id", None)
@@ -950,7 +975,7 @@ class _OutputScan:
self.saw_streaming_or_item_incomplete = response_status in {"queued", "in_progress"}
self.saw_commentary_phase = self.saw_final_answer_phase = self.saw_reasoning_item = False
def scan(self, output: List[Any], issuer_kind: Optional[str]) -> None:
def scan(self, output: List[Any], issuer_kind: Optional[str], issuer_model: Optional[str] = None) -> None:
for item in output:
item_type = getattr(item, "type", None)
item_status = _lower_or_none(getattr(item, "status", None))
@@ -967,7 +992,7 @@ class _OutputScan:
self.reasoning_parts.append(reasoning_text)
# Compaction checkpoints ride the codex_reasoning_items sidecar (persistence,
# replay, cross-issuer guard and kill switch for free).
raw_item = _capture_encrypted_item(item, item_type, issuer_kind)
raw_item = _capture_encrypted_item(item, item_type, issuer_kind, issuer_model)
if raw_item is not None:
self.reasoning_items_raw.append(raw_item)
if item_type == "compaction":
@@ -995,9 +1020,11 @@ class _OutputScan:
))
def _normalize_codex_response(response: Any, *, issuer_kind: Optional[str] = None) -> tuple[Any, str]:
def _normalize_codex_response(
response: Any, *, issuer_kind: Optional[str] = None, issuer_model: Optional[str] = None,
) -> tuple[Any, str]:
"""Normalize a Responses API object to ``(assistant_message, finish_reason)``.
``issuer_kind`` is stamped onto captured reasoning items for cross-issuer replay drops."""
``issuer_kind`` / ``issuer_model`` are stamped onto captured reasoning items for provenance replay drops."""
response_status = _lower_or_none(getattr(response, "status", None))
incomplete_reason = str(_field(getattr(response, "incomplete_details", None), "reason", "") or "").strip().lower()
response_incomplete_content_filter = response_status == "incomplete" and incomplete_reason == "content_filter"
@@ -1021,7 +1048,7 @@ def _normalize_codex_response(response: Any, *, issuer_kind: Optional[str] = Non
if response_status in {"failed", "cancelled"}:
raise RuntimeError(_format_responses_error(getattr(response, "error", None), response_status))
scan = _OutputScan(response_status)
scan.scan(output, issuer_kind)
scan.scan(output, issuer_kind, issuer_model)
tool_calls, reasoning_parts = scan.tool_calls, scan.reasoning_parts
final_text = "\n".join(scan.content_parts).strip()
if not final_text and (scan.saw_final_answer_phase or not scan.saw_commentary_phase):
+3
View File
@@ -790,6 +790,9 @@ def _classify_400(c: _Ctx) -> Verdict:
if code == "invalid_encrypted_content" or "invalid_encrypted_content" in msg or (
"encrypted content for item" in msg and "could not be verified" in msg
) or "could not decrypt the provided encrypted_content" in msg or (
# Custom Responses endpoints wrap a replay rejection in a generic bad_request (#95834).
"encrypted content could not be decrypted or parsed" in msg
) or (
# Azure Foundry (gpt-6-astra) rejects replayed reasoning from several prior responses this way (#105369).
"conflicting authenticated continuation identities" in msg
):
+14 -6
View File
@@ -489,6 +489,7 @@ class ResponsesApiTransport(ProviderTransport):
# Issuer kind of the most recent build_kwargs/convert_messages call (normalize_response fallback).
_last_issuer_kind: Optional[str] = None
_last_issuer_model: Optional[str] = None
# ``{wire_alias: original}`` of the most recent build_kwargs. None = no request built (legacy map).
_last_wire_aliases: Optional[dict[str, str]] = None
@@ -510,13 +511,15 @@ class ResponsesApiTransport(ProviderTransport):
def convert_messages(self, messages: list[dict[str, Any]], **kwargs) -> Any:
"""Convert OpenAI chat messages to Responses API input items."""
from agent.codex_responses_adapter import _chat_messages_to_responses_input
from agent.codex_responses_adapter import _chat_messages_to_responses_input, _wire_model_identity
self._last_issuer_model = _wire_model_identity(kwargs.get("model"))
return _chat_messages_to_responses_input(
messages, is_xai_responses=kwargs.get("is_xai_responses") is True,
is_github_responses=kwargs.get("is_github_responses") is True,
replay_encrypted_reasoning=bool(kwargs.get("replay_encrypted_reasoning", True)),
current_issuer_kind=self._resolve_issuer_kind(kwargs),
current_issuer_model=self._last_issuer_model,
native_compaction_eligible=_native_compaction_active(kwargs.get("context_management")),
)
@@ -580,14 +583,17 @@ class ResponsesApiTransport(ProviderTransport):
# Lazy: provider plugins import this transport during model_metadata init.
from agent.model_metadata import strip_codex_context_variant_suffix as _strip_ctx_variant
request_overrides = params.get("request_overrides") or {}
# An override may rewrite the wire model; provenance must be stamped with what actually goes out.
wire_model = _strip_ctx_variant(request_overrides.get("model", model))
kwargs = {
# ``-900k`` picker variants are Hermes-side aliases; the backend knows only the base slug.
"model": _strip_ctx_variant(model),
"model": wire_model,
"instructions": instructions,
"input": self.convert_messages(
payload_messages, is_xai_responses=is_xai_responses, is_github_responses=is_github_responses,
replay_encrypted_reasoning=replay_encrypted_reasoning, base_url=params.get("base_url"),
is_codex_backend=is_codex_backend, context_management=context_management,
is_codex_backend=is_codex_backend, context_management=context_management, model=wire_model,
),
"store": False,
}
@@ -617,8 +623,9 @@ class ResponsesApiTransport(ProviderTransport):
replay_encrypted_reasoning=replay_encrypted_reasoning,
is_xai_responses=is_xai_responses, is_github_responses=is_github_responses,
))
if params.get("request_overrides"):
kwargs.update(params["request_overrides"])
if request_overrides:
kwargs.update(request_overrides)
kwargs["model"] = wire_model
_sanitize_astra_request_kwargs(kwargs, model, params.get("base_url"))
@@ -677,7 +684,8 @@ class ResponsesApiTransport(ProviderTransport):
from agent.codex_responses_adapter import _normalize_codex_response
msg, finish_reason = _normalize_codex_response(
response, issuer_kind=kwargs.get("issuer_kind") or self._last_issuer_kind
response, issuer_kind=kwargs.get("issuer_kind") or self._last_issuer_kind,
issuer_model=kwargs.get("issuer_model") or self._last_issuer_model,
)
tool_calls = None
@@ -555,6 +555,56 @@ def test_chat_messages_to_responses_input_drops_foreign_id_for_codex_backend():
assert xai_message["id"] == _FOREIGN_ITEM_ID
def _reasoning_history(item):
return [
{"role": "assistant", "content": "done", "codex_reasoning_items": [item]},
{"role": "user", "content": "next"},
]
def test_reasoning_replay_requires_matching_issuer_model_on_same_endpoint():
# Blobs are sealed to the minting model, not just the endpoint: same endpoint + other model must drop.
issuer = "other:https://responses.example.com/v1"
normalized, _ = _normalize_codex_response(
SimpleNamespace(
status="completed",
output=[
SimpleNamespace(type="reasoning", id="rs_a", encrypted_content="model-a-blob", summary=[]),
SimpleNamespace(
type="message", role="assistant", status="completed", id="msg_a",
content=[SimpleNamespace(type="output_text", text="done")],
),
],
),
issuer_kind=issuer, issuer_model="gpt-5.6-sol",
)
captured = normalized.codex_reasoning_items[0]
assert captured["_issuer_model"] == "gpt-5.6-sol"
same = _chat_messages_to_responses_input(
_reasoning_history(captured), current_issuer_kind=issuer, current_issuer_model="gpt-5.6-sol"
)
other = _chat_messages_to_responses_input(
_reasoning_history(captured), current_issuer_kind=issuer, current_issuer_model="gpt-5.7-sol"
)
replayed = [i for i in same if i.get("type") == "reasoning"]
assert [i["encrypted_content"] for i in replayed] == ["model-a-blob"]
assert "_issuer_model" not in replayed[0] and "_issuer_kind" not in replayed[0]
assert not any(i.get("type") == "reasoning" for i in other)
def test_reasoning_replay_drops_endpoint_stamped_legacy_item_without_model():
# Fail closed: an endpoint-only stamp cannot prove the minting model once the request model is known.
issuer = "other:https://responses.example.com/v1"
legacy = {"type": "reasoning", "encrypted_content": "legacy-blob", "_issuer_kind": issuer}
items = _chat_messages_to_responses_input(
_reasoning_history(legacy), current_issuer_kind=issuer, current_issuer_model="gpt-5.6-sol"
)
assert not any(i.get("type") == "reasoning" for i in items)
# Ordinary assistant text stays replayable; only the sealed blob is withheld.
assert any(i.get("role") == "assistant" for i in items)
def test_preflight_codex_api_kwargs_drops_oversized_message_id_end_to_end():
kwargs = _preflight_codex_api_kwargs(
{
@@ -921,7 +921,7 @@ class TestCodexBuildKwargs:
monkeypatch.setattr(
"agent.codex_responses_adapter._normalize_codex_response",
lambda resp, issuer_kind=None: (msg, "tool_calls"),
lambda resp, issuer_kind=None, issuer_model=None: (msg, "tool_calls"),
)
normalized = transport.normalize_response(response)
@@ -1103,7 +1103,7 @@ class TestOpencodeReservedToolAliases:
response = SimpleNamespace(output=[], status="completed")
monkeypatch.setattr(
"agent.codex_responses_adapter._normalize_codex_response",
lambda resp, issuer_kind=None: (msg, "tool_calls"),
lambda resp, issuer_kind=None, issuer_model=None: (msg, "tool_calls"),
)
normalized = transport.normalize_response(response)
names = [tc.name for tc in normalized.tool_calls]
@@ -1207,7 +1207,7 @@ class TestXaiReservedToolSearchAlias:
response = SimpleNamespace(output=[], status="completed")
monkeypatch.setattr(
"agent.codex_responses_adapter._normalize_codex_response",
lambda resp, issuer_kind=None: (msg, "tool_calls"),
lambda resp, issuer_kind=None, issuer_model=None: (msg, "tool_calls"),
)
# Pair the response with a real request so provenance is recorded.
transport.build_kwargs(
@@ -1237,7 +1237,7 @@ class TestXaiReservedToolSearchAlias:
response = SimpleNamespace(output=[], status="completed")
monkeypatch.setattr(
"agent.codex_responses_adapter._normalize_codex_response",
lambda resp, issuer_kind=None: (msg, "tool_calls"),
lambda resp, issuer_kind=None, issuer_model=None: (msg, "tool_calls"),
)
return transport.normalize_response(response)