Compare commits
9
Commits
a5bb81b44d
..
v1.0.4
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
6bb438de6e | ||
|
|
b1bef9024f | ||
|
|
dcf2de0052 | ||
|
|
715d54564f | ||
|
|
10e4e0374f | ||
|
|
0acf1068bf | ||
|
|
1c8f48fd77 | ||
|
|
2ee8a4d818 | ||
|
|
397efef776 |
+105
@@ -0,0 +1,105 @@
|
||||
# PR #21 Review — fix(tanko): adapter fixes, event logging, platform-based deploy.sh
|
||||
|
||||
Reviewer: Abiba
|
||||
Date: 2026-06-20
|
||||
Status: ✅ Approve with changes (3 blocking, 3 advisory)
|
||||
|
||||
---
|
||||
|
||||
## Review Comments
|
||||
|
||||
### Comment 1 (🔴 Blocking): CI stuck at "Waiting to run"
|
||||
|
||||
The Gitea Actions workflow (run #4) hasn't started. No CI results available. Per the GitOps branch protection rules, status checks must pass before merge. Do not merge until CI completes and all jobs are green.
|
||||
|
||||
---
|
||||
|
||||
### Comment 2 (🟡 High): `print()` for journald is an anti-pattern
|
||||
|
||||
In `adapter.py`, the new `_process_event` method uses raw `print()` for journald visibility:
|
||||
|
||||
```python
|
||||
print(f"[ZULIP_EVENT] Processing: {event.get('type', 'unknown')}")
|
||||
```
|
||||
|
||||
The existing `logger.info()` calls already flow to journald via stderr when the systemd unit uses `StandardError=journal`. Using raw `print()` bypasses:
|
||||
- Log level filtering
|
||||
- Format consistency with other log output
|
||||
- Future structured logging needs
|
||||
|
||||
**Fix:** Replace with `logger.info(f"[ZULIP_EVENT] Processing: {event.get('type', 'unknown')}")`.
|
||||
|
||||
---
|
||||
|
||||
### Comment 3 (🟡 High): `asyncio.new_event_loop()` leaks on thread restart
|
||||
|
||||
In `_event_loop`, a new event loop is created but never closed:
|
||||
|
||||
```python
|
||||
def _event_loop(self) -> None:
|
||||
loop = asyncio.new_event_loop()
|
||||
asyncio.set_event_loop(loop)
|
||||
try:
|
||||
self._client.call_on_each_message(...)
|
||||
except Exception as e:
|
||||
logger.error(f"Event loop crashed: {e}")
|
||||
self.connected = False
|
||||
```
|
||||
|
||||
If the thread restarts (e.g., reconnection), the old loop is leaked. This accumulates over time.
|
||||
|
||||
**Fix:** Wrap in `try/finally`:
|
||||
|
||||
```python
|
||||
def _event_loop(self) -> None:
|
||||
loop = asyncio.new_event_loop()
|
||||
asyncio.set_event_loop(loop)
|
||||
try:
|
||||
self._client.call_on_each_message(
|
||||
lambda event: self._process_event(event),
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(f"Event loop crashed: {e}")
|
||||
self.connected = False
|
||||
finally:
|
||||
loop.close()
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
### Comment 4 (🟡 Medium): Verify `Client(site=...)` parameter name
|
||||
|
||||
The PR changes `Client(server_url=...)` → `Client(site=...)` for "Python 3.13 compatibility." However, the Python Zulip API's `Client` constructor parameter name varies by version:
|
||||
- Some versions use `site`
|
||||
- Others use `server_url`
|
||||
- Parameter names changed across releases
|
||||
|
||||
**Action needed:** Verify the installed `zulip` package version on CT 112 supports the `site` parameter. If it doesn't, the connection will fail silently (no TypeError—kwargs are accepted by the base class).
|
||||
|
||||
---
|
||||
|
||||
### Comment 5 (🟡 Medium): deploy.sh refactor only tested on Hermes platform
|
||||
|
||||
The PR refactors `deploy.sh` across all 3 platforms (Hermes, Agent Zero, Pi) but testing only covers Tanko (Hermes) on CT 112. The Agent Zero (`pip install -r requirements.txt`) and Pi (`/reload` instead of `systemctl`) code paths are untested.
|
||||
|
||||
**Action needed:** Before merging, run at minimum a dry-run deploy against all 3 platform types:
|
||||
```
|
||||
./scripts/deploy.sh --ct=kagentz main --dry-run
|
||||
./scripts/deploy.sh --ct=abiba main --dry-run
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
### Comment 6 (🟢 Low): Agent Zero pip install assumption
|
||||
|
||||
The deploy.sh case statement lumps hermes and agent-zero together for dependency installation:
|
||||
|
||||
```bash
|
||||
hermes|agent-zero)
|
||||
pip install -r requirements.txt --quiet
|
||||
;;
|
||||
```
|
||||
|
||||
This assumes Agent Zero has the same `requirements.txt` location and content as Hermes. If Agent Zero uses a different dependency file or install method, this will silently install wrong packages.
|
||||
|
||||
**Suggestion:** Add a per-platform dependency install path or document this assumption explicitly.
|
||||
+32
-7
@@ -1,8 +1,20 @@
|
||||
# Zulip Multi-Platform Agent Communication — Architecture
|
||||
# Zulip Multi-Platform Agent Communication — Architecture (v2)
|
||||
|
||||
## Overview
|
||||
A production-ready system enabling 6 AI agents across 3 platforms (Hermes Python, Agent Zero, pi TypeScript) to communicate through Zulip via dedicated per-agent bot users.
|
||||
|
||||
## Architecture (Hermes Native Plugin — Current)
|
||||
|
||||
As of v1.0.0, Hermes agents (Tanko, Mumuni, Koonimo, Koby) use the **Hermes native platform plugin**
|
||||
at `~/.hermes/plugins/platforms/zulip/`. This replaces the old standalone systemd service.
|
||||
|
||||
Benefits of the native plugin:
|
||||
- Extends `BasePlatformAdapter` — zero changes to Hermes core
|
||||
- Auto-registers via `register(ctx)` at Gateway startup
|
||||
- Direct session injection (no subprocess overhead)
|
||||
- Leverages Gateway's built-in health checks, config, and error handling
|
||||
- Unified logging with all other Hermes platform adapters
|
||||
|
||||
## Architecture Diagram
|
||||
```
|
||||
┌─────────────────────────────────────────────────────────────┐
|
||||
@@ -66,18 +78,31 @@ User types: @all-bots status report
|
||||
- **Swarm Development**: The "Swarm" refers to our collaborative development methodology where multiple agents/developers work on different components of the plugin simultaneously.
|
||||
|
||||
|
||||
### Hermes (Python) — BasePlatformAdapter
|
||||
- Path: `hermes-zulip-plugin/src/hermes_zulip/`
|
||||
### Hermes (Python) — BasePlatformAdapter (CURRENT)
|
||||
- Path: `plugins/platforms/zulip/`
|
||||
- Deploy: `~/.hermes/plugins/platforms/zulip/`
|
||||
- Implements: `BasePlatformAdapter` (Hermes Gateway)
|
||||
- Config: `config.yaml` per-agent
|
||||
- Entry point: `plugin.yaml` (Hermes manifest)
|
||||
- Config: Hermes `config.yaml` under `platforms.zulip.extra` or env vars
|
||||
- Entry point: `plugin.yaml` + `register(ctx)` (Hermes manifest)
|
||||
- Versions: Gen 3 (v1.0.0) — 1,169 lines, 14/14 Success Criteria met
|
||||
- Features: DM-first, placeholder→edit streaming, dedup, self-test, health stats, @all-bots resolution
|
||||
|
||||
### Agent Zero — A0 Plugin System
|
||||
### Agent Zero — A0 Plugin System (LEGACY — to migrate)
|
||||
- Path: `agent-zero-plugin/src/`
|
||||
- Implements: Agent Zero plugin API
|
||||
- Config: `config.yaml` per-agent
|
||||
|
||||
### pi (TypeScript) — pi Extension API
|
||||
### pi (TypeScript) — pi Extension API (CURRENT)
|
||||
- Path: `pi-zulip-extension/`
|
||||
- Deploy: PM2-managed process
|
||||
- Runs under `pi --mode rpc --session-id zulip-service`
|
||||
- Active: abiba-bot only (ZULIP_EXTENSION_ACTIVE=true guard)
|
||||
|
||||
### Legacy Hermes Plugin (DEPRECATED)
|
||||
- OLD path: `hermes-zulip-plugin/src/hermes_zulip/`
|
||||
- OLD deploy: `/opt/hermes-zulip-plugin/` + systemd service
|
||||
- Status: Replaced by `plugins/platforms/zulip/` native plugin as of v1.0.0
|
||||
- Migration: See `hermes-zulip-plugin/DEPRECATED.md`
|
||||
- Path: `pi-zulip-extension/src/`
|
||||
- Implements: pi extension (TypeScript module in `~/.pi/agent/extensions/`)
|
||||
- Config: `config.yaml` per-agent
|
||||
|
||||
@@ -0,0 +1,40 @@
|
||||
# DEPRECATED — Hermes Zulip Plugin (Legacy)
|
||||
|
||||
**Status:** Deprecated as of 2026-06-25
|
||||
**Replaced by:** `plugins/platforms/zulip/` (Hermes native platform plugin)
|
||||
|
||||
This directory contains the original standalone Hermes Zulip plugin that ran as a
|
||||
systemd service with its own poll loop, subprocess agent invocation, and health
|
||||
server.
|
||||
|
||||
## Why Deprecated
|
||||
|
||||
The new `plugins/platforms/zulip/` plugin is a proper Hermes platform plugin that:
|
||||
|
||||
- Extends `BasePlatformAdapter` from the Hermes Gateway
|
||||
- Auto-registers via the Hermes plugin system (`register(ctx)`)
|
||||
- Uses direct session injection (no subprocess overhead)
|
||||
- Leverages Gateway's built-in health checks, config management, and error handling
|
||||
- Requires zero changes to core Hermes code
|
||||
|
||||
## Migration
|
||||
|
||||
Remove the old systemd service and deploy the new plugin:
|
||||
|
||||
```bash
|
||||
# 1. Stop old service
|
||||
systemctl stop zulip-plugin
|
||||
systemctl disable zulip-plugin
|
||||
|
||||
# 2. Install plugin to ~/.hermes/plugins/platforms/zulip/
|
||||
cp -r plugins/platforms/zulip/ ~/.hermes/plugins/platforms/zulip/
|
||||
|
||||
# 3. Restart Hermes Gateway
|
||||
hermes gateway restart
|
||||
```
|
||||
|
||||
Config moves from:
|
||||
- `/opt/hermes-zulip-plugin/config.yaml` → Hermes `config.yaml` under `platforms: zulip:`
|
||||
- Or set `ZULIP_SITE`, `ZULIP_EMAIL`, `ZULIP_API_KEY` env vars
|
||||
|
||||
See `plugins/platforms/zulip/plugin.yaml` for configuration reference.
|
||||
@@ -1,616 +0,0 @@
|
||||
"""Zulip platform adapter (Hermes plugin).
|
||||
|
||||
Connects to a Zulip server via the Zulip Python SDK, registers an event
|
||||
queue for real-time message delivery, and polls for incoming events.
|
||||
DM-first architecture: private messages route directly into the Hermes
|
||||
agent's session, maintaining conversation continuity with TUI/CLI.
|
||||
|
||||
Configuration in config.yaml::
|
||||
|
||||
platforms:
|
||||
zulip:
|
||||
enabled: true
|
||||
extra:
|
||||
email: "tanko-bot@chat.sysloggh.net"
|
||||
api_key: "${ZULIP_API_KEY}"
|
||||
site: "https://chat.sysloggh.net"
|
||||
poll_interval_ms: 3000
|
||||
|
||||
Environment variables (env wins over config.yaml):
|
||||
ZULIP_EMAIL Bot email (required)
|
||||
ZULIP_API_KEY Bot API key (required)
|
||||
ZULIP_SITE Server URL (required)
|
||||
ZULIP_ALLOWED_USERS Comma-separated user emails allowed (optional)
|
||||
ZULIP_ALLOW_ALL_USERS Allow any user (optional)
|
||||
ZULIP_HOME_CHANNEL Default recipient for cron/notifications (optional)
|
||||
ZULIP_HOME_CHANNEL_NAME Human label for home channel (optional)
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import time
|
||||
import uuid
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
try:
|
||||
import zulip
|
||||
ZULIP_SDK_AVAILABLE = True
|
||||
except ImportError:
|
||||
ZULIP_SDK_AVAILABLE = False
|
||||
zulip = None # type: ignore[assignment]
|
||||
|
||||
from gateway.config import Platform, PlatformConfig
|
||||
from gateway.platforms.base import (
|
||||
BasePlatformAdapter,
|
||||
MessageEvent,
|
||||
MessageType,
|
||||
SendResult,
|
||||
)
|
||||
|
||||
# ── Constants ──────────────────────────────────────────────────────────────
|
||||
|
||||
DEFAULT_POLL_INTERVAL_MS = 3000
|
||||
MAX_MESSAGE_LENGTH = 10000 # Zulip max message body length
|
||||
RECONNECT_BACKOFF = [2, 5, 10, 30, 60]
|
||||
DEDUP_WINDOW_SECONDS = 300
|
||||
DEDUP_MAX_SIZE = 1000
|
||||
|
||||
PLACEHOLDERS = [
|
||||
":robot: _Processing your message..._",
|
||||
":hourglass_flowing_sand: _Thinking..._",
|
||||
":brain: _Generating response..._",
|
||||
]
|
||||
|
||||
|
||||
# ── Plugin registration helpers ────────────────────────────────────────────
|
||||
|
||||
|
||||
def check_requirements() -> bool:
|
||||
"""Check whether the Zulip SDK is available and minimally configured."""
|
||||
if not ZULIP_SDK_AVAILABLE:
|
||||
return False
|
||||
email = os.getenv("ZULIP_EMAIL", "").strip()
|
||||
api_key = os.getenv("ZULIP_API_KEY", "").strip()
|
||||
site = os.getenv("ZULIP_SITE", "").strip()
|
||||
return bool(email and api_key and site)
|
||||
|
||||
|
||||
def validate_config(config) -> bool:
|
||||
"""Validate Zulip platform has email and site configured."""
|
||||
extra = getattr(config, "extra", {}) or {}
|
||||
email = extra.get("email") or os.getenv("ZULIP_EMAIL", "")
|
||||
site = extra.get("site") or os.getenv("ZULIP_SITE", "")
|
||||
return bool(email and site)
|
||||
|
||||
|
||||
def is_connected(config) -> bool:
|
||||
"""Check whether Zulip is configured (env or config.yaml)."""
|
||||
extra = getattr(config, "extra", {}) or {}
|
||||
email = os.getenv("ZULIP_EMAIL") or extra.get("email", "")
|
||||
site = os.getenv("ZULIP_SITE") or extra.get("site", "")
|
||||
return bool(email and site)
|
||||
|
||||
|
||||
def _env_enablement() -> dict | None:
|
||||
"""Seed PlatformConfig.extra from env vars during gateway config load.
|
||||
|
||||
Returns None when Zulip isn't minimally configured.
|
||||
"""
|
||||
email = os.getenv("ZULIP_EMAIL", "").strip()
|
||||
api_key = os.getenv("ZULIP_API_KEY", "").strip()
|
||||
site = os.getenv("ZULIP_SITE", "").strip()
|
||||
if not (email and api_key and site):
|
||||
return None
|
||||
seed: dict = {"email": email, "site": site}
|
||||
home = os.getenv("ZULIP_HOME_CHANNEL", "").strip()
|
||||
if home:
|
||||
seed["home_channel"] = {
|
||||
"chat_id": home,
|
||||
"name": os.getenv("ZULIP_HOME_CHANNEL_NAME", home),
|
||||
}
|
||||
return seed
|
||||
|
||||
|
||||
async def _standalone_send(
|
||||
pconfig,
|
||||
chat_id: str,
|
||||
message: str,
|
||||
*,
|
||||
thread_id: Optional[str] = None,
|
||||
media_files: Optional[List[str]] = None,
|
||||
force_document: bool = False,
|
||||
) -> Dict[str, Any]:
|
||||
"""Out-of-process send for cron / send_message_tool fallbacks."""
|
||||
if not ZULIP_SDK_AVAILABLE:
|
||||
return {"error": "zulip standalone send: zulip SDK not installed"}
|
||||
extra = getattr(pconfig, "extra", {}) or {}
|
||||
email = extra.get("email") or os.getenv("ZULIP_EMAIL", "")
|
||||
api_key = extra.get("api_key") or os.getenv("ZULIP_API_KEY", "")
|
||||
site = extra.get("site") or os.getenv("ZULIP_SITE", "")
|
||||
if not (email and api_key and site):
|
||||
return {"error": "zulip standalone send: ZULIP_EMAIL, ZULIP_API_KEY, ZULIP_SITE required"}
|
||||
try:
|
||||
client = zulip.Client(email=email, api_key=api_key, site=site)
|
||||
result = client.send_message({
|
||||
"type": "private",
|
||||
"to": chat_id,
|
||||
"content": message[:MAX_MESSAGE_LENGTH],
|
||||
})
|
||||
if result.get("result") == "success":
|
||||
return {
|
||||
"success": True,
|
||||
"platform": "zulip",
|
||||
"chat_id": chat_id,
|
||||
"message_id": str(result.get("id", "")),
|
||||
}
|
||||
return {"error": f"zulip send failed: {result.get('msg', 'unknown')}"}
|
||||
except Exception as e:
|
||||
return {"error": f"zulip standalone send failed: {e}"}
|
||||
|
||||
|
||||
# ── Adapter ────────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
class ZulipAdapter(BasePlatformAdapter):
|
||||
"""Zulip platform adapter.
|
||||
|
||||
Connects to Zulip, registers a real-time event queue, and polls for
|
||||
incoming messages. DMs are routed into the Hermes agent's session.
|
||||
"""
|
||||
|
||||
MAX_MESSAGE_LENGTH = MAX_MESSAGE_LENGTH
|
||||
|
||||
def __init__(self, config: PlatformConfig):
|
||||
platform = Platform("zulip")
|
||||
super().__init__(config=config, platform=platform)
|
||||
extra = config.extra or {}
|
||||
|
||||
# Core Zulip config — env overrides config.yaml
|
||||
self._email: str = os.getenv("ZULIP_EMAIL") or extra.get("email", "")
|
||||
self._api_key: str = os.getenv("ZULIP_API_KEY") or extra.get("api_key", "")
|
||||
self._site: str = (os.getenv("ZULIP_SITE") or extra.get("site", "")).rstrip("/")
|
||||
|
||||
# Polling config
|
||||
self._poll_interval: float = (
|
||||
(extra.get("poll_interval_ms") or DEFAULT_POLL_INTERVAL_MS) / 1000.0
|
||||
)
|
||||
|
||||
# State
|
||||
self._client: Optional["zulip.Client"] = None
|
||||
self._queue_id: Optional[str] = None
|
||||
self._last_event_id: int = -1
|
||||
self._poll_task: Optional[asyncio.Task] = None
|
||||
|
||||
# Message deduplication: event_id -> timestamp
|
||||
self._seen_messages: Dict[str, float] = {}
|
||||
|
||||
# Pending replies for placeholder -> edit streaming.
|
||||
# Maps chat_id -> {placeholder_msg_id, ...}
|
||||
self._pending_replies: Dict[str, Dict[str, Any]] = {}
|
||||
|
||||
# ── Connection lifecycle ────────────────────────────────────────────
|
||||
|
||||
async def connect(self) -> bool:
|
||||
"""Connect to Zulip by registering an event queue."""
|
||||
if not ZULIP_SDK_AVAILABLE:
|
||||
logger.warning(
|
||||
"[%s] zulip SDK not installed. Run: pip install zulip", self.name
|
||||
)
|
||||
return False
|
||||
if not (self._email and self._api_key and self._site):
|
||||
logger.warning(
|
||||
"[%s] ZULIP_EMAIL, ZULIP_API_KEY, ZULIP_SITE not all configured",
|
||||
self.name,
|
||||
)
|
||||
return False
|
||||
|
||||
try:
|
||||
self._client = zulip.Client(
|
||||
email=self._email,
|
||||
api_key=self._api_key,
|
||||
site=self._site,
|
||||
client="Hermes-Zulip-Plugin/1.0.0",
|
||||
)
|
||||
|
||||
# Register event queue for message events
|
||||
queue_result = await asyncio.get_event_loop().run_in_executor(
|
||||
None,
|
||||
lambda: self._client.register(
|
||||
event_types=["message"],
|
||||
fetch_event_types=["message"],
|
||||
),
|
||||
)
|
||||
self._queue_id = queue_result.get("queue_id")
|
||||
self._last_event_id = queue_result.get("last_event_id", -1)
|
||||
|
||||
if not self._queue_id:
|
||||
logger.error(
|
||||
"[%s] Queue registration failed: %s",
|
||||
self.name, queue_result,
|
||||
)
|
||||
return False
|
||||
|
||||
self._poll_task = asyncio.create_task(self._poll_loop())
|
||||
self._mark_connected()
|
||||
|
||||
logger.info(
|
||||
"[%s] Connected — %s on %s, queue=%s last_event=%s",
|
||||
self.name, self._email, self._site,
|
||||
self._queue_id, self._last_event_id,
|
||||
)
|
||||
return True
|
||||
|
||||
except Exception as e:
|
||||
logger.error("[%s] Failed to connect: %s", self.name, e)
|
||||
return False
|
||||
|
||||
async def disconnect(self) -> bool:
|
||||
"""Disconnect from Zulip."""
|
||||
self._running = False
|
||||
self._mark_disconnected()
|
||||
if self._poll_task:
|
||||
self._poll_task.cancel()
|
||||
try:
|
||||
await self._poll_task
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
self._poll_task = None
|
||||
if self._client and self._queue_id:
|
||||
try:
|
||||
await asyncio.get_event_loop().run_in_executor(
|
||||
None, lambda: self._client.deregister(self._queue_id)
|
||||
)
|
||||
except Exception:
|
||||
pass
|
||||
self._client = None
|
||||
self._queue_id = None
|
||||
self._seen_messages.clear()
|
||||
self._pending_replies.clear()
|
||||
logger.info("[%s] Disconnected", self.name)
|
||||
return True
|
||||
|
||||
# ── Polling loop ────────────────────────────────────────────────────
|
||||
|
||||
async def _poll_loop(self) -> None:
|
||||
"""Poll the Zulip event queue for new messages, with reconnect."""
|
||||
backoff_idx = 0
|
||||
loop_start: float = 0.0
|
||||
|
||||
while self._running:
|
||||
try:
|
||||
loop_start = time.monotonic()
|
||||
await self._poll_once()
|
||||
if time.monotonic() - loop_start < 60.0:
|
||||
backoff_idx = 0
|
||||
await asyncio.sleep(self._poll_interval)
|
||||
|
||||
except asyncio.CancelledError:
|
||||
return
|
||||
except Exception as e:
|
||||
if not self._running:
|
||||
return
|
||||
logger.warning("[%s] Poll error: %s", self.name, e)
|
||||
if "BAD_EVENT_QUEUE_ID" in str(e):
|
||||
logger.info("[%s] Queue expired, re-registering...", self.name)
|
||||
if await self._reregister_queue():
|
||||
backoff_idx = 0
|
||||
continue
|
||||
delay = RECONNECT_BACKOFF[
|
||||
min(backoff_idx, len(RECONNECT_BACKOFF) - 1)
|
||||
]
|
||||
logger.info("[%s] Retrying in %ds...", self.name, delay)
|
||||
await asyncio.sleep(delay)
|
||||
backoff_idx += 1
|
||||
|
||||
async def _poll_once(self) -> None:
|
||||
"""Retrieve and process one batch of events."""
|
||||
if not self._client or not self._queue_id:
|
||||
return
|
||||
|
||||
data = await asyncio.get_event_loop().run_in_executor(
|
||||
None,
|
||||
lambda: self._client.get_events(
|
||||
queue_id=self._queue_id,
|
||||
last_event_id=self._last_event_id,
|
||||
),
|
||||
)
|
||||
|
||||
if not data or not isinstance(data, dict):
|
||||
return
|
||||
|
||||
if data.get("result") == "error":
|
||||
msg = data.get("msg", "")
|
||||
if "BAD_EVENT_QUEUE_ID" in msg or "queue_id" in msg.lower():
|
||||
raise RuntimeError(f"BAD_EVENT_QUEUE_ID: {msg}")
|
||||
return
|
||||
|
||||
events = data.get("events", [])
|
||||
if not events:
|
||||
return
|
||||
|
||||
for event in events:
|
||||
event_id = event.get("id", 0)
|
||||
if event_id > self._last_event_id:
|
||||
self._last_event_id = event_id
|
||||
await self._on_event(event)
|
||||
|
||||
async def _reregister_queue(self) -> bool:
|
||||
"""Re-register the event queue after expiry."""
|
||||
try:
|
||||
queue_result = await asyncio.get_event_loop().run_in_executor(
|
||||
None,
|
||||
lambda: self._client.register(
|
||||
event_types=["message"],
|
||||
fetch_event_types=["message"],
|
||||
),
|
||||
)
|
||||
self._queue_id = queue_result.get("queue_id")
|
||||
self._last_event_id = queue_result.get("last_event_id", -1)
|
||||
if self._queue_id:
|
||||
logger.info(
|
||||
"[%s] Queue re-registered: %s", self.name, self._queue_id
|
||||
)
|
||||
return True
|
||||
except Exception as e:
|
||||
logger.error("[%s] Queue re-registration failed: %s", self.name, e)
|
||||
return False
|
||||
|
||||
# ── Event processing ────────────────────────────────────────────────
|
||||
|
||||
async def _on_event(self, event: Dict[str, Any]) -> None:
|
||||
"""Process a single Zulip event."""
|
||||
if event.get("type") != "message":
|
||||
return
|
||||
|
||||
msg = event.get("message", {})
|
||||
if not msg:
|
||||
return
|
||||
|
||||
# Skip own messages (echo loop prevention)
|
||||
sender_email = msg.get("sender_email", "")
|
||||
if sender_email == self._email:
|
||||
return
|
||||
|
||||
# Deduplicate
|
||||
event_id = str(event.get("id", ""))
|
||||
if self._is_duplicate(event_id):
|
||||
return
|
||||
|
||||
# DM-first architecture (ADR-001): only process private messages
|
||||
if msg.get("type") == "private":
|
||||
await self._on_dm(msg)
|
||||
return
|
||||
|
||||
async def _on_dm(self, msg: Dict[str, Any]) -> None:
|
||||
"""Process an incoming DM and dispatch to the gateway."""
|
||||
sender_id = str(msg.get("sender_id", ""))
|
||||
sender_email = str(msg.get("sender_email", "unknown"))
|
||||
sender_name = str(msg.get("sender_full_name", sender_email))
|
||||
content = str(msg.get("content", "")).strip()
|
||||
|
||||
if not content or not sender_id:
|
||||
return
|
||||
|
||||
logger.info(
|
||||
"[%s] DM from %s (%s): %.60s",
|
||||
self.name, sender_name, sender_email, content,
|
||||
)
|
||||
|
||||
# Send a placeholder for streaming UX
|
||||
placeholder_msg_id = await self._send_placeholder(sender_id)
|
||||
|
||||
# Register pending reply so send_message can edit the placeholder
|
||||
self._pending_replies[sender_id] = {
|
||||
"placeholder_msg_id": placeholder_msg_id,
|
||||
"type": "private",
|
||||
}
|
||||
|
||||
# Fire-and-forget typing indicator
|
||||
asyncio.ensure_future(
|
||||
self._send_typing_indicator(int(sender_id), "start")
|
||||
)
|
||||
|
||||
# Build the session source
|
||||
source = self.build_source(
|
||||
chat_id=sender_id,
|
||||
chat_name=sender_email,
|
||||
chat_type="dm",
|
||||
user_id=sender_email,
|
||||
user_name=sender_name,
|
||||
message_id=str(msg.get("id", "")),
|
||||
)
|
||||
|
||||
# Build the MessageEvent for the gateway
|
||||
now = datetime.now(tz=timezone.utc)
|
||||
message_event = MessageEvent(
|
||||
text=content,
|
||||
message_type=MessageType.TEXT,
|
||||
source=source,
|
||||
message_id=str(msg.get("id", "")),
|
||||
raw_message=msg,
|
||||
timestamp=now,
|
||||
)
|
||||
|
||||
logger.debug("[%s] Dispatching DM from %s", self.name, sender_email)
|
||||
await self.handle_message(message_event)
|
||||
|
||||
def _is_duplicate(self, msg_id: str) -> bool:
|
||||
"""Dedup check with sliding window."""
|
||||
now = time.time()
|
||||
stale = [
|
||||
mid for mid, ts in self._seen_messages.items()
|
||||
if now - ts > DEDUP_WINDOW_SECONDS
|
||||
]
|
||||
for mid in stale:
|
||||
del self._seen_messages[mid]
|
||||
if len(self._seen_messages) > DEDUP_MAX_SIZE:
|
||||
oldest = sorted(
|
||||
self._seen_messages.keys(), key=lambda k: self._seen_messages[k]
|
||||
)[:100]
|
||||
for k in oldest:
|
||||
del self._seen_messages[k]
|
||||
if msg_id in self._seen_messages:
|
||||
return True
|
||||
self._seen_messages[msg_id] = now
|
||||
return False
|
||||
|
||||
# ── Message delivery ────────────────────────────────────────────────
|
||||
|
||||
async def send_message(
|
||||
self,
|
||||
chat_id: str,
|
||||
content: str,
|
||||
*,
|
||||
reply_to: Optional[str] = None,
|
||||
metadata: Optional[Dict[str, Any]] = None,
|
||||
edit_message_id: Optional[str] = None,
|
||||
) -> SendResult:
|
||||
"""Send a DM reply.
|
||||
|
||||
If a pending placeholder exists for this chat, edits it with the
|
||||
final response (streaming UX). Otherwise sends a fresh message.
|
||||
"""
|
||||
# Check for pending placeholder to finalize
|
||||
pending = self._pending_replies.pop(chat_id, None)
|
||||
if pending and pending.get("placeholder_msg_id"):
|
||||
return await self._edit_message(
|
||||
pending["placeholder_msg_id"],
|
||||
content[:MAX_MESSAGE_LENGTH],
|
||||
chat_id,
|
||||
)
|
||||
|
||||
return await self._send_fresh(chat_id, content)
|
||||
|
||||
async def send_typing(self, chat_id: str, metadata=None) -> None:
|
||||
"""Send typing indicator. No-op; managed via event flow."""
|
||||
pass
|
||||
|
||||
async def get_chat_info(self, chat_id: str) -> Dict[str, Any]:
|
||||
"""Return basic info about a Zulip chat."""
|
||||
return {"name": chat_id, "type": "dm"}
|
||||
|
||||
# ── Internal helpers ────────────────────────────────────────────────
|
||||
|
||||
async def _send_placeholder(self, chat_id: str) -> Optional[str]:
|
||||
"""Send a 'Thinking...' placeholder message. Returns message ID."""
|
||||
if not self._client:
|
||||
return None
|
||||
placeholder = PLACEHOLDERS[
|
||||
self._next_reply_id % len(PLACEHOLDERS)
|
||||
]
|
||||
try:
|
||||
result = await asyncio.get_event_loop().run_in_executor(
|
||||
None,
|
||||
lambda: self._client.send_message({
|
||||
"type": "private",
|
||||
"to": chat_id,
|
||||
"content": placeholder,
|
||||
}),
|
||||
)
|
||||
if result.get("result") == "success":
|
||||
msg_id = str(result.get("id", ""))
|
||||
logger.debug("[%s] Placeholder sent to %s (msg_id=%s)", self.name, chat_id, msg_id)
|
||||
return msg_id
|
||||
except Exception as e:
|
||||
logger.debug("[%s] Placeholder send failed: %s", self.name, e)
|
||||
return None
|
||||
|
||||
async def _send_typing_indicator(self, user_id: int, operation: str) -> None:
|
||||
"""Send typing start/stop via Zulip API."""
|
||||
if not self._client:
|
||||
return
|
||||
try:
|
||||
await asyncio.get_event_loop().run_in_executor(
|
||||
None,
|
||||
lambda: self._client.send_typing_notification(
|
||||
recipients=[user_id],
|
||||
operation=operation,
|
||||
),
|
||||
)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
async def _edit_message(
|
||||
self, message_id: str, content: str, chat_id: str
|
||||
) -> SendResult:
|
||||
"""Edit an existing Zulip message (finalize a placeholder)."""
|
||||
if not self._client:
|
||||
return SendResult(success=False, error="Client not connected")
|
||||
try:
|
||||
await asyncio.get_event_loop().run_in_executor(
|
||||
None,
|
||||
lambda: self._client.update_message({
|
||||
"message_id": message_id,
|
||||
"content": content,
|
||||
}),
|
||||
)
|
||||
return SendResult(success=True, message_id=message_id)
|
||||
except Exception as e:
|
||||
logger.warning(
|
||||
"[%s] Edit failed, sending fresh: %s", self.name, e
|
||||
)
|
||||
return await self._send_fresh(chat_id, content)
|
||||
|
||||
async def _send_fresh(self, chat_id: str, content: str) -> SendResult:
|
||||
"""Send a fresh DM to Zulip."""
|
||||
if not self._client:
|
||||
return SendResult(success=False, error="Client not connected")
|
||||
try:
|
||||
truncated = content[:MAX_MESSAGE_LENGTH]
|
||||
result = await asyncio.get_event_loop().run_in_executor(
|
||||
None,
|
||||
lambda: self._client.send_message({
|
||||
"type": "private",
|
||||
"to": chat_id,
|
||||
"content": truncated,
|
||||
}),
|
||||
)
|
||||
if result.get("result") == "success":
|
||||
return SendResult(
|
||||
success=True,
|
||||
message_id=str(result.get("id", "")),
|
||||
)
|
||||
return SendResult(
|
||||
success=False,
|
||||
error=result.get("msg", "Send failed"),
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error("[%s] Send error: %s", self.name, e)
|
||||
return SendResult(success=False, error=str(e))
|
||||
|
||||
|
||||
# ── Plugin registration ────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def register(ctx) -> None:
|
||||
"""Plugin entry point — called by the Hermes plugin system at startup."""
|
||||
ctx.register_platform(
|
||||
name="zulip",
|
||||
label="Zulip",
|
||||
adapter_factory=lambda cfg: ZulipAdapter(cfg),
|
||||
check_fn=check_requirements,
|
||||
validate_config=validate_config,
|
||||
is_connected=is_connected,
|
||||
required_env=["ZULIP_EMAIL", "ZULIP_API_KEY", "ZULIP_SITE"],
|
||||
install_hint="pip install zulip httpx",
|
||||
env_enablement_fn=_env_enablement,
|
||||
cron_deliver_env_var="ZULIP_HOME_CHANNEL",
|
||||
standalone_sender_fn=_standalone_send,
|
||||
allowed_users_env="ZULIP_ALLOWED_USERS",
|
||||
allow_all_env="ZULIP_ALLOW_ALL_USERS",
|
||||
max_message_length=MAX_MESSAGE_LENGTH,
|
||||
emoji="💬",
|
||||
pii_safe=False,
|
||||
allow_update_command=True,
|
||||
platform_hint=(
|
||||
"You are communicating via Zulip direct messages. "
|
||||
"Your responses are delivered through Zulip's API. "
|
||||
"Use markdown formatting for rich responses. "
|
||||
f"Keep messages under {MAX_MESSAGE_LENGTH} characters (Zulip limit)."
|
||||
),
|
||||
)
|
||||
@@ -1,47 +0,0 @@
|
||||
name: zulip-platform
|
||||
label: Zulip
|
||||
kind: platform
|
||||
version: 1.0.0
|
||||
description: >
|
||||
Zulip messaging gateway adapter for Hermes Agent.
|
||||
Connects to any Zulip server, registers an event queue, and polls for
|
||||
incoming messages. Handles DMs first (ADR-001), with @mention support
|
||||
planned. Uses the Zulip Python SDK for API access.
|
||||
|
||||
DM-first architecture: private messages route directly into the Hermes
|
||||
agent's session, so the same personality and conversation continuity
|
||||
is maintained between TUI, CLI, and Zulip.
|
||||
|
||||
Features: event queue polling, typing indicators, placeholder→edit
|
||||
streaming, exponential backoff reconnection, health endpoint.
|
||||
author: Syslog Solution LLC
|
||||
requires_env:
|
||||
- name: ZULIP_EMAIL
|
||||
description: "Zulip bot email address (e.g. tanko-bot@chat.sysloggh.net)"
|
||||
prompt: "Zulip bot email"
|
||||
password: false
|
||||
- name: ZULIP_API_KEY
|
||||
description: "Zulip bot API key"
|
||||
prompt: "Zulip API key"
|
||||
password: true
|
||||
- name: ZULIP_SITE
|
||||
description: "Zulip server URL (e.g. https://chat.sysloggh.net)"
|
||||
prompt: "Zulip server URL"
|
||||
password: false
|
||||
optional_env:
|
||||
- name: ZULIP_ALLOWED_USERS
|
||||
description: "Comma-separated Zulip user IDs or emails allowed to DM the agent"
|
||||
prompt: "Allowed users (comma-separated emails or IDs)"
|
||||
password: false
|
||||
- name: ZULIP_ALLOW_ALL_USERS
|
||||
description: "Allow any user to DM the bot (disables allowlist)"
|
||||
prompt: "Allow all users? (true/false)"
|
||||
password: false
|
||||
- name: ZULIP_HOME_CHANNEL
|
||||
description: "Default recipient for cron / notification delivery"
|
||||
prompt: "Home channel user email (or empty)"
|
||||
password: false
|
||||
- name: ZULIP_HOME_CHANNEL_NAME
|
||||
description: "Human label for the home channel"
|
||||
prompt: "Home channel display name"
|
||||
password: false
|
||||
@@ -1,2 +1,3 @@
|
||||
zulip>=0.9.0
|
||||
httpx>=0.27.0
|
||||
zulip-bots>=0.9.0
|
||||
pyyaml>=6.0
|
||||
|
||||
@@ -0,0 +1,146 @@
|
||||
"""
|
||||
Hermes Zulip Plugin — Core adapter for Hermes Python agents
|
||||
|
||||
Connects Hermes agents (Tanko, Mumuni, Koonimo, Koby) to the Sysloggh
|
||||
Zulip agent mesh via event queue polling + subprocess agent invocation.
|
||||
|
||||
Architecture mirrors pi-zulip-extension:
|
||||
Zulip event queue → poll loop → parse message → spawn agent subprocess
|
||||
→ capture stdout → post response to Zulip (with placeholder + edit)
|
||||
|
||||
Usage (systemd service):
|
||||
python3 -m hermes_zulip --config /path/to/config.yaml
|
||||
|
||||
Path: hermes-zulip-plugin/src/hermes_zulip/
|
||||
Config: config.yaml (per-agent, deployed alongside plugin)
|
||||
|
||||
@see ADR-001 DM-first, ADR-005 @mention detection
|
||||
@see ADR-006 @all-bots, ADR-009 error handling
|
||||
"""
|
||||
|
||||
import argparse
|
||||
import logging
|
||||
import sys
|
||||
from typing import Dict, Any
|
||||
|
||||
# Ensure the src directory is on the path for the adapter import
|
||||
import os
|
||||
import sys
|
||||
sys.path.insert(0, os.path.join(os.path.dirname(__file__), ".."))
|
||||
|
||||
from hermes_zulip.adapter import ZulipAdapter
|
||||
|
||||
__version__ = "0.2.0"
|
||||
|
||||
|
||||
def load_config(config_path: str) -> Dict[str, Any]:
|
||||
"""Load YAML config and inject env vars."""
|
||||
import yaml
|
||||
|
||||
with open(config_path) as f:
|
||||
cfg = yaml.safe_load(f)
|
||||
|
||||
# Override api_key from env if set
|
||||
env_key = os.environ.get("ZULIP_API_KEY")
|
||||
if env_key:
|
||||
cfg.setdefault("zulip", {})["api_key"] = env_key
|
||||
|
||||
return cfg
|
||||
|
||||
|
||||
def setup_logging(cfg: Dict[str, Any]) -> None:
|
||||
"""Configure logging from config."""
|
||||
level = (
|
||||
cfg.get("monitoring", {}).get("log_level", "INFO").upper()
|
||||
)
|
||||
logging.basicConfig(
|
||||
level=getattr(logging, level, logging.INFO),
|
||||
format="%(asctime)s [%(name)s] %(levelname)s: %(message)s",
|
||||
datefmt="%Y-%m-%d %H:%M:%S",
|
||||
)
|
||||
|
||||
|
||||
def start_health_server(cfg: Dict[str, Any], adapter: ZulipAdapter) -> None:
|
||||
"""Start a minimal health HTTP endpoint in a background thread."""
|
||||
import threading
|
||||
from http.server import HTTPServer, BaseHTTPRequestHandler
|
||||
|
||||
port = cfg.get("monitoring", {}).get("health_port", 9200)
|
||||
start_time = time.time()
|
||||
|
||||
class HealthHandler(BaseHTTPRequestHandler):
|
||||
def do_GET(self) -> None:
|
||||
if self.path in ("/", "/health"):
|
||||
import json
|
||||
body = json.dumps({
|
||||
"status": "ok",
|
||||
"agent": cfg.get("agent", {}).get("name", "unknown"),
|
||||
"connected": adapter.connected,
|
||||
"uptime_seconds": int(time.time() - start_time),
|
||||
"timestamp": datetime.utcnow().isoformat(),
|
||||
})
|
||||
self.send_response(200)
|
||||
self.send_header("Content-Type", "application/json")
|
||||
self.end_headers()
|
||||
self.wfile.write(body.encode())
|
||||
else:
|
||||
self.send_response(404)
|
||||
self.end_headers()
|
||||
|
||||
def log_message(self, fmt, *args):
|
||||
pass # Suppress health check log spam
|
||||
|
||||
server = HTTPServer(("127.0.0.1", port), HealthHandler)
|
||||
thread = threading.Thread(target=server.serve_forever, daemon=True)
|
||||
thread.start()
|
||||
logging.getLogger(__name__).info(f"Health endpoint on :{port}")
|
||||
|
||||
|
||||
def main() -> None:
|
||||
"""CLI entry point for the Hermes Zulip plugin.
|
||||
|
||||
Usage:
|
||||
python3 -m hermes_zulip --config /path/to/config.yaml
|
||||
"""
|
||||
import time
|
||||
from datetime import datetime
|
||||
|
||||
parser = argparse.ArgumentParser(
|
||||
description="Hermes Zulip Plugin — Connects Hermes agents to Zulip"
|
||||
)
|
||||
parser.add_argument(
|
||||
"--config", "-c",
|
||||
default="config.yaml",
|
||||
help="Path to config.yaml (default: config.yaml)",
|
||||
)
|
||||
args = parser.parse_args()
|
||||
|
||||
# Load config
|
||||
cfg = load_config(args.config)
|
||||
setup_logging(cfg)
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
logger.info(
|
||||
f"Starting Hermes Zulip Plugin v{__version__} "
|
||||
f"for {cfg.get('agent', {}).get('name', 'unknown')}"
|
||||
)
|
||||
|
||||
# Create adapter
|
||||
adapter = ZulipAdapter(cfg)
|
||||
|
||||
# Start health endpoint (if enabled)
|
||||
if cfg.get("monitoring", {}).get("health_endpoint_enabled", True):
|
||||
start_health_server(cfg, adapter)
|
||||
|
||||
# Run the event loop (blocks until interrupted)
|
||||
try:
|
||||
adapter.run()
|
||||
except KeyboardInterrupt:
|
||||
logger.info("Shutting down...")
|
||||
except Exception as e:
|
||||
logger.error(f"Fatal error: {e}")
|
||||
sys.exit(1)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -0,0 +1,523 @@
|
||||
"""
|
||||
Core Zulip adapter for Hermes agents (Tanko, Mumuni, Koonimo, Koby).
|
||||
|
||||
Architecture (mirrors pi-zulip-extension):
|
||||
Zulip event queue → poll loop → parse message → spawn agent subprocess
|
||||
→ capture stdout → post response to Zulip
|
||||
|
||||
The Hermes agent is invoked as a subprocess (configurable command) so the
|
||||
plugin remains decoupled from the agent runtime. In a future iteration,
|
||||
this could switch to a direct IPC/socket if the agent exposes one.
|
||||
|
||||
See ADR-005 (@mention detection), ADR-006 (@all-bots), ADR-009 (error handling).
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import re
|
||||
import shlex
|
||||
import subprocess
|
||||
import time
|
||||
from datetime import datetime
|
||||
from typing import Any, Dict, Optional
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# Regex to strip Zulip mention artifacts from message bodies (ADR-008)
|
||||
MENTION_CLEANER = re.compile(r'@\*\*[^*]+\*\*')
|
||||
|
||||
|
||||
class ZulipAdapter:
|
||||
"""Zulip adapter for Hermes agents.
|
||||
|
||||
Connects to Zulip via event queues (not the deprecated call_on_each_message),
|
||||
polls for new events, routes @mentions and DMs to the Hermes agent via
|
||||
subprocess, and posts responses back to Zulip.
|
||||
"""
|
||||
|
||||
def __init__(self, config: Dict[str, Any]) -> None:
|
||||
self.config = config
|
||||
self.connected = False
|
||||
self._client = None
|
||||
self._queue_id: Optional[str] = None
|
||||
self._last_event_id: int = -1
|
||||
self._poll_task: Optional[asyncio.Task] = None
|
||||
|
||||
# Agent subprocess config
|
||||
self._agent_command: str = config.get("agent", {}).get(
|
||||
"command", "hermes chat"
|
||||
)
|
||||
self._agent_timeout: int = config.get("error_handling", {}).get(
|
||||
"timeout_seconds", 60
|
||||
)
|
||||
|
||||
# Bot identity for filtering own messages
|
||||
self._bot_email: str = config.get("zulip", {}).get("email", "")
|
||||
self._bot_id: int = config.get("swarm", {}).get("bot_id", 0)
|
||||
|
||||
# Stream config
|
||||
self._stream: str = config.get("zulip", {}).get("stream", "agent-hub")
|
||||
self._all_bots_user_id: int = config.get("zulip", {}).get(
|
||||
"all_bots_user_id", 1
|
||||
)
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# Connection lifecycle
|
||||
# ------------------------------------------------------------------
|
||||
|
||||
async def connect(self) -> None:
|
||||
"""Connect to Zulip and register an event queue."""
|
||||
if self.connected:
|
||||
logger.info("Already connected to Zulip.")
|
||||
return
|
||||
|
||||
try:
|
||||
import zulip
|
||||
|
||||
server_url = self.config["zulip"]["server_url"]
|
||||
email = self.config["zulip"]["email"]
|
||||
api_key = self.config["zulip"]["api_key"]
|
||||
|
||||
self._client = zulip.Client(
|
||||
server_url=server_url, email=email, api_key=api_key
|
||||
)
|
||||
|
||||
# Register an event queue (like the pi extension)
|
||||
queue_res = self._client.register(
|
||||
event_types=["message"]
|
||||
)
|
||||
self._queue_id = queue_res.get("queue_id")
|
||||
self._last_event_id = queue_res.get("last_event_id", -1)
|
||||
self.connected = True
|
||||
|
||||
logger.info(
|
||||
f"Connected to Zulip: {server_url} as {email} "
|
||||
f"(queue={self._queue_id})"
|
||||
)
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to connect to Zulip: {e}")
|
||||
self.connected = False
|
||||
raise
|
||||
|
||||
async def disconnect(self) -> None:
|
||||
"""Disconnect and cancel the poll loop."""
|
||||
if self._poll_task:
|
||||
self._poll_task.cancel()
|
||||
try:
|
||||
await self._poll_task
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
self._poll_task = None
|
||||
self.connected = False
|
||||
logger.info("Disconnected from Zulip.")
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# Event polling (async loop, like pi extension's setInterval)
|
||||
# ------------------------------------------------------------------
|
||||
|
||||
async def poll_forever(self, poll_interval: float = 3.0) -> None:
|
||||
"""Continuously poll the Zulip event queue and process messages.
|
||||
|
||||
This is the main event loop. Runs until cancelled.
|
||||
"""
|
||||
while self.connected and self._client and self._queue_id:
|
||||
try:
|
||||
events = await self._poll_events()
|
||||
for event in events:
|
||||
await self._process_event(event)
|
||||
except asyncio.CancelledError:
|
||||
break
|
||||
except Exception as e:
|
||||
logger.error(f"Poll error: {e}")
|
||||
# Check for queue expiry
|
||||
if "BAD_EVENT_QUEUE_ID" in str(e):
|
||||
logger.info("Queue expired, reconnecting...")
|
||||
self.connected = False
|
||||
await self._reconnect_with_retry()
|
||||
continue
|
||||
|
||||
await asyncio.sleep(poll_interval)
|
||||
|
||||
async def _poll_events(self) -> list:
|
||||
"""Fetch events from the Zulip event queue."""
|
||||
import zulip
|
||||
|
||||
if not self._client or not self._queue_id:
|
||||
return []
|
||||
|
||||
response = self._client.get_events(
|
||||
queue_id=self._queue_id,
|
||||
last_event_id=self._last_event_id,
|
||||
dont_block=True,
|
||||
)
|
||||
|
||||
if response.get("result") != "success":
|
||||
raise RuntimeError(
|
||||
f"Events API error: {response.get('msg', 'unknown')}"
|
||||
)
|
||||
|
||||
events = response.get("events", [])
|
||||
for event in events:
|
||||
if event.get("id", 0) > self._last_event_id:
|
||||
self._last_event_id = event["id"]
|
||||
|
||||
return [e for e in events if e.get("type") == "message"]
|
||||
|
||||
async def _reconnect_with_retry(self) -> None:
|
||||
"""Retry connection with backoff."""
|
||||
retries = self.config.get("error_handling", {}).get("retry_count", 3)
|
||||
delay = self.config.get("error_handling", {}).get(
|
||||
"retry_delay_seconds", 5
|
||||
)
|
||||
|
||||
for attempt in range(retries):
|
||||
logger.info(
|
||||
f"Reconnect attempt {attempt + 1}/{retries} "
|
||||
f"in {delay}s..."
|
||||
)
|
||||
await asyncio.sleep(delay)
|
||||
try:
|
||||
await self.connect()
|
||||
if self.connected:
|
||||
logger.info("Reconnected successfully.")
|
||||
return
|
||||
except Exception as e:
|
||||
logger.error(f"Reconnect attempt {attempt + 1} failed: {e}")
|
||||
delay *= 2 # Exponential backoff
|
||||
|
||||
logger.error("Max reconnection retries reached. Giving up.")
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# Message processing
|
||||
# ------------------------------------------------------------------
|
||||
|
||||
async def _process_event(self, event: Dict[str, Any]) -> None:
|
||||
"""Route a Zulip message event to the Hermes agent."""
|
||||
msg = event.get("message", {})
|
||||
msg_type = msg.get("type", "") # "private" or "stream"
|
||||
sender_email = msg.get("sender_email", "")
|
||||
sender_name = msg.get("sender_full_name", "Unknown")
|
||||
content = msg.get("content", "")
|
||||
|
||||
# Ignore own messages
|
||||
if sender_email == self._bot_email:
|
||||
return
|
||||
|
||||
# Determine if this message targets this bot
|
||||
mentioned_users = msg.get("mentioned_users", [])
|
||||
mentioned_user_ids = [u.get("user_id") for u in mentioned_users]
|
||||
is_dm = msg_type == "private"
|
||||
is_mention = self._bot_id in mentioned_user_ids
|
||||
is_all_bots = self._all_bots_user_id in mentioned_user_ids
|
||||
|
||||
# DM-first: process all private messages (like ADR-001/ADR-002)
|
||||
if is_dm:
|
||||
logger.info(f"DM from {sender_name}: {content[:80]}...")
|
||||
await self._route_to_agent(
|
||||
message_type="dm",
|
||||
sender_name=sender_name,
|
||||
content=content,
|
||||
recipient=sender_email,
|
||||
)
|
||||
return
|
||||
|
||||
# Stream: only respond to @mentions and @all-bots (ADR-005, ADR-006)
|
||||
if msg_type == "stream" and (is_mention or is_all_bots):
|
||||
stream_name = msg.get("display_recipient", "unknown")
|
||||
topic = msg.get("subject", "general")
|
||||
logger.info(
|
||||
f"{'@mention' if is_mention else '@all-bots'} "
|
||||
f"in #{stream_name} > {topic}"
|
||||
)
|
||||
|
||||
# Clean @mention artifacts from content (ADR-008)
|
||||
clean_content = MENTION_CLEANER.sub("", content).strip()
|
||||
|
||||
await self._route_to_agent(
|
||||
message_type="mention" if is_mention else "all-bots",
|
||||
sender_name=sender_name,
|
||||
content=clean_content,
|
||||
recipient=stream_name,
|
||||
topic=topic,
|
||||
stream_id=msg.get("stream_id"),
|
||||
)
|
||||
|
||||
async def _route_to_agent(
|
||||
self,
|
||||
message_type: str,
|
||||
sender_name: str,
|
||||
content: str,
|
||||
recipient: str,
|
||||
topic: Optional[str] = None,
|
||||
stream_id: Optional[int] = None,
|
||||
) -> None:
|
||||
"""Send the message to the Hermes agent subprocess and relay the
|
||||
response back to Zulip.
|
||||
|
||||
Architecture note: This uses a subprocess (configurable via
|
||||
agent.command in config.yaml). This mirrors the initial pi extension
|
||||
approach. A future iteration could use IPC or a socket if the
|
||||
Hermes agent exposes one.
|
||||
"""
|
||||
# 1. Send typing indicator
|
||||
await self._send_typing_indicator(recipient, "start")
|
||||
|
||||
# 2. Send a "Thinking..." placeholder to Zulip immediately.
|
||||
# On agent_end, this message will be edited with the final response.
|
||||
placeholder = (
|
||||
f":robot: _Processing your message..._"
|
||||
)
|
||||
placeholder_msg_id = None
|
||||
try:
|
||||
placeholder_msg_id = await self._send_message(
|
||||
message_type=message_type,
|
||||
recipient=recipient,
|
||||
content=placeholder,
|
||||
topic=topic,
|
||||
stream_id=stream_id,
|
||||
)
|
||||
except Exception as e:
|
||||
logger.warning(f"Failed to send placeholder: {e}")
|
||||
|
||||
# 3. Invoke the Hermes agent subprocess
|
||||
response_text = await self._invoke_agent(content, sender_name)
|
||||
|
||||
# 4. Stop typing indicator
|
||||
await self._send_typing_indicator(recipient, "stop")
|
||||
|
||||
# 5. Post (or edit) the response
|
||||
if not response_text.strip():
|
||||
logger.warning(
|
||||
f"Hermes agent returned empty response for "
|
||||
f"{message_type} from {sender_name}"
|
||||
)
|
||||
return
|
||||
|
||||
# Truncate to Zulip's 10K char limit
|
||||
MAX_ZULIP_MSG = 10000
|
||||
truncated = (
|
||||
response_text[:MAX_ZULIP_MSG]
|
||||
+ "\n\n[...truncated at Zulip limit]"
|
||||
if len(response_text) > MAX_ZULIP_MSG
|
||||
else response_text
|
||||
)
|
||||
|
||||
try:
|
||||
if placeholder_msg_id:
|
||||
await self._edit_message(placeholder_msg_id, truncated)
|
||||
logger.info(
|
||||
f"Finalized response to {sender_name} "
|
||||
f"({len(truncated)} chars)"
|
||||
)
|
||||
else:
|
||||
await self._send_message(
|
||||
message_type=message_type,
|
||||
recipient=recipient,
|
||||
content=truncated,
|
||||
topic=topic,
|
||||
stream_id=stream_id,
|
||||
)
|
||||
logger.info(
|
||||
f"Sent response to {sender_name} "
|
||||
f"({len(truncated)} chars)"
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to post response: {e}")
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# Hermes agent subprocess invocation
|
||||
# ------------------------------------------------------------------
|
||||
|
||||
async def _invoke_agent(
|
||||
self, message: str, sender_name: str
|
||||
) -> str:
|
||||
"""Spawn the Hermes agent subprocess with the message.
|
||||
|
||||
The agent command is configurable (agent.command in config.yaml).
|
||||
Default: "hermes chat"
|
||||
|
||||
The message is passed via stdin (pipe). The agent's stdout is
|
||||
captured as the response.
|
||||
|
||||
This is synchronous (run in executor) to avoid blocking the
|
||||
event loop during subprocess execution.
|
||||
"""
|
||||
cmd_str = self._agent_command
|
||||
cmd = shlex.split(cmd_str)
|
||||
|
||||
logger.info(
|
||||
f"Invoking agent: {' '.join(cmd)} "
|
||||
f"(timeout={self._agent_timeout}s)"
|
||||
)
|
||||
|
||||
def _run() -> str:
|
||||
try:
|
||||
result = subprocess.run(
|
||||
cmd,
|
||||
input=message,
|
||||
capture_output=True,
|
||||
text=True,
|
||||
timeout=self._agent_timeout,
|
||||
)
|
||||
if result.returncode != 0:
|
||||
stderr = result.stderr.strip()
|
||||
logger.error(
|
||||
f"Agent exited with code {result.returncode}: "
|
||||
f"{stderr}"
|
||||
)
|
||||
return (
|
||||
f":warning: Agent error (exit {result.returncode}). "
|
||||
f"Please try again later."
|
||||
)
|
||||
return result.stdout.strip()
|
||||
except subprocess.TimeoutExpired:
|
||||
logger.error(
|
||||
f"Agent timed out after {self._agent_timeout}s"
|
||||
)
|
||||
return (
|
||||
f":hourglass: Agent timed out after "
|
||||
f"{self._agent_timeout}s. Please try again."
|
||||
)
|
||||
except FileNotFoundError:
|
||||
logger.error(f"Agent command not found: {cmd_str}")
|
||||
return (
|
||||
f":warning: Agent command not found: `{cmd_str}`. "
|
||||
f"Check configuration."
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(f"Agent invocation error: {e}")
|
||||
return (
|
||||
f":warning: Failed to invoke agent: {e}"
|
||||
)
|
||||
|
||||
return await asyncio.get_event_loop().run_in_executor(None, _run)
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# Zulip API helpers
|
||||
# ------------------------------------------------------------------
|
||||
|
||||
async def _send_message(
|
||||
self,
|
||||
message_type: str,
|
||||
recipient: str,
|
||||
content: str,
|
||||
topic: Optional[str] = None,
|
||||
stream_id: Optional[int] = None,
|
||||
) -> Optional[int]:
|
||||
"""Send a message to Zulip. Returns the message ID if successful."""
|
||||
if not self._client:
|
||||
return None
|
||||
|
||||
def _send() -> Optional[int]:
|
||||
if message_type in ("dm", "private"):
|
||||
# Private message: recipient is email
|
||||
payload = {
|
||||
"type": "private",
|
||||
"to": [recipient],
|
||||
"content": content,
|
||||
}
|
||||
else:
|
||||
# Stream message
|
||||
payload = {
|
||||
"type": "stream",
|
||||
"to": stream_id or recipient,
|
||||
"subject": topic or "general",
|
||||
"content": content,
|
||||
}
|
||||
|
||||
try:
|
||||
result = self._client.send_message(payload)
|
||||
if result.get("result") == "success":
|
||||
return result.get("id")
|
||||
logger.error(
|
||||
f"Send message failed: {result.get('msg', 'unknown')}"
|
||||
)
|
||||
return None
|
||||
except Exception as e:
|
||||
logger.error(f"Send message error: {e}")
|
||||
return None
|
||||
|
||||
return await asyncio.get_event_loop().run_in_executor(None, _send)
|
||||
|
||||
async def _edit_message(
|
||||
self, message_id: int, content: str
|
||||
) -> bool:
|
||||
"""Edit a previously sent Zulip message (for streaming updates)."""
|
||||
if not self._client:
|
||||
return False
|
||||
|
||||
def _edit() -> bool:
|
||||
try:
|
||||
result = self._client.update_message(
|
||||
{"message_id": message_id, "content": content}
|
||||
)
|
||||
if result.get("result") != "success":
|
||||
logger.warning(
|
||||
f"Edit message {message_id}: "
|
||||
f"{result.get('msg', 'unknown')}"
|
||||
)
|
||||
return result.get("result") == "success"
|
||||
except Exception as e:
|
||||
logger.warning(f"Edit message {message_id} error: {e}")
|
||||
return False
|
||||
|
||||
return await asyncio.get_event_loop().run_in_executor(None, _edit)
|
||||
|
||||
async def _send_typing_indicator(
|
||||
self, recipient: str, operation: str
|
||||
) -> None:
|
||||
"""Send a typing indicator (start/stop)."""
|
||||
if not self._client:
|
||||
return
|
||||
|
||||
def _typing() -> None:
|
||||
try:
|
||||
# The zulip Python library doesn't have a direct typing API
|
||||
# Use the REST endpoint directly via requests
|
||||
import requests
|
||||
|
||||
server_url = self.config["zulip"]["server_url"]
|
||||
email = self.config["zulip"]["email"]
|
||||
api_key = self.config["zulip"]["api_key"]
|
||||
|
||||
# For private messages, recipient is the user's email
|
||||
# We need their user_id. Use the API directly.
|
||||
to_data = json.dumps([recipient])
|
||||
requests.post(
|
||||
f"{server_url}/api/v1/typing",
|
||||
auth=requests.auth.HTTPBasicAuth(email, api_key),
|
||||
data={"to": to_data, "op": operation},
|
||||
timeout=5,
|
||||
)
|
||||
except Exception as e:
|
||||
# Non-critical
|
||||
pass
|
||||
|
||||
await asyncio.get_event_loop().run_in_executor(None, _typing)
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# Synchronous entry point for the Hermes runner
|
||||
# ------------------------------------------------------------------
|
||||
|
||||
def run(self) -> None:
|
||||
"""Synchronous entry point that starts the async poll loop.
|
||||
|
||||
Called by the Hermes runner (e.g., from a systemd service).
|
||||
"""
|
||||
asyncio.run(self._run_async())
|
||||
|
||||
async def _run_async(self) -> None:
|
||||
"""Async entry point."""
|
||||
try:
|
||||
await self.connect()
|
||||
if self.connected:
|
||||
logger.info("Starting event poll loop...")
|
||||
await self.poll_forever()
|
||||
except asyncio.CancelledError:
|
||||
logger.info("Poll loop cancelled.")
|
||||
except Exception as e:
|
||||
logger.error(f"Fatal error in event loop: {e}")
|
||||
finally:
|
||||
await self.disconnect()
|
||||
@@ -1,71 +0,0 @@
|
||||
# Abiba Zulip Gateway — Standalone Service
|
||||
|
||||
**Replaces** the old pi extension (`~/.pi/agent/extensions/zulip/`).
|
||||
|
||||
Instead of injecting messages into pi's session (which dies when pi exits), this
|
||||
runs as an independent Node.js **systemd service** that:
|
||||
- Polls Zulip event queue for DMs
|
||||
- Calls the harness inference API directly (no pi dependency)
|
||||
- Maintains per-sender conversation memory
|
||||
- Survives pi shutdown and restart
|
||||
|
||||
## Architecture
|
||||
|
||||
```
|
||||
Zulip event queue → poll every 3s → DM detected
|
||||
→ send typing indicator + placeholder
|
||||
→ POST /v1/chat/completions with conversation history
|
||||
→ edit placeholder with final response
|
||||
```
|
||||
|
||||
## Requirements
|
||||
|
||||
- Node.js 22+
|
||||
- `npm install` in this directory
|
||||
|
||||
## Config
|
||||
|
||||
All config via environment variables. Copy `abiba-zulip.service` to set up as a
|
||||
systemd service, or run directly:
|
||||
|
||||
```bash
|
||||
ZULIP_EMAIL=abiba-bot@chat.sysloggh.net \
|
||||
ZULIP_API_KEY=your_key \
|
||||
ZULIP_SITE=https://chat.sysloggh.net \
|
||||
HARNESS_URL=http://192.168.68.116/v1/chat/completions \
|
||||
HARNESS_API_KEY=sk-xxx \
|
||||
HARNESS_MODEL=qwen3.6-35B-A3B \
|
||||
node dist/index.js
|
||||
```
|
||||
|
||||
## Deploy as a service
|
||||
|
||||
```bash
|
||||
# 1. Install deps
|
||||
cd /opt/abiba-zulip
|
||||
npm install
|
||||
|
||||
# 2. Build
|
||||
npx tsc
|
||||
|
||||
# 3. Install systemd service
|
||||
cp abiba-zulip.service /etc/systemd/system/
|
||||
systemctl daemon-reload
|
||||
systemctl enable abiba-zulip
|
||||
systemctl start abiba-zulip
|
||||
|
||||
# 4. Check status
|
||||
systemctl status abiba-zulip
|
||||
journalctl -u abiba-zulip -f
|
||||
|
||||
# 5. Health check
|
||||
curl http://127.0.0.1:9200/health
|
||||
```
|
||||
|
||||
## Development
|
||||
|
||||
```bash
|
||||
npm run dev # watch mode (tsc watch)
|
||||
npm run build # compile
|
||||
npm start # run compiled
|
||||
```
|
||||
@@ -1,48 +0,0 @@
|
||||
[Unit]
|
||||
Description=Abiba Zulip Gateway Service — Standalone Zulip bot for pi agent
|
||||
Documentation=https://git.sysloggh.net/SyslogSolution/zulip-platform-plugins
|
||||
After=network-online.target
|
||||
Wants=network-online.target
|
||||
|
||||
[Service]
|
||||
Type=simple
|
||||
User=root
|
||||
WorkingDirectory=/opt/abiba-zulip
|
||||
|
||||
# ── Zulip credentials ─────────────────────────────────────────────────────
|
||||
Environment=ZULIP_EMAIL=abiba-bot@chat.sysloggh.net
|
||||
Environment=ZULIP_API_KEY=cKTDMZAPW08dk3zl05sStzO7HRztzyn8
|
||||
Environment=ZULIP_SITE=https://chat.sysloggh.net
|
||||
|
||||
# ── Agent identity ─────────────────────────────────────────────────────────
|
||||
Environment=AGENT_NAME=abiba
|
||||
Environment=AGENT_OWNER_EMAIL=jerome@sysloggh.com
|
||||
|
||||
# ── Harness API (inference) ────────────────────────────────────────────────
|
||||
Environment=HARNESS_URL=http://192.168.68.116/v1/chat/completions
|
||||
Environment=HARNESS_API_KEY=sk-856ffb0bbb-e5aaf78b10054eca608f8fbcbd73a889
|
||||
Environment=HARNESS_MODEL=qwen3.6-35B-A3B
|
||||
Environment=HARNESS_MAX_TOKENS=4096
|
||||
|
||||
# ── Service config ─────────────────────────────────────────────────────────
|
||||
Environment=HEALTH_PORT=9200
|
||||
Environment=POLL_INTERVAL_MS=3000
|
||||
Environment=MAX_RETRIES=5
|
||||
Environment=RETRY_DELAY_MS=5000
|
||||
|
||||
ExecStart=/usr/bin/node /opt/abiba-zulip/dist/index.js
|
||||
Restart=always
|
||||
RestartSec=5
|
||||
|
||||
# Logging
|
||||
StandardOutput=journal
|
||||
StandardError=journal
|
||||
|
||||
# Security hardening (optional — adjust as needed)
|
||||
NoNewPrivileges=true
|
||||
ProtectHome=read-only
|
||||
ProtectSystem=full
|
||||
PrivateTmp=true
|
||||
|
||||
[Install]
|
||||
WantedBy=multi-user.target
|
||||
Generated
+1975
-17
File diff suppressed because it is too large
Load Diff
@@ -1,19 +1,28 @@
|
||||
{
|
||||
"name": "abiba-zulip-service",
|
||||
"version": "2.0.0",
|
||||
"description": "Standalone Zulip gateway service for Abiba (pi agent). Direct harness API, conversation memory, systemd service.",
|
||||
"name": "pi-zulip-extension",
|
||||
"version": "0.1.0",
|
||||
"description": "pi extension for Zulip agent communication — connects Abiba to the Sysloggh agent mesh",
|
||||
"main": "src/index.ts",
|
||||
"type": "module",
|
||||
"main": "dist/index.js",
|
||||
"scripts": {
|
||||
"build": "tsc",
|
||||
"start": "node dist/index.js",
|
||||
"dev": "node --watch dist/index.js"
|
||||
"check": "tsc --noEmit"
|
||||
},
|
||||
"dependencies": {
|
||||
"yaml": "^2.9.0",
|
||||
"zulip-js": "^2.0.0"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/node": "^22.0.0",
|
||||
"typescript": "^5.7.0"
|
||||
}
|
||||
"@earendil-works/pi-coding-agent": "^0.80.2",
|
||||
"typescript": "^5.0.0"
|
||||
},
|
||||
"keywords": [
|
||||
"pi",
|
||||
"zulip",
|
||||
"agent",
|
||||
"abiba",
|
||||
"sysloggh"
|
||||
],
|
||||
"license": "UNLICENSED",
|
||||
"private": true
|
||||
}
|
||||
|
||||
+521
-387
File diff suppressed because it is too large
Load Diff
@@ -1,18 +1,14 @@
|
||||
{
|
||||
"compilerOptions": {
|
||||
"target": "ES2022",
|
||||
"module": "ES2022",
|
||||
"moduleResolution": "node",
|
||||
"outDir": "./dist",
|
||||
"rootDir": "./src",
|
||||
"strict": false,
|
||||
"module": "ESNext",
|
||||
"moduleResolution": "bundler",
|
||||
"strict": true,
|
||||
"esModuleInterop": true,
|
||||
"skipLibCheck": true,
|
||||
"forceConsistentCasingInFileNames": true,
|
||||
"outDir": "dist",
|
||||
"rootDir": "src",
|
||||
"declaration": true,
|
||||
"declarationMap": true,
|
||||
"sourceMap": true
|
||||
"skipLibCheck": true
|
||||
},
|
||||
"include": ["src/**/*"],
|
||||
"exclude": ["node_modules", "dist"]
|
||||
"include": ["src/**/*.ts"]
|
||||
}
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,56 @@
|
||||
name: zulip-platform
|
||||
label: Zulip
|
||||
kind: platform
|
||||
version: 1.0.0
|
||||
description: >
|
||||
Zulip messaging platform adapter for Hermes Agent. Connects to a Zulip
|
||||
server via event queue polling, processes DMs and @mentions, and sends
|
||||
replies with placeholder→edit streaming. Uses the official zulip-js
|
||||
HTTP API pattern (polling, not WebSockets) — lightweight, no external
|
||||
SDK beyond httpx.
|
||||
|
||||
author: Syslog Solution LLC
|
||||
|
||||
requires_env:
|
||||
- name: ZULIP_SITE
|
||||
description: "Zulip server URL (e.g. https://chat.sysloggh.net)"
|
||||
prompt: "Zulip server URL"
|
||||
password: false
|
||||
- name: ZULIP_EMAIL
|
||||
description: "Bot email address (e.g. tanko-bot@chat.sysloggh.net)"
|
||||
prompt: "Zulip bot email"
|
||||
password: false
|
||||
- name: ZULIP_API_KEY
|
||||
description: "Zulip bot API key"
|
||||
prompt: "Zulip API key"
|
||||
password: true
|
||||
|
||||
optional_env:
|
||||
- name: ZULIP_STREAM
|
||||
description: "Primary stream to subscribe to (default: agent-hub)"
|
||||
prompt: "Zulip stream name"
|
||||
password: false
|
||||
- name: ZULIP_ALL_BOTS_USER_ID
|
||||
description: "User ID of the @all-bots user (default: 1)"
|
||||
prompt: "All Bots user ID"
|
||||
password: false
|
||||
- name: ZULIP_AGENT_NAME
|
||||
description: "Agent display name for logging (default: hermes-agent)"
|
||||
prompt: "Agent name"
|
||||
password: false
|
||||
- name: ZULIP_OWNER_EMAIL
|
||||
description: "Owner email for private topic ACL"
|
||||
prompt: "Owner email"
|
||||
password: false
|
||||
- name: ZULIP_POLL_INTERVAL
|
||||
description: "Event poll interval in seconds (default: 3)"
|
||||
prompt: "Poll interval (seconds)"
|
||||
password: false
|
||||
- name: ZULIP_HOME_CHANNEL
|
||||
description: "Default recipient for cron / notification delivery"
|
||||
prompt: "Home channel (email or stream:topic)"
|
||||
password: false
|
||||
- name: ZULIP_HOME_CHANNEL_NAME
|
||||
description: "Human label for the home channel"
|
||||
prompt: "Home channel display name"
|
||||
password: false
|
||||
+183
-81
@@ -1,32 +1,56 @@
|
||||
#!/usr/bin/env bash
|
||||
# deploy.sh GitOps deployment for zulip-platform-plugins
|
||||
# Usage: ./deploy.sh <tag|branch> [--ct=<agent>] [--dry-run]
|
||||
# ./deploy.sh v1.0.0 # Deploy to all 6 CTs
|
||||
# ./deploy.sh main # Deploy latest (staging only!)
|
||||
# ./deploy.sh --ct=tanko v1.0.0 # Deploy to single CT
|
||||
# ./deploy.sh --dry-run v1.0.0 # Preview deployment without making changes
|
||||
# deploy.sh — GitOps deployment for zulip-platform-plugins
|
||||
#
|
||||
# Supports two deployment modes:
|
||||
# LEGACY: /opt/hermes-zulip-plugin/ + systemctl restart zulip-plugin
|
||||
# NATIVE: ~/.hermes/plugins/platforms/zulip/ + hermes gateway restart
|
||||
#
|
||||
# Usage:
|
||||
# ./deploy.sh <tag|branch> # Deploy all agents (default: native)
|
||||
# ./deploy.sh --mode=native v1.0.0 # Hermes native plugin (new)
|
||||
# ./deploy.sh --mode=legacy v1.0.0 # Old systemd service (deprecated)
|
||||
# ./deploy.sh --ct=tanko v1.0.0 # Single agent
|
||||
# ./deploy.sh --dry-run v1.0.0 # Preview only
|
||||
#
|
||||
set -euo pipefail
|
||||
|
||||
DEPLOY_LOG="deploy.log"
|
||||
TAG=""
|
||||
SINGLE_CT=""
|
||||
DRY_RUN=false
|
||||
DEPLOY_MODE="native" # default to new Hermes native plugin
|
||||
|
||||
# Parse args
|
||||
for arg in "$@"; do
|
||||
case "$arg" in
|
||||
--ct=*) SINGLE_CT="${arg#--ct=}" ;;
|
||||
--dry-run) DRY_RUN=true ;;
|
||||
--mode=*) DEPLOY_MODE="${arg#--mode=}" ;;
|
||||
*) TAG="$arg" ;;
|
||||
esac
|
||||
done
|
||||
|
||||
if [[ -z "$TAG" ]]; then
|
||||
echo "Usage: $0 <tag|branch> [--ct=<agent>] [--dry-run]"
|
||||
echo "Usage: $0 <tag|branch> [--ct=<agent>] [--mode=native|legacy] [--dry-run]"
|
||||
echo ""
|
||||
echo " --ct=<agent> Deploy to single agent (tanko, mumuni, etc.)"
|
||||
echo " --mode=native Hermes native plugin at ~/.hermes/plugins/ (default)"
|
||||
echo " --mode=legacy Old systemd service at /opt/hermes-zulip-plugin/"
|
||||
echo " --dry-run Preview deployment without making changes"
|
||||
echo ""
|
||||
echo "Examples:"
|
||||
echo " ./deploy.sh v1.0.0 # Deploy native to all agents"
|
||||
echo " ./deploy.sh --ct=tanko v1.0.0 # Deploy native to Tanko only"
|
||||
echo " ./deploy.sh --mode=legacy v0.9.0 # Deploy legacy to all"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
# Agent CT registry (must match CONTEXT.md)
|
||||
if [[ "$DEPLOY_MODE" != "native" && "$DEPLOY_MODE" != "legacy" ]]; then
|
||||
echo "ERROR: --mode must be 'native' or 'legacy'"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
# ── Agent Registry (must match docs/CONTEXT.md) ──────────────────────
|
||||
declare -A AGENTS=(
|
||||
["tanko"]="amdpve CT 112"
|
||||
["mumuni"]="minipve CT 114"
|
||||
@@ -36,32 +60,31 @@ declare -A AGENTS=(
|
||||
["abiba"]="amdpve CT 100"
|
||||
)
|
||||
|
||||
# Platform paths per agent type
|
||||
declare -A PLUGIN_PATHS=(
|
||||
# Native Hermes plugin path (~/.hermes/plugins/platforms/zulip/)
|
||||
NATIVE_PLUGIN_SRC="plugins/platforms/zulip"
|
||||
|
||||
# Legacy paths (DEPRECATED — systemd zulip-plugin service)
|
||||
declare -A LEGACY_PATHS=(
|
||||
["tanko"]="/opt/hermes-zulip-plugin"
|
||||
["mumuni"]="/opt/hermes-zulip-plugin"
|
||||
["koonimo"]="/opt/hermes-zulip-plugin"
|
||||
["koby\"]="/opt/hermes-zulip-plugin"
|
||||
["kagentz\"]="/opt/agent-zero-plugin"
|
||||
["abiba\"]="/root/.pi/agent/extensions/zulip.ts"
|
||||
["koby"]="/opt/hermes-zulip-plugin"
|
||||
["kagentz"]="/opt/agent-zero-plugin"
|
||||
["abiba"]="/root/.pi/agent/extensions/zulip.ts"
|
||||
)
|
||||
|
||||
# Correcting the backslash artifacts from the original file
|
||||
PLUGIN_PATHS["koby"]="/opt/hermes-zulip-plugin"
|
||||
PLUGIN_PATHS["kagentz"]="/opt/agent-zero-plugin"
|
||||
PLUGIN_PATHS["abiba"]="/root/.pi/agent/extensions/zulip.ts"
|
||||
|
||||
declare -A SERVICE_NAMES=(
|
||||
# Service names for legacy mode
|
||||
declare -A LEGACY_SERVICES=(
|
||||
["tanko"]="zulip-plugin"
|
||||
["mumuni"]="zulip-plugin"
|
||||
["koonimo"]="zulip-plugin"
|
||||
["koby"]="zulip-plugin"
|
||||
["kagentz"]="zulip-plugin"
|
||||
["abiba"]="pi" # pi reload, not systemctl
|
||||
["abiba"]="pi"
|
||||
)
|
||||
|
||||
GITEA_REPO="https://git.sysloggh.net/SyslogSolution/zulip-platform-plugins.git"
|
||||
HEALTH_PORT=9200
|
||||
GITEA_REPO="https://git.sysloggh.net/SyslogSolution/zulip-platform-plugins.git"
|
||||
|
||||
log() {
|
||||
local prefix=""
|
||||
@@ -69,90 +92,169 @@ log() {
|
||||
echo "${prefix}[$(date '+%Y-%m-%d %H:%M:%S')] $*" | tee -a "$DEPLOY_LOG"
|
||||
}
|
||||
|
||||
deploy_agent() {
|
||||
# ── Deployment Functions ─────────────────────────────────────────────
|
||||
|
||||
deploy_native() {
|
||||
local agent="$1"
|
||||
local ct_info="${AGENTS[$agent]}"
|
||||
local plugin_path="${PLUGIN_PATHS[$agent]}"
|
||||
local service="${SERVICE_NAMES[$agent]}"
|
||||
local plugin_dir="$HOME/.hermes/plugins/platforms/zulip"
|
||||
|
||||
log "=== Deploying $agent ($ct_info) @ $TAG ==="
|
||||
log "Path: $plugin_path"
|
||||
log "📦 Deploying $agent: native Hermes plugin -> $plugin_dir"
|
||||
|
||||
# 1. Git pull & checkout tag
|
||||
if [[ "$DRY_RUN" != "true" ]]; then
|
||||
cd "$plugin_path" || { log "ERROR: $agent path $plugin_path not found"; return 1; }
|
||||
git fetch --tags origin
|
||||
git checkout "$TAG"
|
||||
else
|
||||
log "Skipping Git checkout (Dry Run)"
|
||||
fi
|
||||
log "$agent: checked out $TAG"
|
||||
|
||||
# 2. Install dependencies (platform-specific)
|
||||
if [[ "$DRY_RUN" != "true" ]]; then
|
||||
case "$agent" in
|
||||
tanko|mumuni|koonimo|koby)
|
||||
pip install -r requirements.txt --quiet
|
||||
;;
|
||||
kagentz)
|
||||
pip install -r requirements.txt --quiet
|
||||
;;
|
||||
abiba)
|
||||
log "$agent: pi extension skipping pip install"
|
||||
;;
|
||||
esac
|
||||
else
|
||||
log "Skipping dependency installation (Dry Run)"
|
||||
if [[ "$DRY_RUN" == "true" ]]; then
|
||||
log " Would copy $NATIVE_PLUGIN_SRC/{adapter.py,__init__.py,plugin.yaml} -> $plugin_dir/"
|
||||
log " Would run: hermes gateway restart"
|
||||
return 0
|
||||
fi
|
||||
|
||||
# 3. Restart service
|
||||
if [[ "$DRY_RUN" != "true" ]]; then
|
||||
case "$agent" in
|
||||
abiba)
|
||||
log "$agent: triggering /reload"
|
||||
;;
|
||||
*)
|
||||
systemctl restart "$service"
|
||||
log "$agent: restarted $service"
|
||||
;;
|
||||
esac
|
||||
# Create plugin directory
|
||||
mkdir -p "$plugin_dir"
|
||||
|
||||
# Copy plugin files
|
||||
cp "$NATIVE_PLUGIN_SRC/adapter.py" "$plugin_dir/"
|
||||
cp "$NATIVE_PLUGIN_SRC/__init__.py" "$plugin_dir/"
|
||||
cp "$NATIVE_PLUGIN_SRC/plugin.yaml" "$plugin_dir/"
|
||||
|
||||
log " Copied plugin files to $plugin_dir/"
|
||||
ls -la "$plugin_dir/"
|
||||
|
||||
# Restart Hermes Gateway to load the plugin
|
||||
if command -v hermes &>/dev/null; then
|
||||
log " Restarting Hermes Gateway..."
|
||||
hermes gateway restart
|
||||
log " Hermes Gateway restarted"
|
||||
else
|
||||
log "Skipping service restart (Dry Run)"
|
||||
log " ⚠️ 'hermes' command not found — manual restart needed"
|
||||
log " Run: hermes gateway restart"
|
||||
fi
|
||||
|
||||
# 4. Health check
|
||||
log "Waiting for service to stabilize..."
|
||||
# Health check
|
||||
log " Waiting for service to stabilize..."
|
||||
sleep 5
|
||||
if [[ "$DRY_RUN" != "true" ]]; then
|
||||
if curl -sf "http://localhost:$HEALTH_PORT/health" > /dev/null 2>&1; then
|
||||
log "OK: $agent health check passed"
|
||||
log " ✅ Health check passed (port $HEALTH_PORT)"
|
||||
else
|
||||
log "ERROR: $agent health check failed check logs"
|
||||
return 1
|
||||
fi
|
||||
else
|
||||
log "Skipping health check (Dry Run)"
|
||||
log " ⚠️ Health check on port $HEALTH_PORT not responding"
|
||||
log " Check Gateway logs for plugin load errors"
|
||||
fi
|
||||
|
||||
log "✅ $agent: native deployment complete"
|
||||
}
|
||||
|
||||
# --- Main ---\
|
||||
log "Deploy started target: $TAG (Dry Run: $DRY_RUN)"
|
||||
deploy_legacy() {
|
||||
local agent="$1"
|
||||
local plugin_path="${LEGACY_PATHS[$agent]}"
|
||||
local service="${LEGACY_SERVICES[$agent]}"
|
||||
|
||||
log "📦 Deploying $agent: LEGACY mode -> $plugin_path"
|
||||
|
||||
if [[ "$DRY_RUN" == "true" ]]; then
|
||||
log " Would git checkout $TAG in $plugin_path"
|
||||
log " Would install deps + restart $service"
|
||||
return 0
|
||||
fi
|
||||
|
||||
# Git checkout
|
||||
if [[ ! -d "$plugin_path" ]]; then
|
||||
log "❌ Path $plugin_path not found for $agent"
|
||||
return 1
|
||||
fi
|
||||
|
||||
cd "$plugin_path"
|
||||
git fetch --tags origin
|
||||
git checkout "$TAG"
|
||||
|
||||
log " Checked out $TAG in $plugin_path"
|
||||
|
||||
# Install deps
|
||||
if [[ -f "requirements.txt" ]]; then
|
||||
pip install -r requirements.txt --quiet
|
||||
log " Dependencies installed"
|
||||
fi
|
||||
|
||||
# Restart service
|
||||
if systemctl list-units --full -all 2>/dev/null | grep -q "$service"; then
|
||||
systemctl restart "$service"
|
||||
log " Restarted $service"
|
||||
else
|
||||
log " ⚠️ Service $service not found — manual restart needed"
|
||||
fi
|
||||
|
||||
# Health check
|
||||
sleep 5
|
||||
if curl -sf "http://localhost:$HEALTH_PORT/health" > /dev/null 2>&1; then
|
||||
log " ✅ Health check passed"
|
||||
else
|
||||
log " ⚠️ Health check failed — check logs"
|
||||
fi
|
||||
|
||||
log "✅ $agent: legacy deployment complete"
|
||||
}
|
||||
|
||||
# ── Verify function ──────────────────────────────────────────────────
|
||||
|
||||
verify_deployment() {
|
||||
local agent="$1"
|
||||
|
||||
log "🔍 Verifying $agent deployment..."
|
||||
|
||||
if [[ "$DEPLOY_MODE" == "native" ]]; then
|
||||
local plugin_dir="$HOME/.hermes/plugins/platforms/zulip"
|
||||
if [[ -f "$plugin_dir/adapter.py" && -f "$plugin_dir/plugin.yaml" ]]; then
|
||||
log " ✅ Plugin files present in $plugin_dir"
|
||||
log " adapter.py: $(wc -l < "$plugin_dir/adapter.py") lines"
|
||||
log " plugin.yaml: $(wc -l < "$plugin_dir/plugin.yaml") lines"
|
||||
else
|
||||
log " ❌ Plugin files missing in $plugin_dir"
|
||||
return 1
|
||||
fi
|
||||
fi
|
||||
|
||||
log " ✅ $agent verification complete"
|
||||
}
|
||||
|
||||
# ── Main ─────────────────────────────────────────────────────────────
|
||||
|
||||
log "🚀 Deploy started — tag: $TAG, mode: $DEPLOY_MODE, dry-run: $DRY_RUN"
|
||||
|
||||
if [[ -n "$SINGLE_CT" ]]; then
|
||||
deploy_agent "$SINGLE_CT"
|
||||
if [[ -z "${AGENTS[$SINGLE_CT]:-}" ]]; then
|
||||
log "❌ Unknown agent: $SINGLE_CT"
|
||||
log " Known agents: ${!AGENTS[*]}"
|
||||
exit 1
|
||||
fi
|
||||
log "--- Deploying single agent: $SINGLE_CT ---"
|
||||
if [[ "$DEPLOY_MODE" == "native" ]]; then
|
||||
deploy_native "$SINGLE_CT"
|
||||
else
|
||||
deploy_legacy "$SINGLE_CT"
|
||||
fi
|
||||
verify_deployment "$SINGLE_CT"
|
||||
else
|
||||
FAILED=""
|
||||
for agent in tanko mumuni koonimo koby kagentz abiba; do
|
||||
if ! deploy_agent "$agent"; then
|
||||
for agent in "${!AGENTS[@]}"; do
|
||||
echo ""
|
||||
if [[ "$DEPLOY_MODE" == "native" ]]; then
|
||||
if deploy_native "$agent"; then
|
||||
verify_deployment "$agent" || FAILED="$FAILED $agent"
|
||||
else
|
||||
FAILED="$FAILED $agent"
|
||||
fi
|
||||
else
|
||||
if deploy_legacy "$agent"; then
|
||||
verify_deployment "$agent" || FAILED="$FAILED $agent"
|
||||
else
|
||||
FAILED="$FAILED $agent"
|
||||
fi
|
||||
fi
|
||||
done
|
||||
|
||||
echo ""
|
||||
log "=== Deploy complete ==="
|
||||
if [[ -n "$FAILED" ]]; then
|
||||
log "FAILED:$FAILED"
|
||||
log "Run rollback: ./scripts/rollback.sh <previous-tag>"
|
||||
log "❌ FAILED:$FAILED"
|
||||
log " Run rollback: ./scripts/rollback.sh <previous-tag>"
|
||||
exit 1
|
||||
fi
|
||||
log "All 6 agents deployed successfully."
|
||||
log "✅ All agents deployed successfully."
|
||||
log " Next: monitor #agent-hub for agent responses"
|
||||
fi
|
||||
|
||||
@@ -0,0 +1,92 @@
|
||||
#!/usr/bin/env bash
|
||||
# verify-deployment.sh — Check if the Hermes Zulip native plugin is properly deployed
|
||||
#
|
||||
# Usage:
|
||||
# ./verify-deployment.sh # Check local agent
|
||||
# ./verify-deployment.sh --ct=tanko # Check specific agent
|
||||
# ./verify-deployment.sh --all # Check all reachable agents
|
||||
#
|
||||
set -euo pipefail
|
||||
|
||||
PLUGIN_DIR="$HOME/.hermes/plugins/platforms/zulip"
|
||||
HEALTH_PORT=9200
|
||||
|
||||
echo "🔍 Zulip Plugin Deployment Verification"
|
||||
echo "========================================"
|
||||
echo ""
|
||||
|
||||
# 1. Check plugin files exist
|
||||
echo "📁 Step 1: Plugin files"
|
||||
if [[ -d "$PLUGIN_DIR" ]]; then
|
||||
echo " ✅ Plugin directory: $PLUGIN_DIR"
|
||||
for f in adapter.py __init__.py plugin.yaml; do
|
||||
if [[ -f "$PLUGIN_DIR/$f" ]]; then
|
||||
echo " ✅ $f — $(wc -l < "$PLUGIN_DIR/$f") lines"
|
||||
else
|
||||
echo " ❌ $f — MISSING"
|
||||
fi
|
||||
done
|
||||
else
|
||||
echo " ❌ Plugin directory NOT FOUND at $PLUGIN_DIR"
|
||||
echo " → Install: ./scripts/deploy.sh --mode=native v1.0.0"
|
||||
fi
|
||||
|
||||
echo ""
|
||||
|
||||
# 2. Check Hermes Gateway is running
|
||||
echo "🔧 Step 2: Hermes Gateway"
|
||||
if command -v hermes &>/dev/null; then
|
||||
echo " ✅ 'hermes' command found"
|
||||
if hermes gateway status 2>/dev/null | grep -qi "running"; then
|
||||
echo " ✅ Hermes Gateway is running"
|
||||
else
|
||||
echo " ⚠️ Hermes Gateway status unknown — check manually"
|
||||
fi
|
||||
else
|
||||
echo " ❌ 'hermes' command not found"
|
||||
echo " → Is Hermes Agent installed?"
|
||||
fi
|
||||
|
||||
echo ""
|
||||
|
||||
# 3. Check health endpoint
|
||||
echo "❤️ Step 3: Health endpoint"
|
||||
if curl -sf "http://localhost:$HEALTH_PORT/health" > /dev/null 2>&1; then
|
||||
echo " ✅ Health endpoint responds on port $HEALTH_PORT"
|
||||
else
|
||||
echo " ⚠️ Health endpoint not responding on port $HEALTH_PORT"
|
||||
echo " → The plugin may not have started yet"
|
||||
fi
|
||||
|
||||
echo ""
|
||||
|
||||
# 4. Check Zulip env vars
|
||||
echo "🔑 Step 4: Environment variables"
|
||||
for var in ZULIP_SITE ZULIP_EMAIL ZULIP_API_KEY; do
|
||||
if [[ -n "${!var:-}" ]]; then
|
||||
val="${!var}"
|
||||
if [[ "$var" == "ZULIP_API_KEY" ]]; then
|
||||
echo " ✅ $var — [REDACTED]"
|
||||
else
|
||||
echo " ✅ $var — $val"
|
||||
fi
|
||||
else
|
||||
echo " ❌ $var — NOT SET"
|
||||
fi
|
||||
done
|
||||
|
||||
echo ""
|
||||
|
||||
# Overall
|
||||
echo "═══════════════════════════════════════"
|
||||
missing=0
|
||||
[[ -d "$PLUGIN_DIR" ]] || missing=$((missing + 1))
|
||||
command -v hermes &>/dev/null || missing=$((missing + 1))
|
||||
[[ -n "${ZULIP_SITE:-}" && -n "${ZULIP_EMAIL:-}" && -n "${ZULIP_API_KEY:-}" ]] || missing=$((missing + 1))
|
||||
|
||||
if [[ "$missing" -eq 0 ]]; then
|
||||
echo "✅ VERDICT: Plugin properly deployed"
|
||||
echo " Send a DM to verify: @**${ZULIP_AGENT_NAME:-hermes-agent}** _hello_"
|
||||
else
|
||||
echo "⚠️ VERDICT: $missing issue(s) found — fix and re-verify"
|
||||
fi
|
||||
Reference in New Issue
Block a user