Files
hermes-agent/agent/codex_runtime.py
Teknium 2776813df3 compat(plugins): temporary import-path shims for external plugins — ONE commit, revert on schedule
The Sep 2026 decomposition (PR #102117) makes internal import paths a non-API: names now live in
the focused modules that define them. This commit is the ONLY thing keeping the old paths alive,
so external plugins have time to update. It is deliberately a single, unsquashed commit:

    git revert <this sha>

removes every shim, stub and manifest at once on the announced date. Nothing in-tree may depend on
these pointers: scripts/check_compat_pointers.py (wired into lint.yml) fails CI if it does.

What it adds (see COMPAT_MANIFEST.md, compat_manifest.json):
- 332 facade modules get one delimited `PLUGIN-COMPAT` block appended at the end of the file
- 1,172 moved names resolved lazily via a module `__getattr__` (PEP 562) — never a top-level import,
  so no import cycles; facades that already had `__getattr__` get a chained one
- 592 third-party/stdlib names the old modules used to expose, with their original import statements
- 266 public definitions that had been deleted as unused, restored byte-for-byte from the pre-decomposition
  tree (+40 private helpers and 16 imports pulled in only because a restored definition needs them)
- 3 deleted modules recreated as re-export stubs (gateway/startup_watchdog, hermes_cli/observability/
  relay_runtime, tools/environments/modal_utils)
- private names (`_x`) get no pointer: they were never API (3,792 skipped)

Verified: all 335 touched modules import under a fresh HERMES_HOME and every manifest name resolves;
the lint reports zero in-tree uses; ruff clean; targeted suites unchanged.
2026-09-03 17:13:22 -07:00

990 lines
56 KiB
Python

"""Codex API runtime — App Server and Responses-API streaming paths. Every entry point takes the parent
AIAgent first: ``run_codex_app_server_turn`` drives one ``codex app-server`` subprocess turn;
``run_codex_stream`` runs one streaming Codex Responses call."""
from __future__ import annotations
import json
import logging
import os
import time
from contextlib import suppress
from types import SimpleNamespace
from typing import Any, Callable, Dict, List
from agent.stream_single_writer import claim_stream_writer, stream_writer_is_current
logger = logging.getLogger(__name__)
def _call_guarded(fn: Callable | None, fail_msg: str, *fail_args: Any, args: tuple = (), kwargs: dict | None = None):
"""Invoke an optional display/debug callback; a buggy hook must never tear down the turn."""
if fn is None:
return
try:
fn(*args, **(kwargs or {}))
except Exception:
logger.debug(fail_msg, *fail_args, exc_info=True)
def _codex_request_failure_details(error: BaseException) -> tuple[int | None, str]:
"""(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
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__)
if request_body_bytes is None:
content = None
with suppress(Exception):
content = getattr(getattr(current, "request", None), "content", None)
if isinstance(content, str):
request_body_bytes = len(content.encode("utf-8"))
elif isinstance(content, (bytes, bytearray, memoryview)):
request_body_bytes = len(content)
implicit_chain = current.__cause__ is None and not current.__suppress_context__
current = current.__context__ if implicit_chain else current.__cause__
return request_body_bytes, " <- ".join(exception_classes)
def _coerce_usage_int(value: Any) -> int:
if isinstance(value, bool):
return 0
if isinstance(value, int):
return max(value, 0)
if isinstance(value, float):
return max(int(value), 0)
if isinstance(value, str):
# Only the str->int parse is guarded; a float NaN still raises like it always has.
with suppress(ValueError):
return max(int(value), 0)
return 0
def _queue_token_counts(agent, fail_msg: str, *fail_extra: Any, counts: Callable[[], dict]) -> None:
"""Enqueue per-call accounting for the SessionDB background writer. ``counts`` is built
lazily inside the guarded try so a stub agent without a session DB is never touched."""
if not (agent._session_db and agent.session_id):
return
try:
if not agent._session_db_created:
agent._ensure_db_session()
agent._session_db.queue_token_counts(agent.session_id, **counts())
except Exception as exc:
logger.debug(fail_msg, agent.session_id, *fail_extra, exc)
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."""
agent.session_api_calls += 1
usage = getattr(turn, "token_usage_last", None)
compressor = getattr(agent, "context_compressor", None)
def billing(**extra):
return dict(model=agent.model, billing_provider=agent.provider, billing_base_url=agent.base_url, api_call_count=1, **extra)
if not isinstance(usage, dict) or not usage:
if compressor is not None and getattr(compressor, "awaiting_real_usage_after_compression", False):
# No usage cannot adjudicate the pending compaction; unlatch preflight deferral.
compressor.update_from_response({})
_queue_token_counts(agent, "Codex app-server api-call persistence failed (session=%s): %s",
counts=lambda: billing(billing_mode="subscription_included"))
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")),
cache_read_tokens=_coerce_usage_int(usage.get("cachedInputTokens")), cache_write_tokens=0,
reasoning_tokens=_coerce_usage_int(usage.get("reasoningOutputTokens")), raw_usage=usage,
)
prompt_tokens = canonical_usage.prompt_tokens
total_tokens = _coerce_usage_int(usage.get("totalTokens")) or canonical_usage.total_tokens
token_counts = {f: getattr(canonical_usage, f) for f in
("input_tokens", "output_tokens", "cache_read_tokens", "cache_write_tokens", "reasoning_tokens")}
usage_dict = {"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)
context_window = getattr(turn, "model_context_window", None)
if isinstance(context_window, int) and context_window > 0:
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", ""),
)
cost_usd = float(cost_result.amount_usd) if cost_result.amount_usd is not None else None
if cost_usd is not None:
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: 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}
def _record_codex_app_server_compaction(agent, turn, *, approx_tokens: int | None = None, force: bool = False) -> bool:
"""Record a Codex-native compaction boundary: the app-server owns the compacted thread,
so local transcript rows are NOT rewritten — only session event/usage counters."""
if not force and not getattr(turn, "compacted", False):
return False
thread_id, turn_id = getattr(turn, "thread_id", None) or "", getattr(turn, "turn_id", None) or ""
logger.info("codex app-server compaction observed: session=%s thread=%s turn=%s force=%s",
getattr(agent, "session_id", None) or "none", thread_id, turn_id, force)
if not force:
with suppress(Exception):
from agent.conversation_compression import COMPACTION_STATUS
agent._emit_status(COMPACTION_STATUS)
compressor = getattr(agent, "context_compressor", None)
if compressor is not None:
compressor.compression_count = getattr(compressor, "compression_count", 0) + 1
compressor.last_compression_rough_tokens = approx_tokens or 0
# Codex owns this summary: a prior Hermes deterministic-fallback flag must not leak into it.
record_boundary = getattr(type(compressor), "record_completed_compaction", None)
if callable(record_boundary):
record_boundary(compressor, used_fallback=False)
elif hasattr(compressor, "_verify_compaction_cleared_threshold"):
compressor._verify_compaction_cleared_threshold = True
if not getattr(turn, "token_usage_last", None):
compressor.last_prompt_tokens, compressor.last_completion_tokens = -1, 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
agent._last_compaction_in_place = False
_call_guarded(getattr(agent, "event_callback", None) or None, "event_callback error on codex session:compress",
args=("session:compress", {
"platform": getattr(agent, "platform", None) or "", "session_id": getattr(agent, "session_id", None) or "",
"old_session_id": "", "in_place": False,
"compression_count": getattr(compressor, "compression_count", 0) if compressor is not None else 0,
"runtime": "codex_app_server", "thread_id": thread_id, "turn_id": turn_id,
}))
return True
# --- 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, ...).
# 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 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"}
# Item types whose preview is the first 120 chars of one string field.
_PREVIEW_FIELDS = {"commandExecution": "command", "webSearch": "query"}
def _item_changes(item: dict) -> list[dict]:
return [c for c in (item.get("changes") or []) if isinstance(c, dict)]
def _codex_item_to_tool_name(item: dict) -> str:
"""Synthetic Hermes tool name for a codex item (mirrors CodexEventProjector)."""
item_type = item.get("type") or ""
if item_type == "mcpToolCall":
server, tool = item.get("server") or "mcp", item.get("tool") or "unknown"
return tool if server == _INTERNAL_MCP_SERVER else f"mcp.{server}.{tool}"
if item_type == "dynamicToolCall":
return item.get("tool") or "dynamic"
return _STATIC_TOOL_NAMES.get(item_type) or item_type or "unknown"
def _codex_item_to_args(item: dict) -> dict:
"""Args dict for tool_progress_callback("tool.started"); mirrors the projector shapes."""
item_type = item.get("type") or ""
if item_type == "commandExecution":
return {"command": item.get("command") or "", "cwd": item.get("cwd") or ""}
if item_type == "fileChange":
return {"changes": [{"kind": (c.get("kind") or {}).get("type") or "update", "path": c.get("path") or ""}
for c in _item_changes(item)]}
if item_type in _MCP_LIKE_ITEM_TYPES:
args = item.get("arguments") or {}
return args if isinstance(args, dict) else {"arguments": args}
return {"query": item.get("query") or ""} if item_type == "webSearch" else {}
def _codex_item_to_preview(item: dict) -> Any:
"""Short preview for the tool.started bubble; None when nothing useful (UI tolerates None)."""
item_type = item.get("type") or ""
if item_type in _PREVIEW_FIELDS:
return (item.get(_PREVIEW_FIELDS[item_type]) or "")[:120] or None
if item_type == "fileChange":
paths = [c.get("path") for c in _item_changes(item) if c.get("path")]
return (", ".join(paths[:3]) + (f", +{len(paths) - 3} more" if len(paths) > 3 else "")) if paths else None
if item_type in _MCP_LIKE_ITEM_TYPES:
args = item.get("arguments") or {}
if isinstance(args, dict) and args:
with suppress(TypeError, ValueError):
return json.dumps(args, ensure_ascii=False)[:120]
return None
def _codex_item_completion_payload(item: dict) -> tuple[str, bool]:
"""(result_text, is_error) for a completed tool item; mirrors the projector's tool-result content."""
item_type = item.get("type") or ""
if item_type == "commandExecution":
out, exit_code = item.get("aggregatedOutput") or "", item.get("exitCode")
is_error = bool(exit_code is not None and exit_code != 0)
return (f"[exit {exit_code}]\n{out}" if is_error else out), is_error
if item_type == "fileChange":
status = item.get("status") or "unknown"
n = len(item.get("changes") or [])
return f"apply_patch status={status}, {n} change(s)", status not in {"completed", "applied", "success"}
if item_type == "mcpToolCall":
if error := item.get("error"):
return f"[error] {json.dumps(error, ensure_ascii=False)[:1000]}", True
result = item.get("result")
return (json.dumps(result, ensure_ascii=False)[:4000] if result is not None else ""), False
if item_type == "dynamicToolCall":
content_items, success = item.get("contentItems") or [], item.get("success", True)
has_items = isinstance(content_items, list) and content_items
return (json.dumps(content_items, ensure_ascii=False)[:4000] if has_items else f"success={success}"), not bool(success)
return "", False
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"
prefix = {"mcpToolCall": f"mcp__{item.get('server') or 'mcp'}__{tool}", "dynamicToolCall": f"dyn_{tool}"}.get(item_type)
return _deterministic_call_id(prefix or _STABLE_ID_PREFIXES.get(item_type, name), item.get("id") or "")
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
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]] = {}
def agent_cb(attr: str, fail_msg: str, *fail_args: Any, args: tuple = (), kwargs: dict | None = None) -> None:
_call_guarded(getattr(agent, attr, None), fail_msg, *fail_args, args=args, kwargs=kwargs)
def _fire_tool_started(item: dict) -> None:
item_id, name = item.get("id") or "", _codex_item_to_tool_name(item)
args = _codex_item_to_args(item)
if item_id:
started[item_id] = (name, args, time.monotonic())
agent_cb("tool_progress_callback", "tool_progress_callback raised on tool.started for %s", name,
args=("tool.started", name, _codex_item_to_preview(item), args))
# Stable-ID tool card (TUI/desktop) fires alongside the progress bubble.
agent_cb("tool_start_callback", "tool_start_callback raised for %s", name,
args=(_stable_call_id(item, name), name, args))
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).
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 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),
kwargs={"duration": duration, "is_error": is_error, "result": result})
args = prior[1] if prior is not None else _codex_item_to_args(item)
agent_cb("tool_complete_callback", "tool_complete_callback raised for %s", name,
args=(_stable_call_id(item, name), name, args, result))
def _fire_delta(params: dict, attr: str) -> None:
text = params.get("delta") or params.get("text") or ""
# Single-writer guard (#65991): a superseded stream must not pollute the turn's accumulated text
# (which also feeds the interim-visible-text de-dup comparison), even when a caller reaches this
# directly (the tool-suppressed content path) rather than through _fire_stream_delta.
if isinstance(text, str) and text:
agent_cb(attr, f"{attr} raised", args=(text,))
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 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},))
def _on_item(params: dict, completed: bool) -> None:
item = params.get("item")
if not isinstance(item, dict):
return
item_type = item.get("type") or ""
if item_type in _CODEX_TOOL_ITEM_TYPES:
(_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"),
"item/reasoning/summaryDelta": lambda p: _fire_delta(p, "_fire_reasoning_delta"),
"item/started": lambda p: _on_item(p, completed=False), "item/completed": lambda p: _on_item(p, completed=True),
}
def on_event(note: dict) -> None:
handler = handlers.get(note.get("method") or "") if isinstance(note, dict) else None
if handler is not None:
params = note.get("params")
handler(params if isinstance(params, dict) else {})
return on_event
# --- Codex app-server turn ----------------------------------------------------
def _close_codex_session(agent) -> None:
"""Drop the session so the next turn respawns codex instead of reusing a dead client."""
with suppress(Exception):
agent._codex_session.close()
agent._codex_session = None
def _consume_user_interrupt(agent, active: bool = True) -> tuple[bool, Any]:
"""(user_interrupted, interrupt_message); clears the agent-level interrupt so a hard
stop cannot poison the next turn (mirrors the conversation-loop finalizer)."""
interrupted = bool(active and getattr(agent, "_interrupt_requested", False))
message = getattr(agent, "_interrupt_message", None) if interrupted else None
if interrupted:
agent.clear_interrupt()
return interrupted, message
def _ensure_codex_session(agent) -> None:
"""Lazily spawn one CodexAppServerSession per AIAgent (reused across turns, closed by the _cleanup hook)."""
if getattr(agent, "_codex_session", None) is not 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.
approval_callback = None
with suppress(Exception):
from tools.terminal_tool import _get_approval_callback
approval_callback = _get_approval_callback()
# 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
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)
# Bridge codex JSON-RPC notifications (item/started, item/completed, item/agentMessage/delta, ...) into
# Hermes' gateway UI callbacks (tool_progress_callback, _fire_stream_delta,
# _emit_interim_assistant_message). Without this, Discord/Telegram users see no live tool-progress or
# interim commentary while codex_app_server is running — only the final answer (#33200). Supersedes the
# narrower item/started-only bridge from #38835.
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),
on_event=make_codex_app_server_event_bridge(agent),
)
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."""
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
flush_ok = False
try:
flush_ok = agent._flush_messages_to_session_db(messages)
except Exception:
logger.warning("codex app-server projected-message flush failed", exc_info=True)
if flush_ok is False:
# Output already streamed and agent_persisted cannot flip to False: surface the gap loudly.
logger.warning("codex app-server turn was delivered but could NOT be persisted to the session DB "
"(session=%s) — this turn will be missing after restart/resume", getattr(agent, "session_id", None))
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 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)
# 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:
_call_guarded(getattr(agent, "_sync_external_memory_for_turn", None), "external memory sync raised", kwargs=dict(
original_user_message=original_user_message, final_response=turn.final_text, interrupted=False, messages=messages,
))
# 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):
_call_guarded(getattr(agent, "_spawn_background_review", None), "background review spawn raised", kwargs=dict(
messages_snapshot=list(messages), review_memory=should_review_memory, review_skills=should_review_skills,
))
return usage_result
def run_codex_app_server_turn(agent, *, user_message: str, original_user_message: Any, messages: List[Dict[str, Any]],
effective_task_id: str, 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().
if getattr(agent, "compression_checkpoint_required", False) is True:
from agent.conversation_compression import _checkpoint_blocked
raise _checkpoint_blocked("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:
logger.exception("codex app-server turn failed")
_close_codex_session(agent)
return _turn_result(
_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 it.
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,
# We flushed the projected rows ourselves (agent_persisted); the gateway must skip its own DB write.
final_response=turn.final_text, agent_persisted=True, codex_thread_id=turn.thread_id, codex_turn_id=turn.turn_id,
**usage_result,
)
def _turn_result(interrupt: tuple[bool, Any], messages: List[Dict[str, Any]], *, api_calls: int, completed: bool,
error: Any, final_response: Any, **extra: Any) -> Dict[str, Any]:
"""Result shape shared with the chat_completions path (``partial`` == ``not completed``)."""
user_interrupted, interrupt_message = interrupt
return {
"final_response": final_response, "messages": messages, "api_calls": api_calls,
"completed": completed, "partial": not completed, "interrupted": user_interrupted,
**({"interrupt_message": interrupt_message} if interrupt_message else {}),
"error": error, **extra,
}
# --- 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.
def _event_field(event: Any, name: str, default: Any = None) -> Any:
"""Field access for attr-style (SDK objects) and dict (raw JSON) events/items."""
value = getattr(event, name, None)
if value is None and isinstance(event, dict):
value = event.get(name, default)
return value if value is not None else default
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."""
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"))
def _message_phase(item: Any) -> str | None:
phase = _event_field(item, "phase", None)
return phase.strip().lower() if isinstance(phase, str) else None
def _output_text_of(item: Any) -> str:
"""Concatenated ``output_text`` parts of a message item ("" if content is not a list)."""
content_parts = _event_field(item, "content", [])
parts = content_parts if isinstance(content_parts, list) else []
return "".join(
str(_event_field(part, "text", "") or "") for part in parts if _event_field(part, "type", "") == "output_text"
).strip()
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)."""
has_tool_calls = first_delta_fired = saw_terminal = False
next_output_sequence = 0
active_message_phase: str | None = None
# Reasoning summary parts carry no separator; a summary_index change is where the blank line belongs.
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.
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.
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.
self.pending_function_calls: Dict[str, Dict[str, Any]] = {}
self.announced_output_order: Dict[str, tuple] = {}
def _safe(self, cb: Callable | None, label: str, *args: Any) -> None:
_call_guarded(cb, f"Codex stream {label} raised", args=args)
def _on_item_added(self, event: Any, event_type: str) -> None:
item = _event_field(event, "item")
item_type = _event_field(item, "type", "")
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.
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"))
self.next_output_sequence += 1
if "function_call" in str(item_type):
self.has_tool_calls = True
if item_id:
announced_sequence, announced_index = self.announced_output_order[item_id]
self.pending_function_calls[item_id] = {
"item": item, "arguments": str(_event_field(item, "arguments", "") or ""),
"output_index": announced_index, "sequence": announced_sequence,
}
def _on_text_delta(self, event: Any, event_type: str) -> None:
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.
if self.active_message_phase == "commentary":
self.commentary_text_deltas.append(delta_text)
# Legacy fallback when no first-class commentary consumer is installed.
if self.on_commentary_message is None:
self._safe(self.on_reasoning_delta, "on_reasoning_delta", delta_text)
elif self.active_message_phase == "analysis":
self._safe(self.on_reasoning_delta, "on_reasoning_delta", delta_text)
else:
self.text_deltas.append(delta_text)
if self.has_tool_calls:
return
if not self.first_delta_fired:
self.first_delta_fired = True
self._safe(self.on_first_delta, "on_first_delta")
self._safe(self.on_text_delta, "on_text_delta", delta_text)
def _on_function_call(self, event: Any, event_type: str) -> None:
self.has_tool_calls = True
pending = self.pending_function_calls.get(str(_event_field(event, "item_id", "")))
if pending is None:
return # the item itself lands on output_item.done
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.
if (done_args := _event_field(event, "arguments", None)) is not None:
pending["arguments"] = str(done_args)
def _on_reasoning_delta(self, event: Any, event_type: str) -> None:
reasoning_text = _event_field(event, "delta", "")
if not reasoning_text or self.on_reasoning_delta is None:
return
summary_index = _event_field(event, "summary_index")
if summary_index is not None:
if self.active_summary_index is not None and summary_index != self.active_summary_index:
reasoning_text = f"\n\n{reasoning_text}"
self.active_summary_index = summary_index
self._safe(self.on_reasoning_delta, "on_reasoning_delta", reasoning_text)
def _on_item_done(self, event: Any, event_type: str) -> None:
done_item = _event_field(event, "item")
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.
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:
announced_sequence, self.next_output_sequence = self.next_output_sequence, self.next_output_sequence + 1
self.output_indexes.append(_event_field(event, "output_index", announced_index))
self.output_sequences.append(announced_sequence)
# Confirmed by the authoritative done event; never settle it twice.
self.pending_function_calls.pop(done_id, None)
if _message_phase(done_item) == "commentary" and self.on_commentary_message is not None:
commentary_text = "".join(self.commentary_text_deltas).strip() or _output_text_of(done_item)
if commentary_text:
self._safe(self.on_commentary_message, "on_commentary_message", commentary_text)
self.commentary_text_deltas = []
def _on_terminal(self, event: Any, event_type: str) -> bool:
self.saw_terminal = True
resp_obj = _event_field(event, "response")
if resp_obj is not None:
self.terminal_usage, self.terminal_response_id = _event_field(resp_obj, "usage"), _event_field(resp_obj, "id")
rstatus = _event_field(resp_obj, "status")
if isinstance(rstatus, str):
self.terminal_status = rstatus
if event_type == "response.incomplete":
self.terminal_incomplete_details = _event_field(resp_obj, "incomplete_details")
elif event_type == "response.failed":
self.terminal_error = _event_field(resp_obj, "error")
self.saw_response_completed = self.saw_response_completed or event_type == "response.completed"
self.terminal_status = self.terminal_status or event_type.removeprefix("response.")
return True
# Exact-type handlers first, then substring-matched ones in priority order. ``error`` frames
# carry the provider's real failure reason; raise so the credential pool + classifier see the body.
_EXACT_HANDLERS = {
"error": lambda self, event, event_type: _raise_stream_error(event),
"response.output_item.added": _on_item_added, "response.output_item.done": _on_item_done,
"response.completed": _on_terminal, "response.incomplete": _on_terminal, "response.failed": _on_terminal,
}
_FUZZY_HANDLERS = (
(lambda t: "output_text.delta" in t, _on_text_delta), (lambda t: "function_call" in t, _on_function_call),
(lambda t: "reasoning" in t and "delta" in t, _on_reasoning_delta),
)
def feed(self, event: Any) -> bool:
"""Process one event; True when the stream hit a terminal frame."""
event_type = _event_field(event, "type", "")
event_type = event_type if isinstance(event_type, str) else ""
handler = self._EXACT_HANDLERS.get(event_type) or next((h for m, h in self._FUZZY_HANDLERS if m(event_type)), None)
return bool(handler(self, event, event_type)) if handler is not None else False
def _settled_output(self) -> List[Any]:
"""Merge .done items with settled pending calls, keeping stream order."""
indexed = list(zip(self.output_indexes, self.output_sequences, self.output_items))
for pending in self.pending_function_calls.values():
item = pending["item"]
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.
arguments=(pending["arguments"] or "").strip() or "{}",
)))
# output_index is optional: protocol order only when every entry has one, else wire order.
if all(entry[0] is not None for entry in indexed):
with suppress(TypeError): # non-comparable index values: keep wire order
indexed.sort(key=lambda entry: entry[0])
else:
indexed.sort(key=lambda entry: entry[1])
return [entry[2] for entry in indexed]
def result(self) -> SimpleNamespace:
# With only plain text deltas (no tool calls), synthesize one message item.
output: List[Any] = list(self.output_items)
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.
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,
error=self.terminal_error)
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).
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)
for event in event_iter:
if on_event is not None:
try:
on_event(event)
except (TimeoutError, InterruptedError):
raise # watchdog / cancellation control flow must propagate
except Exception:
logger.debug("Codex stream on_event hook raised", exc_info=True)
if (interrupt_check is not None and interrupt_check()) or assembler.feed(event):
break
return assembler.result()
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."""
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)
if not (callable(backend_predicate) and bool(backend_predicate())):
return sanitized
dropped_from = ["top-level"] if "prompt_cache_retention" in sanitized else []
sanitized.pop("prompt_cache_retention", None)
# Copy before editing (caller's mapping must not mutate); drop when emptied.
extra_body = sanitized.get("extra_body")
if isinstance(extra_body, dict) and "prompt_cache_retention" in extra_body:
sanitized["extra_body"] = {k: v for k, v in extra_body.items() if k != "prompt_cache_retention"}
if not sanitized["extra_body"]:
sanitized.pop("extra_body")
dropped_from.append("extra_body")
if dropped_from:
logger.warning("Dropped unsupported prompt_cache_retention at consumer Codex wire boundary (model=%s, via %s).",
sanitized.get("model", getattr(agent, "model", "unknown")), ", ".join(dropped_from))
return sanitized
# 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; 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):
return all(isinstance(key, str) and _is_plain_json_data(item) for key, item in value.items())
if isinstance(value, list):
return all(_is_plain_json_data(item) for item in value)
return False
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."""
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
if isinstance(stream_kwargs.get(f), (dict, list)) and _is_plain_json_data(stream_kwargs[f])}
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 {}
# An explicit caller-provided extra_body entry keeps precedence (SDK post-transform merge).
bypassed["extra_body"] = {**merged, **{f: v for f, v in moved.items() if f not in merged}}
return bypassed
def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta=None):
"""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
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, 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.
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}
def _request_is_current() -> bool:
return request_token is None or getattr(agent, "_active_codex_stream_request_token", None) is request_token
def _fenced(fn: Callable[[Any], None]) -> Callable[[Any], None]:
"""Wrap a callback so a retired request's late frames never reach the agent."""
return lambda value: fn(value) if _request_is_current() else None
def _on_text_delta(text: str) -> None:
agent._codex_streamed_text_parts.append(text)
agent._fire_stream_delta(text)
def _on_event(event: Any) -> None: # TTFB watchdog and activity touch — once per SSE event.
agent._codex_stream_last_event_ts = time.time()
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.
if not _request_is_current():
raise TimeoutError("Codex Responses stream request retired before terminal response")
return bool(agent._interrupt_requested)
def _open_codex_stream(next_api_kwargs: dict[str, Any]):
stream_kwargs = _sanitize_consumer_codex_request(agent, next_api_kwargs)
stream_kwargs["stream"] = True
return active_client.responses.create(**_bypass_sdk_request_transform(stream_kwargs))
def _log_failure(exc: BaseException) -> None:
request_body_bytes, exception_chain = _codex_request_failure_details(exc)
logger.warning("Codex Responses request failed: serialized_request_body_bytes=%s stream_opened=%s "
"exception_chain=%s model=%s", "unknown" if request_body_bytes is None else request_body_bytes,
str(writer_token["value"] is not None).lower(), exception_chain, getattr(agent, "model", "unknown"))
def _codex_stream_created(_raw_stream: Any) -> None:
# Claim the delta sink for THIS attempt; a newer attempt supersedes this token.
writer_token["value"] = claim_stream_writer(agent)
def _accept_codex_chunk(_chunk: Any) -> bool:
token = writer_token["value"]
if token is None or stream_writer_is_current(agent, token):
return True
logger.warning("Codex streaming attempt superseded by a newer stream; stopping consumption to preserve "
"the single-writer invariant (model=%s).", api_kwargs.get("model", "unknown"))
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.
try:
for _ignored in event_stream:
pass
except (*transport_errors, _APIConnectionError) as exc:
if not isinstance(exc, transport_errors):
_log_failure(exc)
logger.warning("Codex Responses stream transport finalization failed after a terminal response was already "
"received; returning the completed response instead of retrying. %s error=%s",
agent._client_log_context(), exc)
def _close_event_stream(event_stream: Any) -> None:
close_fn = getattr(event_stream, "close", None) # None while connect never succeeded
try:
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.
if client is not None:
agent._abort_request_openai_client(active_client, reason="codex_stream_close_failed")
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")
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"] = event_stream = None
try:
try:
event_stream = relay_llm.stream(
dict(api_kwargs), _open_codex_stream,
session_id=str(getattr(agent, "session_id", "") or ""),
name=str(getattr(agent, "provider", "") or "codex"), model_name=str(model or ""),
finalizer=lambda: _consume_codex_event_stream(list(intercepted_events), model=model),
on_stream_created=_codex_stream_created, on_chunk=intercepted_events.append,
chunk_adapter=lambda chunk: chunk, accept_chunk=_accept_codex_chunk,
completed_response_predicate=lambda r: bool(hasattr(r, "output") and not hasattr(r, "__iter__")),
metadata={"api_mode": "codex_responses", "call_role": call_role, "retry_count": attempt,
"api_request_id": getattr(agent, "_current_api_request_id", None)},
defer_logical_completion=True,
)
final = _consume_codex_event_stream(
event_stream, model=model, on_text_delta=_fenced(_on_text_delta),
on_reasoning_delta=_fenced(lambda text: agent._fire_reasoning_delta(text)),
on_commentary_message=on_commentary_message, on_first_delta=on_first_delta,
on_event=_fenced(_on_event), interrupt_check=_interrupt_or_superseded,
)
except transport_errors as exc:
if attempt >= max_stream_retries:
_log_failure(exc)
raise
logger.debug(
"Codex Responses stream connect failed (attempt %s/%s); retrying. %s error=%s" if event_stream is None
else "Codex Responses stream transport failed mid-iteration (attempt %s/%s); retrying. %s error=%s",
attempt + 1, max_stream_retries + 1, agent._client_log_context(), exc,
)
continue
except RuntimeError:
# "No terminal response"; Relay may still hold a finalizer-assembled response.
if event_stream is not None and event_stream.final_response is not None:
return event_stream.final_response
raise
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 "
"(incomplete_details=%s, error=%s, streamed_chars=%d). %s",
final.status, final.incomplete_details, final.error,
sum(len(p) for p in agent._codex_streamed_text_parts), agent._client_log_context())
return final
finally:
_close_event_stream(event_stream)
__all__ = [
"run_codex_app_server_turn", "run_codex_stream",
"_consume_codex_event_stream", "make_codex_app_server_event_bridge",
]
# ---- BEGIN PLUGIN-COMPAT (revert-scheduled; see COMPAT_MANIFEST.md) ----
# Names external plugins imported from this module before the Sep 2026 decomposition.
# Internal code MUST NOT use these (scripts/check_compat_pointers.py fails CI if it does).
# The whole block is removed by reverting the commit that added it.
def run_codex_create_stream_fallback(agent, api_kwargs: dict, client: Any = None):
"""Backward-compatible alias for the unified event-driven path.
Historically this was the fallback when the SDK's high-level
``responses.stream(...)`` helper raised on shape drift. The primary
path now does exactly what the fallback did, so this just forwards.
Kept as a public symbol because tests and a small number of call sites
still reference it by name.
"""
return run_codex_stream(agent, api_kwargs, client=client)
# ---- END PLUGIN-COMPAT ----