fix: wait sixty seconds before provider silence notices
This commit is contained in:
@@ -138,10 +138,10 @@ class _NonStreamRequest:
|
||||
return
|
||||
silence = self.call_start + elapsed - (
|
||||
activity_ts if activity_ts is not None else self.call_start)
|
||||
if activity_ts is not None and silence < 30.0:
|
||||
if silence < 60.0:
|
||||
self.agent._touch_activity(
|
||||
"waiting for first stream event after reconnect"
|
||||
if retry_started_ts is not None else "receiving stream response")
|
||||
if retry_started_ts is not None else "waiting for provider response")
|
||||
return
|
||||
status = "no response yet"
|
||||
if retry_started_ts is not None:
|
||||
|
||||
@@ -9,7 +9,7 @@ from agent.model_metadata import is_local_endpoint
|
||||
class StreamingWaitMonitor:
|
||||
def _poll_local_load_notice(self, now: float) -> bool:
|
||||
"""Managed local server: surface a cold model's weight-load progress
|
||||
instead of the 30s "provider may be slow" copy. Polled ~1s only while no
|
||||
instead of the 60s "provider may be slow" copy. Polled ~1s only while no
|
||||
REAL chunk arrived for 2s+ (never during healthy token flow); in-memory,
|
||||
no network. True while loading = heartbeat liveness, skip the rest of
|
||||
this iteration (the stale detector's local floor dwarfs any load)."""
|
||||
@@ -34,11 +34,11 @@ class StreamingWaitMonitor:
|
||||
self.agent._emit_wait_notice("")
|
||||
return False
|
||||
|
||||
def _heartbeat(self, waiting_secs: int, interval: float) -> None:
|
||||
def _heartbeat(self, waiting_secs: int) -> None:
|
||||
"""Gateway inactivity heartbeat: the start-to-first-chunk gap (thinking,
|
||||
local prefill) can exceed the gateway timeout."""
|
||||
if waiting_secs >= interval:
|
||||
# No chunks for 30s+: say WHAT the wait is and WHEN recovery kicks in.
|
||||
if waiting_secs >= 60.0:
|
||||
# No chunks for 60s+: say WHAT the wait is and WHEN recovery kicks in.
|
||||
stale = self._stream_stale_timeout
|
||||
_recovery = f"; auto-reconnect at {int(stale)}s" if stale is not None and stale != float("inf") else ""
|
||||
self._mon.wait_notice_started_ts = self._mon.last_heartbeat
|
||||
@@ -69,7 +69,7 @@ class StreamingWaitMonitor:
|
||||
self._mon.wait_notice_started_ts = None
|
||||
if _hb_now - self._mon.last_heartbeat >= _HEARTBEAT_INTERVAL:
|
||||
self._mon.last_heartbeat = _hb_now
|
||||
self._heartbeat(int(_hb_now - self.last_chunk_time["t"]), _HEARTBEAT_INTERVAL)
|
||||
self._heartbeat(int(_hb_now - self.last_chunk_time["t"]))
|
||||
_stale_elapsed = time.time() - self.last_chunk_time["t"]
|
||||
if _stale_elapsed > self._stream_stale_timeout:
|
||||
self._mon.wait_notice_started_ts = None # Reconnect status has its own owner.
|
||||
|
||||
@@ -425,7 +425,7 @@ def test_wait_notice_omits_reconnect_when_all_deadlines_are_non_finite(
|
||||
|
||||
|
||||
def test_moa_heartbeat_survives_infinite_stale_timeout(monkeypatch):
|
||||
"""The full 100-poll MoA heartbeat must leave a healthy call running."""
|
||||
"""A MoA silence notice must leave an unbounded healthy call running."""
|
||||
from agent import chat_completion_helpers as h
|
||||
|
||||
notices: list[str] = []
|
||||
@@ -441,8 +441,11 @@ def test_moa_heartbeat_survives_infinite_stale_timeout(monkeypatch):
|
||||
_emit_wait_notice=notices.append,
|
||||
)
|
||||
|
||||
now = [1000.0]
|
||||
monkeypatch.setattr(h.time, "time", lambda: now[0])
|
||||
|
||||
class HeartbeatThread:
|
||||
"""Keep the synthetic worker alive through one heartbeat."""
|
||||
"""Keep the synthetic worker alive through the first silence notice."""
|
||||
|
||||
def __init__(self, *, target, daemon):
|
||||
self._polls = 0
|
||||
@@ -452,11 +455,11 @@ def test_moa_heartbeat_survives_infinite_stale_timeout(monkeypatch):
|
||||
pass
|
||||
|
||||
def join(self, timeout=None):
|
||||
pass
|
||||
now[0] = round(now[0] + timeout, 1)
|
||||
|
||||
def is_alive(self):
|
||||
self._polls += 1
|
||||
if self._polls == 101:
|
||||
if self._polls == 201:
|
||||
self._target()
|
||||
return False
|
||||
return True
|
||||
@@ -492,6 +495,9 @@ def test_wait_notice_formatting_error_does_not_abort_request(monkeypatch):
|
||||
_emit_wait_notice=lambda _message: None,
|
||||
)
|
||||
|
||||
now = [1000.0]
|
||||
monkeypatch.setattr(h.time, "time", lambda: now[0])
|
||||
|
||||
class HeartbeatThread:
|
||||
def __init__(self, *, target, daemon):
|
||||
self._polls = 0
|
||||
@@ -501,11 +507,11 @@ def test_wait_notice_formatting_error_does_not_abort_request(monkeypatch):
|
||||
pass
|
||||
|
||||
def join(self, timeout=None):
|
||||
pass
|
||||
now[0] = round(now[0] + timeout, 1)
|
||||
|
||||
def is_alive(self):
|
||||
self._polls += 1
|
||||
if self._polls == 101:
|
||||
if self._polls == 201:
|
||||
self._target()
|
||||
return False
|
||||
return True
|
||||
|
||||
@@ -38,11 +38,13 @@ def _request():
|
||||
[
|
||||
(59.0, 59.0, None, None), # Reasoning/text/tool arguments still arriving.
|
||||
(59.0, None, None, None), # Lifecycle traffic is not transport silence either.
|
||||
(10.0, 10.0, None, "50s with no stream events"),
|
||||
(10.0, None, None, "50s with no stream events"),
|
||||
(0.0, 0.0, None, "60s with no stream events"),
|
||||
(0.0, None, None, "60s with no stream events"),
|
||||
(1.0, None, None, None), # 59 seconds of silence is still quiet.
|
||||
(None, None, None, "60s with no response yet"),
|
||||
(10.0, 10.0, 59.0, None), # Internal reconnect gets a fresh first-event wait.
|
||||
(10.0, 10.0, 20.0, "40s with no response after reconnect"),
|
||||
(0.0, 0.0, 0.0, "60s with no response after reconnect"),
|
||||
(0.0, 0.0, 1.0, None),
|
||||
],
|
||||
)
|
||||
def test_wait_notice_tracks_current_attempt_silence(event, progress, retry, expected):
|
||||
@@ -51,6 +53,12 @@ def test_wait_notice_tracks_current_attempt_silence(event, progress, retry, expe
|
||||
state.last_event_ts = None if event is None else request.call_start + event
|
||||
state.last_progress_ts = None if progress is None else request.call_start + progress
|
||||
state.retry_started_ts = None if retry is None else request.call_start + retry
|
||||
request._emit_wait_notice(30.0)
|
||||
request._emit_wait_notice(59.0)
|
||||
assert notices == []
|
||||
assert touches
|
||||
if event is None and retry is None:
|
||||
assert "receiving" not in touches[-1]
|
||||
request._emit_wait_notice(60.0)
|
||||
if expected is None:
|
||||
assert notices == []
|
||||
@@ -74,13 +82,13 @@ def test_resumed_events_clear_only_this_requests_wait_notice(monkeypatch):
|
||||
pass
|
||||
|
||||
def is_alive(self):
|
||||
return ticks[0] < 104
|
||||
return ticks[0] < 204
|
||||
|
||||
def join(self, timeout):
|
||||
ticks[0] += 1
|
||||
if ticks[0] >= 101:
|
||||
if ticks[0] >= 201:
|
||||
request.codex_watchdog_state.last_event_ts = 1000.0 + ticks[0] * 0.3
|
||||
if ticks[0] == 104:
|
||||
if ticks[0] == 204:
|
||||
request.result["response"] = sentinel
|
||||
|
||||
monkeypatch.setattr(h.threading, "Thread", Worker)
|
||||
|
||||
@@ -9,46 +9,49 @@ from agent import chat_completion_helpers as h
|
||||
@pytest.mark.parametrize("local_loading,heartbeat_race", [(False, False), (True, False), (False, True)])
|
||||
def test_resumed_chunks_clear_wait_without_erasing_local_load(monkeypatch, local_loading, heartbeat_race):
|
||||
call = h._StreamingCall.__new__(h._StreamingCall)
|
||||
notices = []
|
||||
notices, touches = [], []
|
||||
now = [1000.0]
|
||||
call.agent = SimpleNamespace(
|
||||
base_url="http://localhost:1234" if local_loading else "https://example.com",
|
||||
_interrupt_requested=False,
|
||||
_emit_wait_notice=lambda text: notices.append((now[0], text)),
|
||||
_touch_activity=lambda text: None,
|
||||
_touch_activity=lambda text: touches.append((now[0], text)),
|
||||
)
|
||||
call.api_kwargs = {"model": "test-model"}
|
||||
call.last_chunk_time = {"t": now[0]}
|
||||
call._stream_stale_timeout = 180.0
|
||||
loading = "Loading local model weights"
|
||||
monkeypatch.setattr(h, "_managed_local_load_notice",
|
||||
lambda *args: loading if now[0] >= 1030.6 else None)
|
||||
lambda *args: loading if now[0] >= 1060.6 else None)
|
||||
monkeypatch.setattr(h.time, "time", lambda: now[0])
|
||||
|
||||
if heartbeat_race:
|
||||
heartbeat = call._heartbeat
|
||||
|
||||
def resume_before_notice(waiting_secs, interval):
|
||||
now[0] += 0.1
|
||||
call.last_chunk_time["t"] = now[0]
|
||||
heartbeat(waiting_secs, interval)
|
||||
def resume_before_notice(waiting_secs):
|
||||
if waiting_secs >= 60:
|
||||
now[0] += 0.1
|
||||
call.last_chunk_time["t"] = now[0]
|
||||
heartbeat(waiting_secs)
|
||||
|
||||
monkeypatch.setattr(call, "_heartbeat", resume_before_notice)
|
||||
|
||||
class Done:
|
||||
def is_set(self):
|
||||
return now[0] >= 1032.0
|
||||
return now[0] >= 1062.0
|
||||
|
||||
def wait(self, timeout):
|
||||
now[0] = round(now[0] + timeout, 1)
|
||||
if not heartbeat_race and now[0] >= 1031.5:
|
||||
if not heartbeat_race and now[0] >= 1061.5:
|
||||
call.last_chunk_time["t"] = now[0]
|
||||
|
||||
call._call_done = Done()
|
||||
call._monitor_loop()
|
||||
assert all(t >= 1060.0 for t, _ in notices)
|
||||
assert touches[0][0] == 1030.0 # Quiet gateway heartbeat is still 30s.
|
||||
assert "waiting on" in notices[0][1]
|
||||
if local_loading:
|
||||
assert notices[-1][1] == loading
|
||||
assert not any(text == "" for _, text in notices)
|
||||
else:
|
||||
assert notices[1:] == [(1030.4 if heartbeat_race else 1031.5, "")]
|
||||
assert notices[1:] == [(1060.4 if heartbeat_race else 1061.5, "")]
|
||||
|
||||
@@ -1190,7 +1190,7 @@ The **stale non-stream detection** kills non-streaming calls that produce no res
|
||||
|
||||
This budget bounds every non-streaming call. A provider that accepts a request and then goes silent — connection held open, no bytes, no error — is aborted at the stale timeout and retried, rather than hanging until the much longer socket read timeout (or, for an unattended cron run, until something external kills the process).
|
||||
|
||||
The Codex Responses **waiting status** describes silence, not total generation time: active stream events (including reasoning) keep it quiet. If events stop, it reports time without stream events instead of claiming no response has arrived; the notice clears when events resume. When a reconnect starts a fresh first-event watchdog phase, the waiting status follows that phase. This display behavior does not extend the separate wall-clock stale-call budget or change watchdog timeouts. Chat-completion streams likewise clear their silence warning promptly when chunks resume, without replacing a local model-loading status.
|
||||
The periodic provider-wait notice appears only after at least **60 seconds of silence**. The Codex Responses **waiting status** describes silence, not total generation time: active stream events (including reasoning) keep it quiet. If events stop, it reports time without stream events instead of claiming no response has arrived; the notice clears when events resume. When a reconnect starts a fresh first-event watchdog phase, the waiting status follows that phase. This display behavior does not extend the separate wall-clock stale-call budget or change watchdog timeouts. Chat-completion streams likewise clear their silence warning promptly when chunks resume, without replacing a local model-loading status.
|
||||
|
||||
Cron jobs and delegated subagents stream too. They run the request inline on their own thread (the interrupt worker other sessions use wedges inside the gateway's nested thread pools), but the wire request is still `stream: true`, so the **stale stream detection** budget above governs them — every token counts as liveness, so a reasoning model that thinks for minutes is not mistaken for a hung provider, and edge proxies that kill silent connections keep seeing bytes.
|
||||
|
||||
|
||||
+1
-1
@@ -738,7 +738,7 @@ Hermes 对流式传输有单独的超时层,以及用于非流式调用的陈
|
||||
|
||||
**陈旧非流检测**终止长时间没有响应的非流式调用。默认情况下,Hermes 在本地端点上禁用此功能,以避免长时间预填充期间的误报。如果您显式设置 `providers.<id>.stale_timeout_seconds`、`providers.<id>.models.<model>.stale_timeout_seconds` 或 `HERMES_API_CALL_STALE_TIMEOUT`,即使在本地端点上也会遵守该显式值。
|
||||
|
||||
Codex Responses 的**等待状态**描述的是没有流事件的时间,而不是生成的总时长:持续收到流事件(包括推理内容)时不会显示等待警告。如果事件停止,提示会显示多久没有收到流事件,而不会声称尚未收到任何响应;事件恢复后提示会清除。当重连进入新的首事件看门狗阶段时,等待状态也会跟随该阶段。这只影响状态显示,不会延长独立的调用总时长限制,也不会更改看门狗超时。Chat Completions 流同样会在数据块恢复时及时清除无输出警告,但不会覆盖本地模型加载状态。
|
||||
周期性的提供商等待提示仅在至少 **60 秒没有响应活动**后显示。Codex Responses 的**等待状态**描述的是没有流事件的时间,而不是生成的总时长:持续收到流事件(包括推理内容)时不会显示等待警告。如果事件停止,提示会显示多久没有收到流事件,而不会声称尚未收到任何响应;事件恢复后提示会清除。当重连进入新的首事件看门狗阶段时,等待状态也会跟随该阶段。这只影响状态显示,不会延长独立的调用总时长限制,也不会更改看门狗超时。Chat Completions 流同样会在数据块恢复时及时清除无输出警告,但不会覆盖本地模型加载状态。
|
||||
|
||||
## 上下文压力警告
|
||||
|
||||
|
||||
Reference in New Issue
Block a user