diff --git a/agent/auxiliary_client.py b/agent/auxiliary_client.py index 8d293f02de..73cb48e51c 100644 --- a/agent/auxiliary_client.py +++ b/agent/auxiliary_client.py @@ -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 diff --git a/agent/codex_responses_adapter.py b/agent/codex_responses_adapter.py index b16e066c27..ffac5e20fa 100644 --- a/agent/codex_responses_adapter.py +++ b/agent/codex_responses_adapter.py @@ -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): diff --git a/agent/error_classifier.py b/agent/error_classifier.py index 5d10d05d80..d0277efe69 100644 --- a/agent/error_classifier.py +++ b/agent/error_classifier.py @@ -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 ): diff --git a/agent/transports/codex.py b/agent/transports/codex.py index a7c1bd0ab4..ac0076dc37 100644 --- a/agent/transports/codex.py +++ b/agent/transports/codex.py @@ -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 diff --git a/tests/agent/test_codex_responses_adapter.py b/tests/agent/test_codex_responses_adapter.py index df011f9516..61894dc517 100644 --- a/tests/agent/test_codex_responses_adapter.py +++ b/tests/agent/test_codex_responses_adapter.py @@ -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( { diff --git a/tests/agent/transports/test_codex_transport.py b/tests/agent/transports/test_codex_transport.py index 1f7eebea38..49eab87e82 100644 --- a/tests/agent/transports/test_codex_transport.py +++ b/tests/agent/transports/test_codex_transport.py @@ -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)