refactor(gateway): watchers — flatten stall-notice send; boundary fallback early return

This commit is contained in:
Teknium
2026-09-02 20:46:03 -07:00
parent e6bc37845f
commit a3bbb80d45
2 changed files with 26 additions and 30 deletions
+18 -21
View File
@@ -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(
+8 -9
View File
@@ -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