From 1f615bbb17680dbfd0ec991fc5202eb57fc74628 Mon Sep 17 00:00:00 2001 From: "Abiba (pi)" Date: Sat, 27 Jun 2026 19:12:47 +0000 Subject: [PATCH] fix(agent-zero): queue expiry detection in adapter Zulip event queues expire after ~600s of inactivity. The adapter now detects 400/BAD_EVENT_QUEUE_ID responses and auto-reconnects instead of silently returning empty results. --- .../src/agent_zero_zulip/adapter.py | 32 ++++++++++++------- 1 file changed, 21 insertions(+), 11 deletions(-) diff --git a/agent-zero-zulip/src/agent_zero_zulip/adapter.py b/agent-zero-zulip/src/agent_zero_zulip/adapter.py index c1d637b..42ca483 100644 --- a/agent-zero-zulip/src/agent_zero_zulip/adapter.py +++ b/agent-zero-zulip/src/agent_zero_zulip/adapter.py @@ -189,13 +189,18 @@ class AgentZeroZulipAdapter: await asyncio.sleep(DEFAULT_POLL_INTERVAL) async def _fetch_events(self) -> list: - """Fetch events from Zulip queue.""" - resp = await self._api_call("GET", "/api/v1/events", params={ + """Fetch events from Zulip queue. Raises on BAD_EVENT_QUEUE_ID.""" + resp, status = await self._api_call("GET", "/api/v1/events", params={ "queue_id": self._queue_id, "last_event_id": str(self._last_event_id), "dont_block": "true", - }) - if resp is None: + }, return_status=True) + + # Detect queue expiry + if status == 400: + raise RuntimeError("BAD_EVENT_QUEUE_ID: queue expired") + + if not resp: return [] events = resp.get("events", []) @@ -365,10 +370,14 @@ class AgentZeroZulipAdapter: # ── Zulip API ── async def _api_call(self, method: str, path: str, - data: dict = None, params: dict = None) -> Optional[dict]: - """Make Zulip API call.""" + data: dict = None, params: dict = None, + return_status: bool = False) -> Any: + """Make Zulip API call. + + Returns: dict response body, or if return_status=True returns (dict, int) + """ if not self._http_client: - return None + return (None, 0) if return_status else None url = f"{self._site}{path}" headers = {"Authorization": self._auth_header} @@ -381,16 +390,17 @@ class AgentZeroZulipAdapter: elif method == "PATCH": resp = await self._http_client.patch(url, data=data, headers=headers) else: - return None + return (None, 0) if return_status else None if resp.status_code >= 400: logger.debug(f"API {method} {path}: {resp.status_code}") - return None + return (None, resp.status_code) if return_status else None - return resp.json() + body = resp.json() + return (body, resp.status_code) if return_status else body except Exception as e: logger.debug(f"API error {method} {path}: {e}") - return None + return (None, 0) if return_status else None async def _send_msg(self, msg_type: str, to: str, content: str) -> Optional[int]: """Send a message to Zulip."""