refactor(agent): trim the watchdog phase split to what the tests pin
- Explicit-vs-implicit idle timeout detection uses env_float with a sentinel instead of a hand-rolled parse (unset and unparseable both mean implicit). - _on_event writes the agent-level timestamps only when there is no request-local watchdog state; with state present they were dual bookkeeping nobody read (the snapshot prefers the state). - _abort_request keeps main's close-then-retire order: the reorder had no test teeth and is a separate concern from #90449.
This commit is contained in:
@@ -1088,13 +1088,10 @@ def _resolve_nonstream_watchdogs(agent, api_kwargs: dict) -> _NonStreamWatchdogs
|
||||
f"{est_tokens:,}")
|
||||
ttfb_timeout = ttfb_cap
|
||||
|
||||
idle_raw = os.getenv("HERMES_CODEX_EVENT_STALE_TIMEOUT_SECONDS", "").strip()
|
||||
try:
|
||||
float(idle_raw)
|
||||
except ValueError:
|
||||
idle_explicit = False
|
||||
else:
|
||||
idle_explicit = True
|
||||
# An operator-set idle timeout keeps first-event semantics; only the implicit
|
||||
# default defers arming until model progress. Sentinel: env_float returns the
|
||||
# default for unset AND unparseable values, so both count as implicit.
|
||||
idle_explicit = env_float("HERMES_CODEX_EVENT_STALE_TIMEOUT_SECONDS", -1.0) != -1.0
|
||||
idle_timeout = env_float("HERMES_CODEX_EVENT_STALE_TIMEOUT_SECONDS", idle_default)
|
||||
return _NonStreamWatchdogs(stale_timeout=stale_timeout, codex=codex, est_tokens=est_tokens,
|
||||
ttfb_enabled=ttfb_enabled, ttfb_timeout=ttfb_timeout, idle_enabled=codex and idle_timeout > 0,
|
||||
@@ -1211,11 +1208,9 @@ class _NonStreamRequest:
|
||||
"""Watchdog/interrupt kill: abort the request client (kind-aware, #67142)
|
||||
and retire the codex token; the worker sees its own forced close via
|
||||
the cancel flags."""
|
||||
# Retire before closing the socket: close can wake the worker immediately,
|
||||
# and a retired worker must not open a new physical retry.
|
||||
self._retire_codex_request_token()
|
||||
with contextlib.suppress(Exception):
|
||||
self.clients.close_once(reason)
|
||||
self._retire_codex_request_token()
|
||||
|
||||
def _await_worker_after_kill(self, timeout_message: str) -> None:
|
||||
# Wait briefly for the worker to notice the closed connection.
|
||||
|
||||
@@ -881,21 +881,18 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta
|
||||
def _on_event(event: Any) -> None: # TTFB/activity touch — once per SSE event.
|
||||
now = time.time()
|
||||
has_progress = _codex_event_has_content(event)
|
||||
reset_progress = False
|
||||
if watchdog_state is not None:
|
||||
with watchdog_state.lock:
|
||||
reset_progress = watchdog_state.retry_started_ts is not None
|
||||
if reset_progress:
|
||||
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._codex_stream_last_event_ts = now
|
||||
if reset_progress:
|
||||
agent._codex_stream_last_progress_ts = None
|
||||
if has_progress:
|
||||
agent._codex_stream_last_progress_ts = now
|
||||
else: # legacy callers without a request-local state (aux/streaming fallback)
|
||||
agent._codex_stream_last_event_ts = now
|
||||
if has_progress:
|
||||
agent._codex_stream_last_progress_ts = now
|
||||
agent._touch_activity("receiving stream response")
|
||||
|
||||
def _interrupt_or_superseded() -> bool:
|
||||
|
||||
Reference in New Issue
Block a user