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