Compare commits
8
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b6fa3da540 | ||
|
|
04f78f3f01 | ||
|
|
3fbd31fbff | ||
|
|
a5dc880af1 | ||
|
|
9b9845fd14 | ||
|
|
106048999f | ||
|
|
6afc46734a | ||
|
|
9ee1919985 |
@@ -17,7 +17,7 @@ jobs:
|
|||||||
- run: python3 --version
|
- run: python3 --version
|
||||||
- run: node --version
|
- run: node --version
|
||||||
- run: echo "Runner works!"
|
- run: echo "Runner works!"
|
||||||
- run: git clone --depth 1 http://abiba-bot:mipjoq-tybbox-2ciHru@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
- run: git clone --depth 1 http://abiba-bot:***REMOVED***@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
||||||
|
|
||||||
- name: Python syntax check
|
- name: Python syntax check
|
||||||
run: |
|
run: |
|
||||||
@@ -46,7 +46,7 @@ print('✅ config.yaml.example valid')
|
|||||||
needs: [validate]
|
needs: [validate]
|
||||||
steps:
|
steps:
|
||||||
- run: echo "🚀 Deploy tag $(echo $GITHUB_REF_NAME)"
|
- run: echo "🚀 Deploy tag $(echo $GITHUB_REF_NAME)"
|
||||||
- run: git clone --depth 1 http://abiba-bot:mipjoq-tybbox-2ciHru@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
- run: git clone --depth 1 http://abiba-bot:***REMOVED***@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
||||||
|
|
||||||
- name: Deploy to Tanko (canary)
|
- name: Deploy to Tanko (canary)
|
||||||
run: ssh -o StrictHostKeyChecking=no -o ConnectTimeout=10 jerome@192.168.68.122 "cd /root && bash -s" < scripts/deploy.sh --ct=tanko --mode=native "$GITHUB_REF_NAME" 2>&1 || echo "⚠️ Tanko deploy skipped"
|
run: ssh -o StrictHostKeyChecking=no -o ConnectTimeout=10 jerome@192.168.68.122 "cd /root && bash -s" < scripts/deploy.sh --ct=tanko --mode=native "$GITHUB_REF_NAME" 2>&1 || echo "⚠️ Tanko deploy skipped"
|
||||||
|
|||||||
@@ -21,7 +21,7 @@ jobs:
|
|||||||
runs-on: ubuntu-latest
|
runs-on: ubuntu-latest
|
||||||
steps:
|
steps:
|
||||||
- name: Checkout
|
- name: Checkout
|
||||||
run: git clone --depth 1 http://abiba-bot:mipjoq-tybbox-2ciHru@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
run: git clone --depth 1 http://abiba-bot:***REMOVED***@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
||||||
- name: Python syntax
|
- name: Python syntax
|
||||||
run: |
|
run: |
|
||||||
python3 -m py_compile plugins/platforms/zulip/adapter.py 2>/dev/null
|
python3 -m py_compile plugins/platforms/zulip/adapter.py 2>/dev/null
|
||||||
@@ -37,7 +37,7 @@ jobs:
|
|||||||
environment: canary
|
environment: canary
|
||||||
steps:
|
steps:
|
||||||
- name: Checkout
|
- name: Checkout
|
||||||
run: git clone --depth 1 http://abiba-bot:mipjoq-tybbox-2ciHru@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
run: git clone --depth 1 http://abiba-bot:***REMOVED***@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
||||||
- name: Deploy to Tanko
|
- name: Deploy to Tanko
|
||||||
env:
|
env:
|
||||||
TAG: ${{ github.ref_name }}
|
TAG: ${{ github.ref_name }}
|
||||||
@@ -66,7 +66,7 @@ jobs:
|
|||||||
environment: production
|
environment: production
|
||||||
steps:
|
steps:
|
||||||
- name: Checkout
|
- name: Checkout
|
||||||
run: git clone --depth 1 http://abiba-bot:mipjoq-tybbox-2ciHru@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
run: git clone --depth 1 http://abiba-bot:***REMOVED***@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
||||||
- name: Deploy to all Hermes agents
|
- name: Deploy to all Hermes agents
|
||||||
env:
|
env:
|
||||||
TAG: ${{ github.ref_name }}
|
TAG: ${{ github.ref_name }}
|
||||||
@@ -97,7 +97,7 @@ jobs:
|
|||||||
environment: agent-zero
|
environment: agent-zero
|
||||||
steps:
|
steps:
|
||||||
- name: Checkout
|
- name: Checkout
|
||||||
run: git clone --depth 1 http://abiba-bot:mipjoq-tybbox-2ciHru@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
run: git clone --depth 1 http://abiba-bot:***REMOVED***@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
||||||
- name: Deploy to kagentz
|
- name: Deploy to kagentz
|
||||||
env:
|
env:
|
||||||
TAG: ${{ github.ref_name }}
|
TAG: ${{ github.ref_name }}
|
||||||
|
|||||||
@@ -203,22 +203,6 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
logger.error("Zulip POST %s network error: %s", path, exc)
|
logger.error("Zulip POST %s network error: %s", path, exc)
|
||||||
return {"result": "error", "msg": str(exc)}
|
return {"result": "error", "msg": str(exc)}
|
||||||
|
|
||||||
async def _api_patch(self, path: str, payload: Dict[str, Any]) -> Dict[str, Any]:
|
|
||||||
import aiohttp
|
|
||||||
url = f"{self._site}/api/v1/{path.lstrip('/')}"
|
|
||||||
try:
|
|
||||||
async with self._session.patch(
|
|
||||||
url, data=payload, auth=self._auth(),
|
|
||||||
timeout=aiohttp.ClientTimeout(total=15),
|
|
||||||
) as resp:
|
|
||||||
data = await resp.json()
|
|
||||||
if resp.status >= 400:
|
|
||||||
logger.debug("Zulip PATCH %s -> %s: %s", path, resp.status, str(data)[:200])
|
|
||||||
return data
|
|
||||||
except (aiohttp.ClientError, asyncio.TimeoutError) as exc:
|
|
||||||
logger.debug("Zulip PATCH %s network error: %s", path, exc)
|
|
||||||
return {"result": "error", "msg": str(exc)}
|
|
||||||
|
|
||||||
async def _api_delete(self, path: str, params: Optional[Dict[str, Any]] = None) -> Dict[str, Any]:
|
async def _api_delete(self, path: str, params: Optional[Dict[str, Any]] = None) -> Dict[str, Any]:
|
||||||
import aiohttp
|
import aiohttp
|
||||||
url = f"{self._site}/api/v1/{path.lstrip('/')}"
|
url = f"{self._site}/api/v1/{path.lstrip('/')}"
|
||||||
@@ -342,23 +326,6 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
|
|
||||||
return SendResult(success=True, message_id=str(last_id) if last_id else None)
|
return SendResult(success=True, message_id=str(last_id) if last_id else None)
|
||||||
|
|
||||||
async def edit_message(
|
|
||||||
self,
|
|
||||||
chat_id: str,
|
|
||||||
message_id: str,
|
|
||||||
content: str,
|
|
||||||
**kwargs,
|
|
||||||
) -> SendResult:
|
|
||||||
"""Edit a previously-sent message. Used by the gateway stream consumer
|
|
||||||
for progressive message updates during agent streaming."""
|
|
||||||
formatted = self.format_message(content)
|
|
||||||
data = await self._api_patch(f"messages/{message_id}", {"content": formatted})
|
|
||||||
if data.get("result") != "success":
|
|
||||||
msg = str(data.get("msg", "edit failed"))
|
|
||||||
logger.debug("Zulip edit_message(%s) -> %s", message_id, msg)
|
|
||||||
return SendResult(success=False, error=msg)
|
|
||||||
return SendResult(success=True, message_id=message_id)
|
|
||||||
|
|
||||||
async def get_chat_info(self, chat_id: str) -> Dict[str, Any]:
|
async def get_chat_info(self, chat_id: str) -> Dict[str, Any]:
|
||||||
kind, to, topic = _parse_target(chat_id, self._default_topic)
|
kind, to, topic = _parse_target(chat_id, self._default_topic)
|
||||||
if kind == "direct":
|
if kind == "direct":
|
||||||
|
|||||||
@@ -77,7 +77,7 @@ SELFTEST_TIMEOUT = 30 # seconds to wait for self-test response
|
|||||||
|
|
||||||
# Gen 4: Connection pool health
|
# Gen 4: Connection pool health
|
||||||
HEARTBEAT_INTERVAL = 300 # Log heartbeat every 5 minutes
|
HEARTBEAT_INTERVAL = 300 # Log heartbeat every 5 minutes
|
||||||
MAX_CONSECUTIVE_EMPTY_POLLS = 500 # Reset client after this many empty polls (~25min at 3s interval). Raised from 20 to prevent idle bots from cycling reconnections.
|
MAX_CONSECUTIVE_EMPTY_POLLS = 20 # Reset client after this many empty polls
|
||||||
CLIENT_REUSE_LIMIT = 500 # Create fresh client every N polls to prevent connection leaks
|
CLIENT_REUSE_LIMIT = 500 # Create fresh client every N polls to prevent connection leaks
|
||||||
POLL_SILENCE_WARN_INTERVAL = 60 # Warn if no events received for this many seconds
|
POLL_SILENCE_WARN_INTERVAL = 60 # Warn if no events received for this many seconds
|
||||||
MAX_ERROR_LOG_LEN = 80 # Truncate error response bodies to avoid HTML spam
|
MAX_ERROR_LOG_LEN = 80 # Truncate error response bodies to avoid HTML spam
|
||||||
@@ -99,23 +99,6 @@ def _build_auth_header(email: str, api_key: str) -> str:
|
|||||||
return f"Basic {token}"
|
return f"Basic {token}"
|
||||||
|
|
||||||
|
|
||||||
def _strip_html(text: str) -> str:
|
|
||||||
"""Strip HTML tags and decode HTML entities from Zulip message content.
|
|
||||||
|
|
||||||
Zulip delivers message content as rendered HTML (e.g. <p>/approve</p>).
|
|
||||||
The gateway slash command parser expects plain text, so HTML tags
|
|
||||||
prevent command matching. This helper strips tags and decodes entities.
|
|
||||||
"""
|
|
||||||
if not text:
|
|
||||||
return text
|
|
||||||
# Strip HTML tags
|
|
||||||
text = re.sub(r"<[^>]+>", "", text)
|
|
||||||
# Decode common HTML entities
|
|
||||||
text = text.replace("&", "&").replace("<", "<").replace(">", ">")
|
|
||||||
text = text.replace(""", '"').replace("'", "'").replace(" ", " ")
|
|
||||||
return text.strip()
|
|
||||||
|
|
||||||
|
|
||||||
def _truncate(text: str, limit: int = MAX_ZULIP_MESSAGE) -> str:
|
def _truncate(text: str, limit: int = MAX_ZULIP_MESSAGE) -> str:
|
||||||
"""Truncate to Zulip's message limit with notice."""
|
"""Truncate to Zulip's message limit with notice."""
|
||||||
if len(text) <= limit:
|
if len(text) <= limit:
|
||||||
@@ -319,9 +302,6 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
# 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()
|
||||||
|
|
||||||
# Gen 5: Ensure bot is subscribed to primary stream
|
|
||||||
await self._ensure_stream_subscriptions()
|
|
||||||
|
|
||||||
self._mark_connected()
|
self._mark_connected()
|
||||||
logger.info(
|
logger.info(
|
||||||
"[%s] Connected to %s as %s (queue=%s, bot_id=%s, all_bots_id=%s)",
|
"[%s] Connected to %s as %s (queue=%s, bot_id=%s, all_bots_id=%s)",
|
||||||
@@ -469,10 +449,6 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
except Exception:
|
except Exception:
|
||||||
pass
|
pass
|
||||||
|
|
||||||
# Clear stale errors on successful poll cycle
|
|
||||||
self._health_stats["last_error_msg"] = None
|
|
||||||
self._health_stats["last_error_at"] = None
|
|
||||||
|
|
||||||
# Gen 4: Heartbeat logging — prove poll loop is alive
|
# Gen 4: Heartbeat logging — prove poll loop is alive
|
||||||
now = time.time()
|
now = time.time()
|
||||||
if now - self._last_heartbeat_time > HEARTBEAT_INTERVAL:
|
if now - self._last_heartbeat_time > HEARTBEAT_INTERVAL:
|
||||||
@@ -701,64 +677,6 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
# Gen 2: Dynamic @all-bots resolution
|
# Gen 2: Dynamic @all-bots resolution
|
||||||
# ------------------------------------------------------------------
|
# ------------------------------------------------------------------
|
||||||
|
|
||||||
async def _ensure_stream_subscriptions(self) -> None:
|
|
||||||
"""Ensure the bot is subscribed to the primary streams.
|
|
||||||
|
|
||||||
Checks current subscriptions and subscribes to the configured primary
|
|
||||||
stream (and general-chat) if missing. Without subscription, the bot
|
|
||||||
receives DMs but no stream events — @mentions in streams are invisible.
|
|
||||||
|
|
||||||
This is a Gen 5 improvement based on lessons from the pi Zulip extension,
|
|
||||||
where the bot had zero stream subscriptions and couldn't receive
|
|
||||||
@mentions in stream topics.
|
|
||||||
"""
|
|
||||||
try:
|
|
||||||
# Check current subscriptions
|
|
||||||
resp, _ = await self._api_call(
|
|
||||||
"GET", "/api/v1/users/me/subscriptions",
|
|
||||||
)
|
|
||||||
if not resp:
|
|
||||||
return
|
|
||||||
|
|
||||||
subscribed = resp.get("subscriptions", [])
|
|
||||||
subscribed_names = [s.get("name", "") for s in subscribed if isinstance(s, dict)]
|
|
||||||
|
|
||||||
# Streams the bot should be subscribed to
|
|
||||||
required_streams = [self._stream, "general chat"]
|
|
||||||
missing = [s for s in required_streams if s not in subscribed_names]
|
|
||||||
|
|
||||||
if not missing:
|
|
||||||
logger.info(
|
|
||||||
"[%s] Stream subscriptions OK: %s",
|
|
||||||
self.name, ", ".join(subscribed_names),
|
|
||||||
)
|
|
||||||
return
|
|
||||||
|
|
||||||
# Subscribe to missing streams
|
|
||||||
import json as _json
|
|
||||||
payload = {"subscriptions": _json.dumps(
|
|
||||||
[{"name": s} for s in missing]
|
|
||||||
)}
|
|
||||||
sub_resp, _ = await self._api_call(
|
|
||||||
"POST", "/api/v1/users/me/subscriptions",
|
|
||||||
data=payload,
|
|
||||||
)
|
|
||||||
if sub_resp and sub_resp.get("result") == "success":
|
|
||||||
subscribed_str = ", ".join(missing)
|
|
||||||
logger.info(
|
|
||||||
"[%s] Subscribed to: %s", self.name, subscribed_str,
|
|
||||||
)
|
|
||||||
self._health_stats["subscriptions_fixed"] = (
|
|
||||||
self._health_stats.get("subscriptions_fixed", 0) + 1
|
|
||||||
)
|
|
||||||
else:
|
|
||||||
logger.warning(
|
|
||||||
"[%s] Failed to subscribe to %s",
|
|
||||||
self.name, ", ".join(missing),
|
|
||||||
)
|
|
||||||
except Exception as e:
|
|
||||||
logger.warning("[%s] Subscription check failed: %s", self.name, e)
|
|
||||||
|
|
||||||
async def _resolve_all_bots_user_id(self) -> None:
|
async def _resolve_all_bots_user_id(self) -> None:
|
||||||
"""Try to resolve the @all-bots user ID from the Zulip server.
|
"""Try to resolve the @all-bots user ID from the Zulip server.
|
||||||
|
|
||||||
@@ -801,8 +719,6 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
sender_name = msg.get("sender_full_name", "Unknown")
|
sender_name = msg.get("sender_full_name", "Unknown")
|
||||||
sender_id = msg.get("sender_id")
|
sender_id = msg.get("sender_id")
|
||||||
content = msg.get("content", "")
|
content = msg.get("content", "")
|
||||||
# Strip Zulip HTML for slash command matching (zulip-approval-fix contract)
|
|
||||||
content = _strip_html(content)
|
|
||||||
|
|
||||||
# Echo-loop prevention: skip own messages
|
# Echo-loop prevention: skip own messages
|
||||||
if sender_email == self._bot_email:
|
if sender_email == self._bot_email:
|
||||||
@@ -1180,35 +1096,6 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
"detail": f"@all-bots user_id: {self._all_bots_user_id}",
|
"detail": f"@all-bots user_id: {self._all_bots_user_id}",
|
||||||
}
|
}
|
||||||
|
|
||||||
# 9. Stream subscriptions (verifies bot is subscribed to at least one stream)
|
|
||||||
try:
|
|
||||||
from urllib.parse import urlencode
|
|
||||||
import json as _json
|
|
||||||
key_str = f"{self._email}:{self._api_key}"
|
|
||||||
auth_b64 = base64.b64encode(key_str.encode()).decode()
|
|
||||||
resp, _ = await self._api_call(
|
|
||||||
"GET", "/api/v1/users/me/subscriptions",
|
|
||||||
)
|
|
||||||
if resp and "subscriptions" in resp:
|
|
||||||
stream_count = len(resp["subscriptions"])
|
|
||||||
stream_names = [s.get("name", "?") for s in resp["subscriptions"]][:5]
|
|
||||||
checks["stream_subscriptions"] = {
|
|
||||||
"status": stream_count > 0,
|
|
||||||
"detail": f"{stream_count} stream(s): {', '.join(stream_names)}"
|
|
||||||
if stream_count > 0
|
|
||||||
else "No stream subscriptions — bot will not receive stream events",
|
|
||||||
}
|
|
||||||
else:
|
|
||||||
checks["stream_subscriptions"] = {
|
|
||||||
"status": False,
|
|
||||||
"detail": "Failed to check subscriptions",
|
|
||||||
}
|
|
||||||
except Exception as sub_err:
|
|
||||||
checks["stream_subscriptions"] = {
|
|
||||||
"status": False,
|
|
||||||
"detail": f"Subscription check failed: {sub_err}",
|
|
||||||
}
|
|
||||||
|
|
||||||
# Overall verdict
|
# Overall verdict
|
||||||
critical = ["connected", "queue_registered", "http_client", "poll_loop"]
|
critical = ["connected", "queue_registered", "http_client", "poll_loop"]
|
||||||
passed = sum(1 for c in checks.values() if c["status"])
|
passed = sum(1 for c in checks.values() if c["status"])
|
||||||
@@ -1424,7 +1311,7 @@ def check_requirements() -> bool:
|
|||||||
"""Check whether the Zulip adapter is installable."""
|
"""Check whether the Zulip adapter is installable."""
|
||||||
if not HTTPX_AVAILABLE:
|
if not HTTPX_AVAILABLE:
|
||||||
return False
|
return False
|
||||||
site = os.getenv("ZULIP_SITE", "").strip() or os.getenv("ZULIP_URL", "").strip()
|
site = os.getenv("ZULIP_SITE", "").strip()
|
||||||
email = os.getenv("ZULIP_EMAIL", "").strip()
|
email = os.getenv("ZULIP_EMAIL", "").strip()
|
||||||
api_key = os.getenv("ZULIP_API_KEY", "").strip()
|
api_key = os.getenv("ZULIP_API_KEY", "").strip()
|
||||||
return bool(site and email and api_key)
|
return bool(site and email and api_key)
|
||||||
@@ -1433,7 +1320,7 @@ def check_requirements() -> bool:
|
|||||||
def validate_config(config) -> bool:
|
def validate_config(config) -> bool:
|
||||||
"""Validate that the configured Zulip platform has credentials."""
|
"""Validate that the configured Zulip platform has credentials."""
|
||||||
extra = getattr(config, "extra", {}) or {}
|
extra = getattr(config, "extra", {}) or {}
|
||||||
site = extra.get("site") or os.getenv("ZULIP_SITE", "") or os.getenv("ZULIP_URL", "")
|
site = extra.get("site") or os.getenv("ZULIP_SITE", "")
|
||||||
email = extra.get("email") or os.getenv("ZULIP_EMAIL", "")
|
email = extra.get("email") or os.getenv("ZULIP_EMAIL", "")
|
||||||
api_key = extra.get("api_key") or os.getenv("ZULIP_API_KEY", "")
|
api_key = extra.get("api_key") or os.getenv("ZULIP_API_KEY", "")
|
||||||
return bool(site and email and api_key)
|
return bool(site and email and api_key)
|
||||||
@@ -1446,7 +1333,7 @@ def is_connected(config) -> bool:
|
|||||||
|
|
||||||
def _env_enablement() -> Optional[dict]:
|
def _env_enablement() -> Optional[dict]:
|
||||||
"""Seeds PlatformConfig.extra from env vars for env-only setups."""
|
"""Seeds PlatformConfig.extra from env vars for env-only setups."""
|
||||||
site = os.getenv("ZULIP_SITE", "").strip() or os.getenv("ZULIP_URL", "").strip()
|
site = os.getenv("ZULIP_SITE", "").strip()
|
||||||
email = os.getenv("ZULIP_EMAIL", "").strip()
|
email = os.getenv("ZULIP_EMAIL", "").strip()
|
||||||
api_key = os.getenv("ZULIP_API_KEY", "").strip()
|
api_key = os.getenv("ZULIP_API_KEY", "").strip()
|
||||||
if not (site and email and api_key):
|
if not (site and email and api_key):
|
||||||
|
|||||||
Reference in New Issue
Block a user