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.
This commit is contained in:
@@ -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."""
|
||||
|
||||
Reference in New Issue
Block a user