diff --git a/agent/codex_runtime.py b/agent/codex_runtime.py index a1b3e8fc12..1750b63640 100644 --- a/agent/codex_runtime.py +++ b/agent/codex_runtime.py @@ -31,10 +31,8 @@ def _call_guarded(fn: Callable | None, fail_msg: str, *fail_args: Any, args: tup def _codex_request_failure_details(error: BaseException) -> tuple[int | None, str]: - """(serialized request bytes, exception class chain) for a failed request. - - OpenAI connection exceptions retain the final ``httpx.Request``; its buffered - content gives the exact byte count without logging payloads or URLs.""" + """(serialized request bytes, exception class chain); the buffered ``httpx.Request`` content + on OpenAI connection errors gives the exact byte count without logging payloads or URLs.""" request_body_bytes: int | None = None exception_classes: list[str] = [] current: BaseException | None = error @@ -78,10 +76,8 @@ def _queue_token_counts(agent, fail_msg: str, *fail_extra: Any, counts: Callable def _record_codex_app_server_usage(agent, turn) -> dict[str, Any]: - """Translate Codex app-server token usage into Hermes accounting. - - Prompt bucket = uncached + cached input; the protocol exposes no cache-write - tokens. A turn with no usage still counts as one API call.""" + """Translate Codex app-server token usage into Hermes accounting. Prompt bucket = uncached + cached + input (the protocol exposes no cache-write tokens); a turn with no usage still counts as one API call.""" agent.session_api_calls += 1 usage = getattr(turn, "token_usage_last", None) compressor = getattr(agent, "context_compressor", None) @@ -129,9 +125,8 @@ def _record_codex_app_server_usage(agent, turn) -> dict[str, Any]: 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: billing( - **token_counts, **cost_fields, billing_mode="subscription_included" if cost_result.status == "included" else None, - ), + counts=lambda: billing(**token_counts, **cost_fields, + billing_mode="subscription_included" if cost_result.status == "included" else None), ) return {**usage_dict, "last_prompt_tokens": prompt_tokens, **cost_fields} @@ -179,20 +174,15 @@ def _record_codex_app_server_compaction(agent, turn, *, approx_tokens: int | Non # --- Codex app-server → Hermes UI bridge ------------------------------------- -# The app-server bypasses the Hermes tool loop, so the bridge translates JSON-RPC -# notifications into the callbacks the standard runtime fires -# (tool_progress_callback, _fire_stream_delta, _emit_interim_assistant_message). +# The app-server bypasses the Hermes tool loop, so the bridge translates JSON-RPC notifications +# into the callbacks the standard runtime fires (tool_progress_callback, _fire_stream_delta, ...). -# Item types that project to a Hermes tool_call (keep in sync with -# agent/transports/codex_event_projector.py so UI names match recorded names). -# webSearch is codex's built-in tool: no projector entry, still gets a bubble. +# Item types that project to a Hermes tool_call (keep in sync with agent/transports/codex_event_projector.py +# so UI names match recorded names). webSearch is codex's built-in tool: no projector entry, still gets a bubble. _CODEX_TOOL_ITEM_TYPES = frozenset({"commandExecution", "fileChange", "mcpToolCall", "dynamicToolCall", "webSearch"}) - -# Internal MCP server wrapping Hermes' native tools: its inner dispatch has no -# tool_progress_callback, so the codex-level mcpToolCall IS the display event and -# the mcp.hermes-tools.* prefix is stripped (users think of these as Hermes tools). +# Internal MCP server wrapping Hermes' native tools: its inner dispatch has no tool_progress_callback, so the +# codex-level mcpToolCall IS the display event and the mcp.hermes-tools.* prefix is stripped (users see Hermes tools). _INTERNAL_MCP_SERVER = "hermes-tools" - _STATIC_TOOL_NAMES = {"commandExecution": "exec_command", "fileChange": "apply_patch", "webSearch": "web_search"} _STABLE_ID_PREFIXES = {"commandExecution": "exec", "fileChange": "apply_patch"} _MCP_LIKE_ITEM_TYPES = {"mcpToolCall", "dynamicToolCall"} @@ -282,10 +272,9 @@ def make_codex_app_server_event_bridge(agent) -> Callable[[dict], None]: """Build the ``on_event`` callback for ``CodexAppServerSession(on_event=...)``. Tool items fire ``tool_progress_callback`` plus the stable-ID ``tool_start_callback`` / - ``tool_complete_callback`` card hooks; deltas go to ``_fire_stream_delta`` / - ``_fire_reasoning_delta``; a completed agentMessage goes to - ``_emit_interim_assistant_message`` (the gateway's ``already_streamed`` check dedupes - against streamed deltas). Every callback is guarded so a buggy display hook cannot + ``tool_complete_callback`` card hooks; deltas go to ``_fire_stream_delta`` / ``_fire_reasoning_delta``; + a completed agentMessage goes to ``_emit_interim_assistant_message`` (the gateway's ``already_streamed`` + check dedupes against streamed deltas). Every callback is guarded so a buggy display hook cannot tear down the turn loop.""" # item_id -> (tool_name, args, started_monotonic); duration even when codex omits durationMs. started: dict[str, tuple[str, dict, float]] = {} @@ -307,11 +296,11 @@ def make_codex_app_server_event_bridge(agent) -> Callable[[dict], None]: def _fire_tool_completed(item: dict) -> None: name = _codex_item_to_tool_name(item) prior = started.pop(item_id, None) if (item_id := item.get("id") or "") else None - # Prefer codex's durationMs; else our started timestamp; else None - # (some codex versions only emit completed for fast items). + # Prefer codex's durationMs; else our started timestamp; else None (some codex + # versions only emit completed for fast items). codex_ms = item.get("durationMs") has_codex_ms = isinstance(codex_ms, (int, float)) and codex_ms >= 0 - duration: Any = codex_ms / 1000.0 if has_codex_ms else (time.monotonic() - prior[2] if prior is not None else None) + duration: Any = codex_ms / 1000.0 if has_codex_ms else (time.monotonic() - prior[2] if prior else None) result, is_error = _codex_item_completion_payload(item) agent_cb("tool_progress_callback", "tool_progress_callback raised on tool.completed for %s", name, args=("tool.completed", name, None, None), @@ -327,8 +316,7 @@ def make_codex_app_server_event_bridge(agent) -> Callable[[dict], None]: def _fire_agent_message_completed(item: dict) -> None: text = item.get("text") or "" - # display.show_commentary=false keeps mid-turn narration off the interim - # path here too (same contract as codex_responses commentary). + # display.show_commentary=false keeps mid-turn narration off the interim path too (codex_responses contract). if isinstance(text, str) and text.strip() and getattr(agent, "show_commentary", True): agent_cb("_emit_interim_assistant_message", "_emit_interim_assistant_message raised", args=({"role": "assistant", "content": text},)) @@ -390,9 +378,8 @@ def _ensure_codex_session(agent) -> None: approval_callback = _get_approval_callback() except Exception: approval_callback = None - # Gateway/cron have no UI for codex approval requests, so exec/apply_patch fail - # closed by default. Only an explicit approval bypass (approvals.mode: off, /yolo, - # --yolo, HERMES_YOLO_MODE) hands policy to codex's own sandbox profile. + # Gateway/cron have no UI for codex approval requests, so exec/apply_patch fail closed by default. Only an + # explicit approval bypass (approvals.mode: off, /yolo, --yolo, HERMES_YOLO_MODE) hands policy to codex's sandbox. auto_approve_requests = False try: from tools.approval import is_approval_bypass_active @@ -409,9 +396,9 @@ def _ensure_codex_session(agent) -> None: def _persist_projected_messages(agent, turn, messages: List[Dict[str, Any]]) -> None: """Splice the projected messages into ``messages`` and flush them to the session DB. - Bypasses conversation_loop's per-step _persist_session(); the flush dedups via - _DB_PERSISTED_MARKER so only the new codex rows are written. The agent stays the - sole persister (agent_persisted=True): a gateway re-write would re-INSERT the user turn.""" + Bypasses conversation_loop's per-step _persist_session(); the flush dedups via _DB_PERSISTED_MARKER so + only the new codex rows are written. The agent stays the sole persister (agent_persisted=True): a + gateway re-write would re-INSERT the user turn.""" if not turn.projected_messages: return from agent.message_metadata import append_message @@ -437,8 +424,7 @@ def _finish_codex_turn( agent, turn, messages: List[Dict[str, Any]], *, original_user_message: Any, should_review_memory: bool, ) -> dict[str, Any]: """Post-turn bookkeeping mirroring the chat_completions loop; returns usage fields.""" - # run_conversation() already bumped _turns_since_memory / _user_turn_count; only - # _iters_since_skill (per tool iteration in the bypassed loop) is ours. + # run_conversation() already bumped _turns_since_memory / _user_turn_count; only _iters_since_skill is ours. 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) @@ -465,11 +451,9 @@ def run_codex_app_server_turn( should_review_memory: bool = False, ) -> Dict[str, Any]: """Hand the turn to a ``codex app-server`` subprocess and project its events into ``messages``. - - Returns the chat_completions result shape. The user message is ALREADY in - ``messages`` — never append it again.""" - # Defense in depth for compression.checkpoint_required: agent init refuses the - # combination, but api_mode is mutable. Explicit-True check matches compress_context(). + Returns the chat_completions result shape. The user message is ALREADY in ``messages`` — never append it again.""" + # Defense in depth for compression.checkpoint_required: agent init refuses the combination, but + # api_mode is mutable. Explicit-True check matches compress_context(). if getattr(agent, "compression_checkpoint_required", False) is True: from agent.conversation_compression import _checkpoint_blocked raise _checkpoint_blocked( @@ -487,8 +471,7 @@ def run_codex_app_server_turn( 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. + # Wedged client (deadline blown, watchdog tripped, OAuth refresh died, subprocess exited): retire it. if getattr(turn, "should_retire", False): logger.warning("codex app-server session retired (turn error: %s)", turn.error) _close_codex_session(agent) @@ -519,10 +502,9 @@ def _turn_result( # --- Event-driven Responses streaming ----------------------------------------- -# The SDK's ``responses.stream(...)`` helper rebuilds a typed Response from -# ``response.completed.response.output`` and crashes when it is null. We consume raw -# ``responses.create(stream=True)`` SSE events and assemble the final response from -# ``output_item.done``, so the terminal ``output`` may be null / [] / a string / absent. +# The SDK's ``responses.stream(...)`` helper rebuilds a typed Response from ``response.completed.response.output`` +# and crashes when it is null. We consume raw ``responses.create(stream=True)`` SSE events and assemble the final +# response from ``output_item.done``, so the terminal ``output`` may be null / [] / a string / absent. def _event_field(event: Any, name: str, default: Any = None) -> Any: @@ -534,10 +516,8 @@ def _event_field(event: Any, name: str, default: Any = None) -> Any: def _raise_stream_error(event: Any) -> None: - """Raise ``_StreamErrorEvent`` from a ``type=error`` SSE frame. - - The spec puts code/message/param at the top level, but the SDK and several - proxies nest them under ``error``; read top-level first, then the envelope.""" + """Raise ``_StreamErrorEvent`` from a ``type=error`` SSE frame. The spec puts code/message/param at the + top level, but the SDK and several proxies nest them under ``error``; read top-level first, then the envelope.""" from run_agent import _StreamErrorEvent nested = _event_field(event, "error") @@ -566,10 +546,9 @@ def _output_text_of(item: Any) -> str: class _CodexResponseAssembler: """Assemble a Response-shaped ``SimpleNamespace`` from raw Responses SSE events. - Only ``usage`` / ``status`` / ``id`` are read from the terminal frame — never - ``response.output``. Output items come from ``output_item.done``, or are synthesized - from text deltas, or settled from function calls announced via ``output_item.added`` - but never confirmed (some backends omit per-item done events on success).""" + Only ``usage`` / ``status`` / ``id`` are read from the terminal frame — never ``response.output``. Output + items come from ``output_item.done``, or are synthesized from text deltas, or settled from function calls + announced via ``output_item.added`` but never confirmed (some backends omit per-item done events on success).""" has_tool_calls = first_delta_fired = saw_terminal = False next_output_sequence = 0 @@ -578,21 +557,19 @@ class _CodexResponseAssembler: active_summary_index: Any = None terminal_status: str = "completed" terminal_usage = terminal_response_id = terminal_incomplete_details = terminal_error = None - # terminal_status defaults to "completed", so settlement needs an explicitly - # observed response.completed frame (not EOF/interrupt). + # terminal_status defaults to "completed", so settlement needs an explicitly observed response.completed frame. saw_response_completed = False def __init__(self, *, model, on_text_delta, on_reasoning_delta, on_commentary_message, on_first_delta): self.model, self.on_text_delta, self.on_reasoning_delta = model, on_text_delta, on_reasoning_delta self.on_commentary_message, self.on_first_delta = on_commentary_message, on_first_delta self.output_items: List[Any] = [] - # output_index / first-observed sequence per output item, in lockstep, so - # settled pending calls merge back in stream order. + # output_index / first-observed sequence per output item, in lockstep, so settled pending calls merge + # back in stream order. self.output_indexes, self.output_sequences = [], [] self.text_deltas, self.commentary_text_deltas = [], [] - # pending_function_calls: announced-but-unconfirmed function calls keyed by item id. - # announced_output_order: first-observed (sequence, output_index) per announced item id - # so a later .done keeps its announced position when merged with settled calls. + # pending_function_calls: announced-but-unconfirmed function calls keyed by item id. announced_output_order: + # first-observed (sequence, output_index) per announced item id so a later .done keeps its announced position. self.pending_function_calls: Dict[str, Dict[str, Any]] = {} self.announced_output_order: Dict[str, tuple] = {} @@ -605,8 +582,8 @@ class _CodexResponseAssembler: self.active_message_phase = _message_phase(item) if item_type == "message" else None if self.active_message_phase == "commentary": self.commentary_text_deltas = [] - # Record first-observed ordering for EVERY announced item; .done must reuse it or a - # mixed announced/pending stream without output_index values reorders the calls. + # Record first-observed ordering for EVERY announced item; .done must reuse it or a mixed + # announced/pending stream without output_index values reorders the calls. item_id = str(_event_field(item, "id", "")) if item_id and item_id not in self.announced_output_order: self.announced_output_order[item_id] = (self.next_output_sequence, _event_field(event, "output_index")) @@ -624,8 +601,8 @@ class _CodexResponseAssembler: delta_text = _event_field(event, "delta", "") if not delta_text: return - # Harmony commentary/analysis text is mid-turn narration, never the final - # answer: route to the reasoning callback, keep only the item for replay. + # Harmony commentary/analysis text is mid-turn narration, never the final answer: route to the + # reasoning callback, keep only the item for replay. if self.active_message_phase == "commentary": self.commentary_text_deltas.append(delta_text) # Legacy fallback when no first-class commentary consumer is installed. @@ -650,8 +627,8 @@ class _CodexResponseAssembler: if "delta" in event_type: pending["arguments"] += _event_field(event, "delta", "") or "" elif event_type.endswith("function_call_arguments.done"): - # Authoritative for the accumulated string; an explicit "" (zero-arg - # call) counts, only a missing field keeps the streamed deltas. + # Authoritative for the accumulated string; an explicit "" (zero-arg call) counts, only a + # missing field keeps the streamed deltas. if (done_args := _event_field(event, "arguments", None)) is not None: pending["arguments"] = str(done_args) @@ -671,8 +648,8 @@ class _CodexResponseAssembler: if done_item is None: return self.output_items.append(done_item) - # Reuse the announced position when known (fresh tail sequence for unannounced - # items); the .done event's own output_index wins over the announced one. + # Reuse the announced position when known (fresh tail sequence for unannounced items); the .done + # event's own output_index wins over the announced one. done_id = str(_event_field(done_item, "id", "")) announced_sequence, announced_index = self.announced_output_order.get(done_id, (None, None)) if announced_sequence is None: @@ -730,8 +707,8 @@ class _CodexResponseAssembler: indexed.append((pending.get("output_index"), pending["sequence"], SimpleNamespace( type="function_call", id=_event_field(item, "id", None), call_id=_event_field(item, "call_id", None), name=_event_field(item, "name", None), status="completed", - # Empty/whitespace arguments become "{}" so zero-delta calls stay - # executable; malformed non-empty JSON passes through untouched. + # Empty/whitespace arguments become "{}" so zero-delta calls stay executable; malformed + # non-empty JSON passes through untouched. arguments=(pending["arguments"] or "").strip() or "{}", ))) # output_index is optional: protocol order only when every entry has one, else wire order. @@ -748,8 +725,8 @@ 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. + # 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. @@ -766,17 +743,16 @@ def _consume_codex_event_stream( event_iter: Any, *, model: str, on_text_delta=None, on_reasoning_delta=None, on_commentary_message=None, on_first_delta=None, on_event=None, interrupt_check=None, ) -> SimpleNamespace: - """Consume a Codex Responses SSE stream into a Response-shaped ``SimpleNamespace`` - (see :class:`_CodexResponseAssembler`; ``status`` is ``completed`` when the stream - ended with content but no terminal frame; ``model`` comes from kwargs). + """Consume a Codex Responses SSE stream into a Response-shaped ``SimpleNamespace`` (see + :class:`_CodexResponseAssembler`; ``status`` is ``completed`` when the stream ended with content but no + terminal frame; ``model`` comes from kwargs). - Callbacks: ``on_text_delta`` per output_text delta, suppressed once a function_call - is seen; ``on_reasoning_delta`` for reasoning and ``phase=analysis`` deltas (also - commentary without a commentary callback); ``on_commentary_message`` once per completed - ``phase=commentary`` message, before any following tool item; ``on_first_delta`` - one-shot; ``on_event`` every event before any processing; ``interrupt_check()`` True - breaks the loop and may raise ``TimeoutError`` / ``InterruptedError`` for request - retirement that must not become a partial final response.""" + Callbacks: ``on_text_delta`` per output_text delta, suppressed once a function_call is seen; + ``on_reasoning_delta`` for reasoning and ``phase=analysis`` deltas (also commentary without a commentary + callback); ``on_commentary_message`` once per completed ``phase=commentary`` message, before any following + tool item; ``on_first_delta`` one-shot; ``on_event`` every event before any processing; ``interrupt_check()`` + True breaks the loop and may raise ``TimeoutError`` / ``InterruptedError`` for request retirement that + must not become a partial final response.""" assembler = _CodexResponseAssembler( model=model, on_text_delta=on_text_delta, on_reasoning_delta=on_reasoning_delta, on_commentary_message=on_commentary_message, on_first_delta=on_first_delta, @@ -795,9 +771,9 @@ def _consume_codex_event_stream( def _sanitize_consumer_codex_request(agent: Any, request: dict[str, Any]) -> dict[str, Any]: - """Drop fields the ChatGPT OAuth Codex endpoint rejects, at the final wire boundary - (after Relay / middleware / ``request_overrides``): a late ``prompt_cache_retention``, - top-level or nested in ``extra_body``, would otherwise HTTP 400 a valid follow-up.""" + """Drop fields the ChatGPT OAuth Codex endpoint rejects, at the final wire boundary (after Relay / + middleware / ``request_overrides``): a late ``prompt_cache_retention``, top-level or nested in + ``extra_body``, would otherwise HTTP 400 a valid follow-up.""" sanitized = dict(request) # getattr: run_codex_stream is also driven with stand-in agents carrying only the attrs a path needs. backend_predicate = getattr(agent, "_is_codex_backend", None) @@ -822,14 +798,12 @@ def _sanitize_consumer_codex_request(agent: Any, request: dict[str, Any]) -> dic return sanitized -# Bulk request fields carrying the conversation payload; the rest is scalar -# config the SDK transform handles in microseconds. +# Bulk request fields carrying the conversation payload; the rest is scalar config the SDK transform handles fast. _SDK_TRANSFORM_BYPASS_FIELDS = ("input", "tools") def _is_plain_json_data(value: Any) -> bool: - """True when ``value`` is purely JSON wire types; anything else (pydantic models, - generators) must keep the typed SDK path.""" + """True when ``value`` is purely JSON wire types; pydantic models / generators must keep the typed SDK path.""" if value is None or isinstance(value, (str, int, float, bool)): return True if isinstance(value, dict): @@ -842,11 +816,10 @@ def _is_plain_json_data(value: Any) -> bool: def _bypass_sdk_request_transform(stream_kwargs: dict) -> dict: """Route bulk payload fields around the SDK's ``maybe_transform``. - ``responses.create`` re-walks the whole body against the ResponseCreateParams union - with the GIL held — multi-MB conversations can wedge for hours, pre-network, where no - watchdog socket kill helps. The SDK merges ``extra_body`` AFTER the transform, so - moving wire-format bulk fields there yields a byte-identical request without the - walk. HERMES_CODEX_SDK_TRANSFORM=1 disables.""" + ``responses.create`` re-walks the whole body against the ResponseCreateParams union with the GIL held — + multi-MB conversations can wedge for hours, pre-network, where no watchdog socket kill helps. The SDK + merges ``extra_body`` AFTER the transform, so moving wire-format bulk fields there yields a byte-identical + request without the walk. HERMES_CODEX_SDK_TRANSFORM=1 disables.""" if os.environ.get("HERMES_CODEX_SDK_TRANSFORM", "").strip().lower() in {"1", "true", "yes", "on"}: return stream_kwargs moved = {f: stream_kwargs[f] for f in _SDK_TRANSFORM_BYPASS_FIELDS @@ -862,8 +835,7 @@ def _bypass_sdk_request_transform(stream_kwargs: dict) -> dict: def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta=None): - """Execute one streaming Responses API request (raw ``responses.create(stream=True)`` - events, see module notes) and return the final response.""" + """One streaming Responses API request over raw ``responses.create(stream=True)`` events.""" import httpx as _httpx from openai import APIConnectionError as _APIConnectionError from agent import relay_llm @@ -872,9 +844,9 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta max_stream_retries, model = 1, api_kwargs.get("model") # Accumulate streamed text so callers / compat shims can read it. agent._codex_streamed_text_parts: list = [] - # Retirement token for THIS request (installed by ``interruptible_api_call``). A - # watchdog that kills the connection clears the agent-level token, so a worker still - # draining frames can tell it was retired. ``None`` = no watchdog; every check passes. + # Retirement token for THIS request (installed by ``interruptible_api_call``). A watchdog that kills the + # connection clears the agent-level token, so a worker still draining frames can tell it was retired. + # ``None`` = no watchdog; every check passes. request_token = getattr(agent, "_active_codex_stream_request_token", None) # Delta-sink claim for the CURRENT physical attempt (None until the stream opens). writer_token = {"value": None} @@ -895,8 +867,8 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta agent._touch_activity("receiving stream response") def _interrupt_or_superseded() -> bool: - # A retired request must NOT break out of the consume loop (that returns a partial - # ``final`` with status "completed"); raise so the watchdog's TimeoutError is seen. + # A retired request must NOT break out of the consume loop (that returns a partial ``final`` with + # status "completed"); raise so the watchdog's TimeoutError is seen. if not _request_is_current(): raise TimeoutError("Codex Responses stream request retired before terminal response") return bool(agent._interrupt_requested) @@ -931,8 +903,8 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta return False def _drain_for_finalizer(event_stream: Any) -> None: - # ``final`` is already assembled; draining only lets Relay run its finalizer. A - # transport error here must NOT discard the completed, already-billed response. + # ``final`` is already assembled; draining only lets Relay run its finalizer. A transport error + # here must NOT discard the completed, already-billed response. try: for _ignored in event_stream: pass @@ -951,12 +923,13 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta if callable(close_fn): close_fn() except Exception: - # A failed close can leave this connection checked out of the httpx pool while - # the caller reuse-caches the client; poison the slot so close really closes the - # pool. ``client is None`` is the shared primary client — never force-shut. + # A failed close can leave this connection checked out of the httpx pool while the caller + # reuse-caches the client; poison the slot so close really closes the pool. ``client is None`` + # is the shared primary client — never force-shut. if client is not None: agent._abort_request_openai_client(active_client, reason="codex_stream_close_failed") - wants_commentary = getattr(agent, "interim_assistant_callback", None) is not None and getattr(agent, "show_commentary", True) + show_commentary = getattr(agent, "show_commentary", True) + wants_commentary = getattr(agent, "interim_assistant_callback", None) is not None and show_commentary on_commentary_message = _fenced(lambda text: agent._fire_streamed_codex_commentary(text)) if wants_commentary else None call_role = ("delegated" if getattr(agent, "is_subagent", False) else "fallback" if int(getattr(agent, "_fallback_index", 0) or 0) > 0 else "primary")