From a3bbb80d451d6164788fafa3cbc605bb9ec5f91b Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 20:46:03 -0700 Subject: [PATCH] =?UTF-8?q?refactor(gateway):=20watchers=20=E2=80=94=20fla?= =?UTF-8?q?tten=20stall-notice=20send;=20boundary=20fallback=20early=20ret?= =?UTF-8?q?urn?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- gateway/run_watchers.py | 39 ++++++++++++++++++-------------------- gateway/stream_consumer.py | 17 ++++++++--------- 2 files changed, 26 insertions(+), 30 deletions(-) diff --git a/gateway/run_watchers.py b/gateway/run_watchers.py index 79ba0d7533..e3492e6f5e 100644 --- a/gateway/run_watchers.py +++ b/gateway/run_watchers.py @@ -188,11 +188,9 @@ class GatewaySessionWatchersMixin: agent = (getattr(self, "_running_agents", None) or {}).get(session_key) if agent is None or agent is _AGENT_PENDING_SENTINEL: return None - if not hasattr(agent, "get_activity_summary"): - return None try: summary = agent.get_activity_summary() - except Exception: + except Exception: # incl. AttributeError: agent without an activity summary return None return summary if isinstance(summary, dict) else None @@ -269,27 +267,26 @@ class GatewaySessionWatchersMixin: notice = format_session_stall_notification(idle_seconds) # Bound the send: a wedged adapter transport (network hang, dead websocket) must not # block the watcher pass — siblings would go unevaluated and the watcher stop. - try: - result = await asyncio.wait_for( - adapter.send(str(source.chat_id), notice, metadata=metadata), - timeout=_STALL_NOTIFY_SEND_TIMEOUT_SECONDS, - ) - except asyncio.TimeoutError: - logger.warning( - "Session stall notify send timed out after %.0fs for %s; will retry next tick", - _STALL_NOTIFY_SEND_TIMEOUT_SECONDS, session_key, - ) - return False - # Adapters often return SendResult(success=False) instead of raising. - if result is not None and getattr(result, "success", True) is False: - logger.warning( - "Session stall notify failed for %s: %s", - session_key, getattr(result, "error", "send returned success=False"), - ) - return False + result = await asyncio.wait_for( + adapter.send(str(source.chat_id), notice, metadata=metadata), + timeout=_STALL_NOTIFY_SEND_TIMEOUT_SECONDS, + ) + except asyncio.TimeoutError: + logger.warning( + "Session stall notify send timed out after %.0fs for %s; will retry next tick", + _STALL_NOTIFY_SEND_TIMEOUT_SECONDS, session_key, + ) + return False except Exception as exc: logger.warning("Session stall notify failed for %s: %s", session_key, exc) return False + # Adapters often return SendResult(success=False) instead of raising. + if result is not None and getattr(result, "success", True) is False: + logger.warning( + "Session stall notify failed for %s: %s", + session_key, getattr(result, "error", "send returned success=False"), + ) + return False return True async def _notify_session_stall( diff --git a/gateway/stream_consumer.py b/gateway/stream_consumer.py index 3c84d9c3a6..8c33b67303 100644 --- a/gateway/stream_consumer.py +++ b/gateway/stream_consumer.py @@ -497,19 +497,18 @@ class GatewayStreamConsumer( "falling back to send() for pre-prompt text (chat=%s)", _reason, self.chat_id, ) - fallback_ok = False try: send_result = await self.adapter.send(self.chat_id, finalize_text) - fallback_ok = getattr(send_result, "success", False) + if getattr(send_result, "success", False): + return True except Exception as send_err: logger.warning("%s boundary: fallback send also failed: %s", _reason, send_err) - if not fallback_ok: - logger.error( - "%s boundary: both finalize and fallback send failed " - "(chat=%s) — pre-prompt text may not have been delivered", - _reason, self.chat_id, - ) - return fallback_ok + logger.error( + "%s boundary: both finalize and fallback send failed " + "(chat=%s) — pre-prompt text may not have been delivered", + _reason, self.chat_id, + ) + return False def on_delta(self, text: str) -> None: """Thread-safe callback from the agent's worker thread. ``None`` signals a tool