diff --git a/agent/bedrock_adapter.py b/agent/bedrock_adapter.py index 74292e4ccf..f5f65ddb8a 100644 --- a/agent/bedrock_adapter.py +++ b/agent/bedrock_adapter.py @@ -154,7 +154,6 @@ class BedrockOpenAISigV4Auth(httpx.Auth): import botocore.session from botocore.auth import SigV4Auth from botocore.awsrequest import AWSRequest - credentials = botocore.session.get_session().get_credentials() if credentials is None: raise RuntimeError( @@ -215,7 +214,6 @@ _STALE_LIB_MODULE_PREFIXES = ("urllib3.", "botocore.", "boto3.") def _stale_error_types() -> tuple: """botocore + urllib3 transport-failure exception classes (best-effort import).""" import importlib - types: list = [] for module, names in ( ("botocore.exceptions", ("ConnectionError", "HTTPClientError")), @@ -656,19 +654,16 @@ def _assistant_blocks(msg: Dict, content) -> List[Dict]: content_blocks = _replay_ordered_blocks(ordered_blocks) if content_blocks: return content_blocks - content_blocks = [] for detail in (msg.get("reasoning_details") or []): if isinstance(detail, dict) and detail.get("type") == "redacted_thinking": redacted = _decode_redacted(detail.get("data") or detail.get("redactedContentBase64")) if redacted is not None: content_blocks.append({"reasoningContent": {"redactedContent": redacted}}) - if isinstance(content, str) and content.strip(): content_blocks.append({"text": content}) elif isinstance(content, list): content_blocks.extend(_convert_content_to_converse(content)) - for tc in (msg.get("tool_calls", []) or []): fn = tc.get("function", {}) content_blocks.append(_tool_use_block(tc.get("id", ""), fn.get("name", ""), _parse_tool_args(fn.get("arguments", "{}")))) @@ -691,7 +686,6 @@ def convert_messages_to_converse(messages: List[Dict]) -> Tuple[Optional[List[Di converse_msgs[-1]["content"].extend(blocks) else: converse_msgs.append({"role": role, "content": blocks}) - for msg in messages: role = msg.get("role", "") content = msg.get("content") @@ -707,7 +701,6 @@ def convert_messages_to_converse(messages: List[Dict]) -> Tuple[Optional[List[Di append_turn("assistant", _assistant_blocks(msg, content) or [dict(_PLACEHOLDER_BLOCK)]) elif role == "user": append_turn("user", _convert_content_to_converse(content)) - if converse_msgs and converse_msgs[0]["role"] != "user": converse_msgs.insert(0, {"role": "user", "content": [dict(_PLACEHOLDER_BLOCK)]}) if converse_msgs and converse_msgs[-1]["role"] != "user": @@ -881,7 +874,6 @@ def stream_converse_with_callbacks( if encoded: parts.add_redacted(encoded) current_block({"reasoningContent": {}}).setdefault("reasoningContent", {})["redactedContentBase64"] = encoded - for event in event_stream.get("stream", []): if on_event is not None: try: @@ -890,7 +882,6 @@ def stream_converse_with_callbacks( pass if on_interrupt_check and on_interrupt_check(): break - if "contentBlockStart" in event: start_event = event["contentBlockStart"] current_block_index = start_event.get("contentBlockIndex", len(stream_blocks)) @@ -902,7 +893,6 @@ def stream_converse_with_callbacks( stream_blocks[current_block_index] = _tool_use_block(current_tool["toolUseId"], current_tool["name"], {}) if on_tool_start: on_tool_start(current_tool["name"]) - elif "contentBlockDelta" in event: delta = event["contentBlockDelta"].get("delta", {}) if "text" in delta: @@ -916,7 +906,6 @@ def stream_converse_with_callbacks( current_tool["input_json"] += delta["toolUse"].get("input", "") elif "reasoningContent" in delta: on_reasoning(delta["reasoningContent"]) - elif "contentBlockStop" in event: if current_tool is not None: input_dict = _parse_tool_args(current_tool["input_json"]) if current_tool["input_json"] else {} @@ -926,17 +915,14 @@ def stream_converse_with_callbacks( current_tool = None else: flush_text() - elif "messageStop" in event: stop_reason = event["messageStop"].get("stopReason", "end_turn") - elif "metadata" in event: meta_usage = event["metadata"].get("usage", {}) usage_data = { key: meta_usage.get(key, 0) for key in ("inputTokens", "outputTokens", "cacheReadInputTokens", "cacheWriteInputTokens") } - flush_text() return parts.build([stream_blocks[i] for i in sorted(stream_blocks)], usage_data, stop_reason, "") @@ -966,17 +952,13 @@ def build_converse_kwargs( def cache_here(placement: str) -> bool: return cache_enabled and cache_point_allowed(model, placement) - inference_config: Dict[str, Any] = {} if max_tokens is not None: inference_config["maxTokens"] = max_tokens kwargs: Dict[str, Any] = {"modelId": model, "messages": converse_messages, "inferenceConfig": inference_config} - if system_prompt: kwargs["system"] = system_prompt + [dict(_CACHE_POINT)] if cache_here("system") else system_prompt - from agent.anthropic_adapter import _forbids_sampling_params - if not _forbids_sampling_params(model): if temperature is not None: inference_config["temperature"] = temperature @@ -984,7 +966,6 @@ def build_converse_kwargs( inference_config["topP"] = top_p if stop_sequences: inference_config["stopSequences"] = stop_sequences - converse_tools = convert_tools_to_converse(tools) if tools else [] if converse_tools: # Non-tool-calling models (e.g. DeepSeek R1) reject toolConfig with a @@ -998,12 +979,10 @@ def build_converse_kwargs( "Model %s does not support tool calling — tools stripped. " "The agent will operate in text-only mode.", model ) - if cache_here("messages") and len(converse_messages) >= 2: content = converse_messages[-2].get("content") if isinstance(content, list) and content: content.append(dict(_CACHE_POINT)) - if guardrail_config: kwargs["guardrailConfig"] = guardrail_config if not inference_config: @@ -1100,7 +1079,6 @@ def _list_inference_profiles(client, filter_set: set, models: List[Dict[str, Any next_token = response.get("nextToken") if not next_token: break - seen_ids = {m["id"].lower() for m in models} for profile in profiles: profile_id = (profile.get("inferenceProfileId") or "").strip() @@ -1123,18 +1101,15 @@ def discover_bedrock_models(region: str, provider_filter: Optional[List[str]] = Returns [] when the client cannot be built. """ import time - cache_key = f"{region}:{','.join(sorted(provider_filter or []))}" cached = _discovery_cache.get(cache_key) if cached and (time.time() - cached["timestamp"]) < _DISCOVERY_CACHE_TTL_SECONDS: return cached["models"] - try: client = _get_bedrock_control_client(region) except Exception as e: logger.warning("Failed to create Bedrock client for model discovery: %s", e) return [] - models: List[Dict[str, Any]] = [] filter_set = {f.lower() for f in (provider_filter or [])} try: @@ -1145,7 +1120,6 @@ def discover_bedrock_models(region: str, provider_filter: Optional[List[str]] = _list_inference_profiles(client, filter_set, models) except Exception as e: logger.debug("Skipping inference profile discovery: %s", e) - models.sort(key=lambda m: (0 if m["id"].startswith("global.") else 1, m["name"].lower())) _discovery_cache[cache_key] = {"timestamp": time.time(), "models": models} return models @@ -1220,7 +1194,6 @@ def probe_bedrock_context_length(model_id: str, region: str) -> Optional[int]: except Exception as exc: # boto3 missing / credential resolution failure logger.debug("Bedrock context probe skipped for %s: %s", model_id, exc) return None - last_error = "" for tier_tokens in _BEDROCK_PROBE_TIERS: oversized = "data " * int(tier_tokens / _WORDS_PER_TOKEN) @@ -1242,7 +1215,6 @@ def probe_bedrock_context_length(model_id: str, region: str) -> Optional[int]: logger.info("Probed Bedrock context window for %s: %s tokens", model_id, f"{limit:,}") return limit # Opaque server error / auth / throttle at this tier — try the next. - logger.debug("Bedrock context probe for %s returned no parseable limit: %s", model_id, last_error[:200]) return None diff --git a/agent/codex_headers.py b/agent/codex_headers.py index bec0c8fb0e..d967ae2e8a 100644 --- a/agent/codex_headers.py +++ b/agent/codex_headers.py @@ -43,7 +43,6 @@ def codex_cloudflare_headers(access_token: str, *, base_url: str = CODEX_AUX_BAS """ if is_official_codex_base_url(base_url): from hermes_cli import __version__ - headers = {"User-Agent": f"HermesAgent/{__version__}", "originator": "hermes-agent"} else: headers = {"User-Agent": "codex_cli_rs/0.0.0 (Hermes Agent)", "originator": "codex_cli_rs"} diff --git a/agent/codex_responses_adapter.py b/agent/codex_responses_adapter.py index da3e23c774..8f55793c16 100644 --- a/agent/codex_responses_adapter.py +++ b/agent/codex_responses_adapter.py @@ -126,11 +126,9 @@ def _neutralize_harmony_tokens(text: str) -> str: """Keep Harmony source readable without emitting reserved wire tokens.""" if not text or "<" not in text or "|" not in text: return text - replacement = rf"<{_FULLWIDTH_PIPE}\1{_FULLWIDTH_PIPE}>" if not any(unicodedata.category(char) == "Cf" for char in text): return _HARMONY_CONTROL_TOKEN_RE.sub(replacement, text) - # The backend strips Unicode format controls (e.g. U+200B) before its # reserved-token check, so match on the visible text and rewrite the # original spans — any Cf-hidden variant is neutralized the same way. @@ -324,7 +322,6 @@ def _derive_responses_function_call_id(call_id: str, response_item_id: Optional[ """Build a valid Responses `function_call.id` (must start with `fc_`).""" if isinstance(response_item_id, str) and response_item_id.strip().startswith("fc_"): return response_item_id.strip() - source = (call_id or "").strip() sanitized = re.sub(r"[^A-Za-z0-9_-]", "", source) for candidate in (source, sanitized): @@ -494,7 +491,6 @@ def _tool_output_item(msg: Dict[str, Any]) -> Optional[Dict[str, Any]]: call_id = raw_tool_call_id.strip() if not _nonblank(call_id): return None - # ``output`` may be a string or an ``input_text``/``input_image`` array. tool_content = msg.get("content") output_value: Any = ( @@ -546,7 +542,6 @@ def _chat_messages_to_responses_input( def emit(new_items: List[Dict[str, Any]], msg: Dict[str, Any]) -> None: items.extend(new_items) item_sources.extend([msg] * len(new_items)) - for msg in messages: if not isinstance(msg, dict): continue @@ -558,7 +553,6 @@ def _chat_messages_to_responses_input( continue if role not in {"user", "assistant"}: continue - content = msg.get("content", "") content_parts = _chat_content_to_responses_parts(content, role=role) # [] unless a list if isinstance(content, list): @@ -566,11 +560,9 @@ def _chat_messages_to_responses_input( content_text = "".join(p["text"] for p in content_parts if p["type"] == text_type) else: content_text = _str_or_empty(content) - if role == "user": emit([{"role": role, "content": content_parts or content_text}], msg) 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, @@ -578,7 +570,6 @@ def _chat_messages_to_responses_input( emit(reasoning_items, msg) message_items = _replay_message_items(msg, is_github_responses=is_github_responses) emit(message_items, msg) - if not message_items: if content_parts: emit([{"role": "assistant", "content": content_parts}], msg) @@ -587,9 +578,7 @@ def _chat_messages_to_responses_input( elif reasoning_items: # Every reasoning item needs a following item (else missing_following_item). emit([{"role": "assistant", "content": ""}], msg) - emit(_replay_tool_call_items(msg, start_index=len(items)), msg) - # Native server-side compaction renders nothing placed before a compaction # item, so pre-checkpoint history is dead upload weight and the user's # plaintext asks / merged local summaries silently vanish. Keep the newest @@ -597,9 +586,7 @@ def _chat_messages_to_responses_input( # messages within a token budget, leave the tail untouched. if not native_compaction_eligible: return items - from agent.native_compaction import prune_pre_checkpoint_items - return prune_pre_checkpoint_items(items, item_sources=item_sources) @@ -622,7 +609,6 @@ def classify_responses_route(agent: Any) -> ResponsesRouteFlags: ``https://evil.com/models.github.ai`` must not classify as GitHub. """ from utils import base_url_hostname - provider = getattr(agent, "provider", None) base_url = str(getattr(agent, "base_url", "") or "") hostname = str(getattr(agent, "_base_url_hostname", "") or "").lower() or base_url_hostname(base_url) @@ -630,7 +616,6 @@ def classify_responses_route(agent: Any) -> ResponsesRouteFlags: def _host_is(domain: str) -> bool: return hostname == domain or hostname.endswith("." + domain) - return ResponsesRouteFlags( is_codex_backend=provider == "openai-codex" or (_host_is("chatgpt.com") and "/backend-api/codex" in lower), is_xai_responses=provider in {"xai", "xai-oauth"} or hostname == "api.x.ai", @@ -654,17 +639,13 @@ def estimate_native_responses_preflight_tokens( """ if getattr(agent, "api_mode", None) != "codex_responses" or not isinstance(messages, list): return None - is_codex_backend, is_xai_responses, is_github_responses = classify_responses_route(agent) - from agent.native_compaction import native_compaction_context_management - if not native_compaction_context_management( agent, is_codex_backend=is_codex_backend, is_xai_responses=is_xai_responses, is_github_responses=is_github_responses, ): return None - try: items = _chat_messages_to_responses_input( messages, is_xai_responses=is_xai_responses, is_github_responses=is_github_responses, @@ -680,9 +661,7 @@ def estimate_native_responses_preflight_tokens( return None if not isinstance(items, list): return None - from agent.model_metadata import estimate_request_tokens_rough - return estimate_request_tokens_rough(items, system_prompt=system_prompt or "", tools=tools) @@ -781,7 +760,6 @@ def _preflight_role_message(item: Dict[str, Any], idx: int, role: str, ctx: _Pre content = item.get("content", "") if not isinstance(content, list): return {"role": role, "content": ctx.sanitize_text(_str_or_empty(content))} - # Parts are already Responses-shaped; validate and re-type text for the role. # Unlike history conversion, empty text / empty image urls are kept, not dropped. text_type = _text_type_for(role) @@ -825,7 +803,6 @@ def _preflight_codex_input_items( ) -> List[Dict[str, Any]]: if not isinstance(raw_items, list): raise ValueError("Codex Responses input must be a list of input items.") - ctx = _PreflightCtx( sanitize_text=_neutralize_harmony_tokens if sanitize_harmony_tokens else (lambda text: text), sanitize_harmony_tokens=sanitize_harmony_tokens, @@ -906,19 +883,15 @@ def _preflight_codex_api_kwargs( ) -> Dict[str, Any]: if not isinstance(api_kwargs, dict): raise ValueError("Codex Responses request must be a dict.") - missing = [key for key in ("model", "instructions", "input") if key not in api_kwargs] if missing: raise ValueError(f"Codex Responses request missing required field(s): {', '.join(sorted(missing))}.") - model = api_kwargs.get("model") if not _nonblank(model): raise ValueError("Codex Responses request 'model' must be a non-empty string.") - instructions = _str_or_empty(api_kwargs.get("instructions")).strip() or DEFAULT_AGENT_IDENTITY if sanitize_harmony_tokens: instructions = _neutralize_harmony_tokens(instructions) - normalized: Dict[str, Any] = { "model": model.strip(), "instructions": instructions, @@ -929,7 +902,6 @@ def _preflight_codex_api_kwargs( ), "store": False, } - tools = api_kwargs.get("tools") if tools is not None: if not isinstance(tools, list): @@ -938,15 +910,12 @@ def _preflight_codex_api_kwargs( if sanitize_harmony_tokens: normalized_tools = _neutralize_harmony_structure(normalized_tools) normalized["tools"] = normalized_tools - if api_kwargs.get("store", False) is not False: raise ValueError("Codex Responses contract requires 'store' to be false.") - for key, accept, coerce in _PREFLIGHT_OPTIONAL_FIELDS: value = api_kwargs.get(key) if accept(value): normalized[key] = coerce(value) if coerce else value - extra_headers = api_kwargs.get("extra_headers") if extra_headers is not None: if not isinstance(extra_headers, dict): @@ -956,7 +925,6 @@ def _preflight_codex_api_kwargs( normalized_headers = {key.strip(): str(value) for key, value in extra_headers.items() if value is not None} if normalized_headers: normalized["extra_headers"] = normalized_headers - extra_body = api_kwargs.get("extra_body") if extra_body is not None: if not isinstance(extra_body, dict): @@ -965,7 +933,6 @@ def _preflight_codex_api_kwargs( # the SDK serializes extra_body without per-field checks. if extra_body: normalized["extra_body"] = dict(extra_body) - allowed_keys = set(_PREFLIGHT_ALLOWED_KEYS) if allow_stream: stream = api_kwargs.get("stream") @@ -976,7 +943,6 @@ def _preflight_codex_api_kwargs( allowed_keys.add("stream") elif "stream" in api_kwargs: raise ValueError("Codex Responses stream flag is only allowed in fallback streaming requests.") - # Defense-in-depth slash-enum strip for xAI (rejects ``Qwen/Qwen3.5`` style # enum values). Gated on the model name because native Codex accepts slashes. is_xai_model = str(api_kwargs.get("model") or "").lower().startswith(("grok-", "x-ai/grok-")) @@ -986,7 +952,6 @@ def _preflight_codex_api_kwargs( normalized["tools"], _ = strip_slash_enum(normalized["tools"]) except Exception: pass # Best-effort — the caller-level sanitization should have handled it - unexpected = sorted(key for key in api_kwargs if key not in allowed_keys) if unexpected: raise ValueError(f"Codex Responses request has unsupported field(s): {', '.join(unexpected)}.") @@ -1025,7 +990,6 @@ def _format_responses_error(error_obj: Any, response_status: str) -> str: def field(name: str) -> str: value = _field(error_obj, name) return str(value).strip() if isinstance(value, str) or value else "" - code_str, message_str = field("code"), field("message") if code_str and message_str: return f"{code_str}: {message_str}" @@ -1113,7 +1077,6 @@ class _OutputScan: if item_status in _INCOMPLETE_STATUSES and item_type not in _SERVER_SIDE_TOOL_CALL_TYPES: self.has_incomplete_items = True self.saw_streaming_or_item_incomplete = True - if item_type == "message": self._message(item, item_status) elif item_type == "reasoning": @@ -1170,7 +1133,6 @@ def _normalize_codex_response(response: Any, *, issuer_kind: Optional[str] = Non response_incomplete_content_filter = ( response_status == "incomplete" and str(incomplete_reason or "").strip().lower() == "content_filter" ) - output = getattr(response, "output", None) if not isinstance(output, list) or not output: # Codex can deliver the whole answer via stream events and return an @@ -1189,20 +1151,16 @@ def _normalize_codex_response(response: Any, *, issuer_kind: Optional[str] = Non else: raise RuntimeError("Responses API returned no output items") response.output = output - 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) tool_calls, reasoning_parts = scan.tool_calls, scan.reasoning_parts - final_text = "\n".join(scan.content_parts).strip() if not final_text and hasattr(response, "output_text") and (scan.saw_final_answer_phase or not scan.saw_commentary_phase): out_text = getattr(response, "output_text", "") if isinstance(out_text, str): final_text = out_text.strip() - # Tool-call leak recovery: gpt-5.x sometimes emits the intended # ``function_call`` as plain Harmony text (``to=functions.foo {json}``) with # no structured item. Treat as incomplete so the continuation path @@ -1216,7 +1174,6 @@ def _normalize_codex_response(response: Any, *, issuer_kind: Optional[str] = Non "Leaked snippet: %r", final_text[:300], ) final_text = "" - # Reasoning-channel answer salvage (xAI grok): grok-4.x sometimes puts the # final answer inside the reasoning item after its ```` delimiter. # Without salvage the reasoning-only rule marks the turn incomplete, and since @@ -1235,7 +1192,6 @@ def _normalize_codex_response(response: Any, *, issuer_kind: Optional[str] = Non final_text = salvaged reasoning_prefix = joined_reasoning[:marker].strip() reasoning_parts = [reasoning_prefix] if reasoning_prefix else [] - assistant_message = SimpleNamespace( content=final_text, tool_calls=tool_calls, @@ -1245,7 +1201,6 @@ def _normalize_codex_response(response: Any, *, issuer_kind: Optional[str] = Non codex_reasoning_items=scan.reasoning_items_raw or None, codex_message_items=scan.message_items_raw or None, ) - if tool_calls: finish_reason = "tool_calls" elif response_incomplete_content_filter: diff --git a/agent/codex_runtime.py b/agent/codex_runtime.py index 9fe119694b..d59530785b 100644 --- a/agent/codex_runtime.py +++ b/agent/codex_runtime.py @@ -41,7 +41,6 @@ def _codex_request_failure_details(error: BaseException) -> tuple[int | None, st exception_classes: list[str] = [] current: BaseException | None = error seen: set[int] = set() - while current is not None and id(current) not in seen and len(seen) < 8: seen.add(id(current)) exception_classes.append(type(current).__name__) @@ -102,7 +101,6 @@ def _record_codex_app_server_usage(agent, turn) -> dict[str, Any]: no usage still counts as one API call for session/status accounting. """ agent.session_api_calls += 1 - usage = getattr(turn, "token_usage_last", None) compressor = getattr(agent, "context_compressor", None) if not isinstance(usage, dict) or not usage: @@ -118,9 +116,7 @@ def _record_codex_app_server_usage(agent, turn) -> dict[str, Any]: ), ) return {} - from agent.usage_pricing import CanonicalUsage, estimate_usage_cost - canonical_usage = CanonicalUsage( input_tokens=_coerce_usage_int(usage.get("inputTokens")), output_tokens=_coerce_usage_int(usage.get("outputTokens")), @@ -139,7 +135,6 @@ def _record_codex_app_server_usage(agent, turn) -> dict[str, Any]: "prompt_tokens": prompt_tokens, "completion_tokens": canonical_usage.output_tokens, "total_tokens": total_tokens, **token_counts, } - if compressor is not None: try: compressor.update_from_response(usage_dict) @@ -148,10 +143,8 @@ def _record_codex_app_server_usage(agent, turn) -> dict[str, Any]: compressor.context_length = context_window except Exception: logger.debug("codex app-server usage update failed", exc_info=True) - for key, value in usage_dict.items(): setattr(agent, f"session_{key}", getattr(agent, f"session_{key}") + value) - cost_result = estimate_usage_cost( agent.model, canonical_usage, provider=agent.provider, base_url=agent.base_url, api_key=getattr(agent, "api_key", ""), @@ -161,7 +154,6 @@ def _record_codex_app_server_usage(agent, turn) -> dict[str, Any]: agent.session_estimated_cost_usd += cost_usd agent.session_cost_status, agent.session_cost_source = cost_result.status, cost_result.source cost_fields = {"estimated_cost_usd": cost_usd, "cost_status": cost_result.status, "cost_source": cost_result.source} - _queue_token_counts( agent, "Codex app-server token persistence failed (session=%s, tokens=%d): %s", total_tokens, counts=lambda: dict( @@ -182,7 +174,6 @@ def _record_codex_app_server_compaction(agent, turn, *, approx_tokens: int | Non """ if not force and not getattr(turn, "compacted", False): return False - thread_id = getattr(turn, "thread_id", None) or "" turn_id = getattr(turn, "turn_id", None) or "" logger.info( @@ -195,7 +186,6 @@ def _record_codex_app_server_compaction(agent, turn, *, approx_tokens: int | Non agent._emit_status(COMPACTION_STATUS) except Exception: pass - compressor = getattr(agent, "context_compressor", None) if compressor is not None: compressor.compression_count = getattr(compressor, "compression_count", 0) + 1 @@ -212,7 +202,6 @@ def _record_codex_app_server_compaction(agent, turn, *, approx_tokens: int | Non compressor.last_prompt_tokens = -1 compressor.last_completion_tokens = 0 compressor.awaiting_real_usage_after_compression = True - # Provider-side context was rewritten; the usage anchor's transcript snapshot no longer matches. agent._usage_anchor = None agent._turn_base_usage_anchor = None @@ -333,7 +322,6 @@ def _codex_item_completion_payload(item: dict) -> tuple[str, bool]: def _stable_call_id(item: dict, name: str) -> str: """Deterministic tool_call id mirroring CodexEventProjector (live TUI card correlates with projected history).""" from agent.transports.codex_event_projector import _deterministic_call_id - item_type = item.get("type") or "" tool = item.get("tool") or "unknown" if item_type == "mcpToolCall": @@ -414,7 +402,6 @@ def make_codex_app_server_event_bridge(agent) -> Callable[[dict], None]: (_fire_tool_completed if completed else _fire_tool_started)(item) elif completed and item_type == "agentMessage": _fire_agent_message_completed(item) - handlers: dict[str, Callable[[dict], None]] = { "item/agentMessage/delta": lambda p: _fire_delta(p, "_fire_stream_delta"), "item/reasoning/delta": lambda p: _fire_delta(p, "_fire_reasoning_delta"), @@ -428,7 +415,6 @@ def make_codex_app_server_event_bridge(agent) -> Callable[[dict], None]: if handler is not None: params = note.get("params") or {} handler(params if isinstance(params, dict) else {}) - return on_event @@ -461,7 +447,6 @@ def _ensure_codex_session(agent) -> None: return from agent.runtime_cwd import resolve_agent_cwd from agent.transports.codex_app_server_session import CodexAppServerSession, _ServerRequestRouting - # Approval callback: Hermes' standard prompt flow when a CLI thread installed one. try: from tools.terminal_tool import _get_approval_callback @@ -478,7 +463,6 @@ def _ensure_codex_session(agent) -> None: auto_approve_requests = is_approval_bypass_active() except Exception: logger.debug("codex app-server: approval-bypass lookup failed; keeping fail-closed default", exc_info=True) - agent._codex_session = CodexAppServerSession( cwd=getattr(agent, "session_cwd", None) or str(resolve_agent_cwd()), approval_callback=approval_callback, request_routing=_ServerRequestRouting(auto_approve_exec=auto_approve_requests, auto_approve_apply_patch=auto_approve_requests), @@ -498,10 +482,8 @@ def _persist_projected_messages(agent, turn, messages: List[Dict[str, Any]]) -> if not turn.projected_messages: return from agent.message_metadata import append_message - for projected_message in turn.projected_messages: append_message(messages, projected_message) - if getattr(agent, "_session_db", None) is None: return try: @@ -528,14 +510,12 @@ def _finish_codex_turn( agent._iters_since_skill = getattr(agent, "_iters_since_skill", 0) + turn.tool_iterations _record_codex_app_server_compaction(agent, turn) usage_result = _record_codex_app_server_usage(agent, turn) - # Skill nudge check AFTER iters were incremented (same as chat_completions). should_review_skills = ( 0 < agent._skill_nudge_interval <= agent._iters_since_skill and "skill_manage" in agent.valid_tool_names ) if should_review_skills: agent._iters_since_skill = 0 - # External memory sync skipped on interrupt/error (no partial transcripts). if not turn.interrupted and turn.error is None: try: @@ -545,7 +525,6 @@ def _finish_codex_turn( ) except Exception: logger.debug("external memory sync raised", exc_info=True) - # Background review fork: only when a trigger tripped AND a real final response exists. if turn.final_text and not turn.interrupted and (should_review_memory or should_review_skills): try: @@ -554,7 +533,6 @@ def _finish_codex_turn( ) except Exception: logger.debug("background review spawn raised", exc_info=True) - return usage_result @@ -581,9 +559,7 @@ def run_codex_app_server_turn( "codex_app_server owns the authoritative thread and compacts it " "without a truthful pre-compaction transcript boundary" ) - _ensure_codex_session(agent) - try: turn = agent._codex_session.run_turn(user_input=user_message) except Exception as exc: @@ -593,20 +569,16 @@ def run_codex_app_server_turn( _consume_user_interrupt(agent), messages, api_calls=0, completed=False, error=str(exc), final_response=f"Codex app-server turn failed: {exc}. Fall back to default runtime with `/codex-runtime auto`.", ) - interrupt = _consume_user_interrupt(agent, turn.interrupted) - # Wedged client (deadline blown, watchdog tripped, OAuth refresh died, # subprocess exited): retire the session so the next turn respawns codex. if getattr(turn, "should_retire", False): logger.warning("codex app-server session retired (turn error: %s)", turn.error) _close_codex_session(agent) - _persist_projected_messages(agent, turn, messages) usage_result = _finish_codex_turn( agent, turn, messages, original_user_message=original_user_message, should_review_memory=should_review_memory, ) - return _turn_result( interrupt, messages, api_calls=1, completed=not turn.interrupted and turn.error is None, error=turn.error, final_response=turn.final_text, @@ -659,13 +631,11 @@ def _raise_stream_error(event: Any) -> None: ``run_agent`` is imported lazily to keep this module importable standalone. """ from run_agent import _StreamErrorEvent - nested = _event_field(event, "error") def _error_field(name: str) -> Any: value = _event_field(event, name) return _event_field(nested, name) if value is None and nested is not None else value - raw_message = _error_field("message") message = (str(raw_message) if raw_message is not None else "stream emitted error event").strip() or "stream emitted error event" raise _StreamErrorEvent(message, code=_error_field("code"), param=_error_field("param")) @@ -879,7 +849,6 @@ class _CodexResponseAssembler: # executable; malformed non-empty JSON passes through untouched. arguments=(pending["arguments"] or "").strip() or "{}", ))) - # output_index is optional and a partial ordering over mixed indexed/unindexed # entries is ill-defined: protocol order only when every entry has an index, else wire order. if all(entry[0] is not None for entry in indexed): @@ -898,17 +867,14 @@ class _CodexResponseAssembler: if not output and self.text_deltas and not self.has_tool_calls: content = [SimpleNamespace(type="output_text", text="".join(self.text_deltas))] output = [SimpleNamespace(type="message", role="assistant", status="completed", content=content)] - # Done items stay authoritative; settlement only fills the gap left by # backends that omit per-item done events on a successful completion. if self.pending_function_calls and self.saw_response_completed: output = self._settled_output() - # No terminal frame AND no usable content = truncated / rejected stream, # distinct from "completed with empty body" (what the SDK helper raised as RuntimeError). if not self.saw_terminal and not output: raise RuntimeError("Codex Responses stream did not emit a terminal response") - return SimpleNamespace( output=output, output_text="".join(self.text_deltas), usage=self.terminal_usage, status=self.terminal_status, id=self.terminal_response_id, model=self.model, incomplete_details=self.terminal_incomplete_details, @@ -1021,7 +987,6 @@ def _bypass_sdk_request_transform(stream_kwargs: dict) -> dict: """ if os.environ.get("HERMES_CODEX_SDK_TRANSFORM", "").strip().lower() in {"1", "true", "yes", "on"}: return stream_kwargs - moved = { field: stream_kwargs[field] for field in _SDK_TRANSFORM_BYPASS_FIELDS @@ -1029,7 +994,6 @@ def _bypass_sdk_request_transform(stream_kwargs: dict) -> dict: } if not moved: return stream_kwargs - bypassed = {key: value for key, value in stream_kwargs.items() if key not in moved} extra_body = bypassed.get("extra_body") merged = dict(extra_body) if isinstance(extra_body, dict) else {} @@ -1049,9 +1013,7 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta """ import httpx as _httpx from openai import APIConnectionError as _APIConnectionError - from agent import relay_llm - transport_errors = (_httpx.RemoteProtocolError, _httpx.ReadTimeout, _httpx.ConnectError, ConnectionError) active_client = client or agent._ensure_primary_openai_client(reason="codex_stream_direct") max_stream_retries = 1 @@ -1144,7 +1106,6 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta # primary client, which is never reuse-cached and must not be force-shut. if client is not None: agent._abort_request_openai_client(active_client, reason="codex_stream_close_failed") - on_commentary_message = ( _fenced(lambda text: agent._fire_streamed_codex_commentary(text)) if getattr(agent, "interim_assistant_callback", None) is not None and getattr(agent, "show_commentary", True) @@ -1155,11 +1116,9 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta else "fallback" if int(getattr(agent, "_fallback_index", 0) or 0) > 0 else "primary" ) - for attempt in range(max_stream_retries + 1): if agent._interrupt_requested: raise InterruptedError("Agent interrupted before Codex stream retry") - intercepted_events: list = [] writer_token["value"] = None event_stream = None @@ -1209,10 +1168,8 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta except _APIConnectionError as exc: _log_failure(exc) raise - if not agent._interrupt_requested: _drain_for_finalizer(event_stream) - if final.status in {"incomplete", "failed"}: logger.warning( "Codex Responses stream terminal status=%s "