File size: 61,515 Bytes
c61f7ed | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522 523 524 525 526 527 528 529 530 531 532 533 534 535 536 537 538 539 540 541 542 543 544 545 546 547 548 549 550 551 552 553 554 555 556 557 558 559 560 561 562 563 564 565 566 567 568 569 570 571 572 573 574 575 576 577 578 579 580 581 582 583 584 585 586 587 588 589 590 591 592 593 594 595 596 597 598 599 600 601 602 603 604 605 606 607 608 609 610 611 612 613 614 615 616 617 618 619 620 621 622 623 624 625 626 627 628 629 630 631 632 633 634 635 636 637 638 639 640 641 642 643 644 645 646 647 648 649 650 651 652 653 654 655 656 657 658 659 660 661 662 663 664 665 666 667 668 669 670 671 672 673 674 675 676 677 678 679 680 681 682 683 684 685 686 687 688 689 690 691 692 693 694 695 696 697 698 699 700 701 702 703 704 705 706 707 708 709 710 711 712 713 714 715 716 717 718 719 720 721 722 723 724 725 726 727 728 729 730 731 732 733 734 735 736 737 738 739 740 741 742 743 744 745 746 747 748 749 750 751 752 753 754 755 756 757 758 759 760 761 762 763 764 765 766 767 768 769 770 771 772 773 774 775 776 777 778 779 780 781 782 783 784 785 786 787 788 789 790 791 792 793 794 795 796 797 798 799 800 801 802 803 804 805 806 807 808 809 810 811 812 813 814 815 816 817 818 819 820 821 822 823 824 825 826 827 828 829 830 831 832 833 834 835 836 837 838 839 840 841 842 843 844 845 846 847 848 849 850 851 852 853 854 855 856 857 858 859 860 861 862 863 864 865 866 867 868 869 870 871 872 873 874 875 876 877 878 879 880 881 882 883 884 885 886 887 888 889 890 891 892 893 894 895 896 897 898 899 900 901 902 903 904 905 906 907 908 909 910 911 912 913 914 915 916 917 918 919 920 921 922 923 924 925 926 927 928 929 930 931 932 933 934 935 936 937 938 939 940 941 942 943 944 945 946 947 948 949 950 951 952 953 954 955 956 957 958 959 960 961 962 963 964 965 966 967 968 969 970 971 972 973 974 975 976 977 978 979 980 981 982 983 984 985 986 987 988 989 990 991 992 993 994 995 996 997 998 999 1000 1001 1002 1003 1004 1005 1006 1007 1008 1009 1010 1011 1012 1013 1014 1015 1016 1017 1018 1019 1020 1021 1022 1023 1024 1025 1026 1027 1028 1029 1030 1031 1032 1033 1034 1035 1036 1037 1038 1039 1040 1041 1042 1043 1044 1045 1046 1047 1048 1049 1050 1051 1052 1053 1054 1055 1056 1057 1058 1059 1060 1061 1062 1063 1064 1065 | """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 contextvars
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
from agent.usage_anchor import set_usage_anchor
logger = logging.getLogger(__name__)
_codex_watchdog_state_var: contextvars.ContextVar[Any | None] = contextvars.ContextVar(
"codex_watchdog_state", default=None
)
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, messages=None) -> 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.
``messages`` (the transcript mirror) lets real usage anchor the next preflight: this runtime bypasses
the main loop's capture, and the mirror is never compacted natively, so without an anchor the rough
estimate grows monotonically and hermes-mode fires thread compaction on tiny threads (#100381)."""
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({})
if compressor is not None and callable(getattr(compressor, "note_usage_less_response", None)):
compressor.note_usage_less_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)
if isinstance(messages, list):
from agent.usage_anchor import capture_usage_anchor, set_usage_anchor
anchor = capture_usage_anchor(prompt_tokens, canonical_usage.output_tokens, messages)
if anchor is not None:
set_usage_anchor(agent, anchor)
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.
set_usage_anchor(agent, 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
projected_messages = turn.projected_messages
# Turn-start persistence owns the accepted input. Codex's leading user item
# only echoes its coerced wire text; later/nonmatching events remain real.
submitted_user_text = getattr(turn, "submitted_user_text", None)
first = projected_messages[0]
if (submitted_user_text is not None and first.get("role") == "user"
and first.get("content") == submitted_user_text):
projected_messages = projected_messages[1:]
for projected_message in 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, messages=messages)
# 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
_CODEX_PROGRESS_DELTA_TYPES = frozenset({
"response.output_text.delta", "response.reasoning_summary_text.delta", "response.text.delta",
"response.audio.delta", "response.function_call_arguments.delta", "response.reasoning_text.delta",
})
def _codex_event_has_content(event: Any) -> bool:
"""Whether a Codex Responses event carries substantive forward progress.
Lifecycle/keepalive frames and empty structural deltas prove transport
liveness, but do not mean the model has begun producing its response.
"""
event_type = _event_field(event, "type")
if event_type in _CODEX_PROGRESS_DELTA_TYPES:
return bool(_event_field(event, "delta"))
if event_type == "response.output_item.added":
item = _event_field(event, "item")
return "function_call" in str(_event_field(item, "type") or "") and any(
bool(_event_field(item, field)) for field in ("id", "call_id", "name", "arguments"))
return False
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.ReadError, _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.
watchdog_state = _codex_watchdog_state_var.get()
request_token = (
watchdog_state.token
if watchdog_state is not None
else 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/activity touch — once per SSE event.
now = time.time()
has_progress = _codex_event_has_content(event)
if watchdog_state is not None:
with watchdog_state.lock:
if watchdog_state.retry_started_ts is not None:
watchdog_state.retry_started_ts = None
watchdog_state.last_progress_ts = None
watchdog_state.last_event_ts = now
if has_progress:
watchdog_state.last_progress_ts = now
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]):
from hermes_cli.providers import is_actual_route
if is_actual_route(
getattr(agent, "provider", ""),
str(getattr(active_client, "base_url", "") or ""),
):
raise ValueError(
"Actual requests require Chat Completions; refusing to call /responses."
)
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 not _request_is_current():
raise TimeoutError("Codex Responses stream request retired before retry")
if agent._interrupt_requested:
raise InterruptedError("Agent interrupted before Codex stream retry")
if attempt > 0 and watchdog_state is not None and watchdog_state.phase_aware:
# A physical reconnect has its own no-event TTFB phase. Its first parsed
# event clears this marker and starts a fresh model-progress phase.
with watchdog_state.lock:
watchdog_state.retry_started_ts = time.time()
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 ----
|