Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b1bef9024f | ||
|
|
dcf2de0052 | ||
|
|
715d54564f |
@@ -43,7 +43,7 @@ except ImportError:
|
|||||||
HTTPX_AVAILABLE = False
|
HTTPX_AVAILABLE = False
|
||||||
httpx = None
|
httpx = None
|
||||||
|
|
||||||
from gateway.config import PlatformConfig
|
from gateway.config import Platform, PlatformConfig
|
||||||
from gateway.platforms.base import (
|
from gateway.platforms.base import (
|
||||||
BasePlatformAdapter,
|
BasePlatformAdapter,
|
||||||
MessageEvent,
|
MessageEvent,
|
||||||
@@ -118,8 +118,8 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
"""
|
"""
|
||||||
|
|
||||||
def __init__(self, config: PlatformConfig):
|
def __init__(self, config: PlatformConfig):
|
||||||
platform_name = "zulip"
|
platform = Platform("zulip")
|
||||||
super().__init__(config=config, platform=platform_name)
|
super().__init__(config=config, platform=platform)
|
||||||
|
|
||||||
extra = config.extra or {}
|
extra = config.extra or {}
|
||||||
|
|
||||||
@@ -242,7 +242,7 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
self._health_stats["started_at"] = datetime.now(timezone.utc).isoformat()
|
self._health_stats["started_at"] = datetime.now(timezone.utc).isoformat()
|
||||||
|
|
||||||
# Register event queue
|
# Register event queue
|
||||||
queue_resp = await self._api_call(
|
queue_resp, _ = await self._api_call(
|
||||||
"POST", "/api/v1/register",
|
"POST", "/api/v1/register",
|
||||||
data={
|
data={
|
||||||
"event_types": '["message"]',
|
"event_types": '["message"]',
|
||||||
@@ -259,6 +259,11 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
self._last_event_id = data.get("last_event_id", -1)
|
self._last_event_id = data.get("last_event_id", -1)
|
||||||
self._bot_user_id = data.get("user_id")
|
self._bot_user_id = data.get("user_id")
|
||||||
|
|
||||||
|
# Resolve bot user_id if register endpoint didn't provide it
|
||||||
|
# (common on some Zulip versions)
|
||||||
|
if self._bot_user_id is None:
|
||||||
|
await self._resolve_bot_user_id()
|
||||||
|
|
||||||
# Gen 2: Try to resolve @all-bots user ID dynamically
|
# Gen 2: Try to resolve @all-bots user ID dynamically
|
||||||
await self._resolve_all_bots_user_id()
|
await self._resolve_all_bots_user_id()
|
||||||
|
|
||||||
@@ -381,11 +386,15 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
await asyncio.sleep(self._poll_interval)
|
await asyncio.sleep(self._poll_interval)
|
||||||
|
|
||||||
async def _fetch_events(self) -> List[Dict[str, Any]]:
|
async def _fetch_events(self) -> List[Dict[str, Any]]:
|
||||||
"""Fetch events from the Zulip event queue."""
|
"""Fetch events from the Zulip event queue.
|
||||||
|
|
||||||
|
Raises RuntimeError with BAD_EVENT_QUEUE_ID when the queue
|
||||||
|
expires, so _poll_forever can reconnect.
|
||||||
|
"""
|
||||||
if not self._queue_id:
|
if not self._queue_id:
|
||||||
return []
|
return []
|
||||||
|
|
||||||
resp = await self._api_call(
|
resp, raw_status = await self._api_call(
|
||||||
"GET", "/api/v1/events",
|
"GET", "/api/v1/events",
|
||||||
params={
|
params={
|
||||||
"queue_id": self._queue_id,
|
"queue_id": self._queue_id,
|
||||||
@@ -393,6 +402,9 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
"dont_block": "true",
|
"dont_block": "true",
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
|
# Detect queue expiry from HTTP response
|
||||||
|
if raw_status == 400:
|
||||||
|
raise RuntimeError("BAD_EVENT_QUEUE_ID: queue expired")
|
||||||
if not resp:
|
if not resp:
|
||||||
return []
|
return []
|
||||||
|
|
||||||
@@ -409,7 +421,7 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
"""Re-register the event queue."""
|
"""Re-register the event queue."""
|
||||||
self._queue_id = None
|
self._queue_id = None
|
||||||
try:
|
try:
|
||||||
resp = await self._api_call(
|
resp, _ = await self._api_call(
|
||||||
"POST", "/api/v1/register",
|
"POST", "/api/v1/register",
|
||||||
data={
|
data={
|
||||||
"event_types": '["message"]',
|
"event_types": '["message"]',
|
||||||
@@ -496,6 +508,32 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
"[%s] @all-bots refresh error (non-fatal)", self.name
|
"[%s] @all-bots refresh error (non-fatal)", self.name
|
||||||
)
|
)
|
||||||
|
|
||||||
|
# ------------------------------------------------------------------
|
||||||
|
# Gen 5: Bot user ID resolution
|
||||||
|
# ------------------------------------------------------------------
|
||||||
|
|
||||||
|
async def _resolve_bot_user_id(self) -> None:
|
||||||
|
"""Fetch the bot's own user ID from /api/v1/users/me.
|
||||||
|
|
||||||
|
The /register endpoint does not always return user_id on
|
||||||
|
all Zulip server versions. This ensures @mention detection
|
||||||
|
works by resolving it from the identity endpoint.
|
||||||
|
"""
|
||||||
|
try:
|
||||||
|
resp, _ = await self._api_call("GET", "/api/v1/users/me")
|
||||||
|
if resp:
|
||||||
|
uid = resp.get("user_id")
|
||||||
|
if uid is not None:
|
||||||
|
self._bot_user_id = int(uid)
|
||||||
|
logger.info(
|
||||||
|
"[%s] Resolved bot user_id=%s from /users/me",
|
||||||
|
self.name, uid,
|
||||||
|
)
|
||||||
|
except Exception as e:
|
||||||
|
logger.debug(
|
||||||
|
"[%s] Could not resolve bot user_id: %s", self.name, e
|
||||||
|
)
|
||||||
|
|
||||||
# ------------------------------------------------------------------
|
# ------------------------------------------------------------------
|
||||||
# Gen 2: Dynamic @all-bots resolution
|
# Gen 2: Dynamic @all-bots resolution
|
||||||
# ------------------------------------------------------------------
|
# ------------------------------------------------------------------
|
||||||
@@ -506,7 +544,7 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
Falls back to the configured or default value if resolution fails.
|
Falls back to the configured or default value if resolution fails.
|
||||||
"""
|
"""
|
||||||
try:
|
try:
|
||||||
resp = await self._api_call("GET", "/api/v1/users")
|
resp, _ = await self._api_call("GET", "/api/v1/users")
|
||||||
if not resp:
|
if not resp:
|
||||||
return
|
return
|
||||||
members = resp.get("members", [])
|
members = resp.get("members", [])
|
||||||
@@ -766,7 +804,7 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
"content": truncated,
|
"content": truncated,
|
||||||
}
|
}
|
||||||
try:
|
try:
|
||||||
resp = await self._api_call("PATCH", "/api/v1/messages", data=payload)
|
resp, _ = await self._api_call("PATCH", "/api/v1/messages", data=payload)
|
||||||
return resp is not None
|
return resp is not None
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.warning("[%s] Edit message %s error: %s", self.name, message_id, e)
|
logger.warning("[%s] Edit message %s error: %s", self.name, message_id, e)
|
||||||
@@ -782,6 +820,7 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
"POST", "/api/v1/typing",
|
"POST", "/api/v1/typing",
|
||||||
data={"to": user_ids, "op": "start"},
|
data={"to": user_ids, "op": "start"},
|
||||||
)
|
)
|
||||||
|
# typing indicator failure is non-critical
|
||||||
except Exception:
|
except Exception:
|
||||||
pass # Non-critical
|
pass # Non-critical
|
||||||
|
|
||||||
@@ -795,6 +834,7 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
"POST", "/api/v1/typing",
|
"POST", "/api/v1/typing",
|
||||||
data={"to": user_ids, "op": "stop"},
|
data={"to": user_ids, "op": "stop"},
|
||||||
)
|
)
|
||||||
|
# typing indicator failure is non-critical
|
||||||
except Exception:
|
except Exception:
|
||||||
pass # Non-critical
|
pass # Non-critical
|
||||||
|
|
||||||
@@ -1009,10 +1049,15 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
path: str,
|
path: str,
|
||||||
data: Optional[Dict] = None,
|
data: Optional[Dict] = None,
|
||||||
params: Optional[Dict] = None,
|
params: Optional[Dict] = None,
|
||||||
) -> Optional[Dict]:
|
) -> Tuple[Optional[Dict], int]:
|
||||||
"""Make an API call to Zulip."""
|
"""Make an API call to Zulip.
|
||||||
|
|
||||||
|
Returns (response_json, status_code). status_code is useful for
|
||||||
|
callers that need to detect specific HTTP errors like 400
|
||||||
|
(BAD_EVENT_QUEUE_ID). Returns (None, 0) on transport errors.
|
||||||
|
"""
|
||||||
if not self._http_client:
|
if not self._http_client:
|
||||||
return None
|
return None, 0
|
||||||
|
|
||||||
url = f"{self._site}{path}"
|
url = f"{self._site}{path}"
|
||||||
headers = {
|
headers = {
|
||||||
@@ -1033,7 +1078,7 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
url, data=data, headers=headers,
|
url, data=data, headers=headers,
|
||||||
)
|
)
|
||||||
else:
|
else:
|
||||||
return None
|
return None, 0
|
||||||
|
|
||||||
if response.status_code >= 400:
|
if response.status_code >= 400:
|
||||||
logger.warning(
|
logger.warning(
|
||||||
@@ -1041,20 +1086,21 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
self.name, method, path,
|
self.name, method, path,
|
||||||
response.status_code, response.text[:200],
|
response.status_code, response.text[:200],
|
||||||
)
|
)
|
||||||
return None
|
return None, response.status_code
|
||||||
|
|
||||||
return response.json()
|
return response.json(), response.status_code
|
||||||
|
|
||||||
except httpx.TimeoutException:
|
except httpx.TimeoutException:
|
||||||
logger.debug("[%s] Timeout on %s %s", self.name, method, path)
|
logger.debug("[%s] Timeout on %s %s", self.name, method, path)
|
||||||
return None
|
return None, 0
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.warning("[%s] API error %s %s: %s", self.name, method, path, e)
|
logger.warning("[%s] API error %s %s: %s", self.name, method, path, e)
|
||||||
return None
|
return None, 0
|
||||||
|
|
||||||
async def _send_api_call(self, payload: Dict) -> Optional[Dict]:
|
async def _send_api_call(self, payload: Dict) -> Optional[Dict]:
|
||||||
"""Send a message to Zulip."""
|
"""Send a message to Zulip."""
|
||||||
return await self._api_call("POST", "/api/v1/messages", data=payload)
|
result, _ = await self._api_call("POST", "/api/v1/messages", data=payload)
|
||||||
|
return result
|
||||||
|
|
||||||
# ------------------------------------------------------------------
|
# ------------------------------------------------------------------
|
||||||
# Utilities
|
# Utilities
|
||||||
|
|||||||
Reference in New Issue
Block a user