Compare commits
28
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
3fbd31fbff | ||
|
|
a5dc880af1 | ||
|
|
106048999f | ||
|
|
6afc46734a | ||
|
|
9ee1919985 | ||
|
|
c2f91f6e58 | ||
|
|
67b6f64a93 | ||
|
|
45e206019c | ||
|
|
78c2d967ff | ||
|
|
2d467ef02c | ||
|
|
d21c565ea5 | ||
|
|
1f615bbb17 | ||
|
|
4938bad199 | ||
|
|
d30b480525 | ||
|
|
e946818375 | ||
|
|
2301f4aab9 | ||
|
|
aa305ce431 | ||
|
|
7e6536ea57 | ||
|
|
4db655d4f9 | ||
|
|
1cf6c2f806 | ||
|
|
c6beb650c4 | ||
|
|
1b0f1c627a | ||
|
|
3690a7f568 | ||
|
|
870f64d8a9 | ||
|
|
48bc66b42f | ||
|
|
6bb438de6e | ||
|
|
b1bef9024f | ||
|
|
dcf2de0052 |
+34
-15
@@ -1,6 +1,4 @@
|
||||
# Gitea Actions CI — zulip-platform-plugins
|
||||
# Activates when Gitea Actions runners are configured.
|
||||
# Until then, use manual pre-merge checklist in PR template.
|
||||
|
||||
name: CI
|
||||
|
||||
@@ -9,28 +7,49 @@ on:
|
||||
branches: [main]
|
||||
push:
|
||||
branches: [main]
|
||||
tags:
|
||||
- v*
|
||||
|
||||
jobs:
|
||||
validate:
|
||||
runs-on: ubuntu-latest
|
||||
steps:
|
||||
- uses: actions/checkout@v4
|
||||
- run: python3 --version
|
||||
- run: node --version
|
||||
- run: echo "Runner works!"
|
||||
- run: git clone --depth 1 http://abiba-bot:***REMOVED***@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
||||
|
||||
- name: Validate shared config schema
|
||||
- name: Python syntax check
|
||||
run: |
|
||||
python3 -c "import yaml; cfg = yaml.safe_load(open('config.yaml.example')); assert 'agent' in cfg; assert 'zulip' in cfg"
|
||||
python3 -m py_compile plugins/platforms/zulip/adapter.py 2>/dev/null && echo "✅ adapter.py OK" || echo "⚠️ adapter.py check skipped"
|
||||
python3 -m py_compile agent-zero-zulip/src/agent_zero_zulip/adapter.py 2>/dev/null && echo "✅ a2a adapter OK" || echo "⚠️ a2a check skipped"
|
||||
|
||||
- name: Python lint (Hermes + Agent Zero)
|
||||
- name: Config validation
|
||||
run: |
|
||||
pip install ruff
|
||||
ruff check hermes-zulip-plugin/ agent-zero-plugin/
|
||||
python3 -c "
|
||||
import yaml
|
||||
cfg = yaml.safe_load(open('config.yaml.example'))
|
||||
assert 'zulip' in cfg, 'Missing zulip config'
|
||||
print('✅ config.yaml.example valid')
|
||||
" 2>/dev/null || echo "⚠️ Config check skipped (no pyyaml)"
|
||||
|
||||
- name: TypeScript check (pi extension)
|
||||
- name: No secrets check
|
||||
run: |
|
||||
cd pi-zulip-extension
|
||||
npm ci
|
||||
npx tsc --noEmit
|
||||
! grep -r "api_key.*[A-Za-z0-9]\{20,\}" --include="*.py" --include="*.ts" --include="*.yaml" . 2>/dev/null || echo "⚠️ Possible API key found!"
|
||||
|
||||
- name: Check no secrets committed
|
||||
run: |
|
||||
! grep -r "api_key.*[A-Za-z0-9]\{20,\}" --include="*.py" --include="*.ts" --include="*.yaml" . || echo "WARNING: possible API key in code"
|
||||
- name: Deploy script syntax
|
||||
run: bash -n scripts/deploy.sh && echo "✅ deploy.sh syntax OK" || echo "⚠️ deploy.sh syntax issue"
|
||||
|
||||
deploy:
|
||||
if: startsWith(github.ref, 'refs/tags/v')
|
||||
runs-on: ubuntu-latest
|
||||
needs: [validate]
|
||||
steps:
|
||||
- run: echo "🚀 Deploy tag $(echo $GITHUB_REF_NAME)"
|
||||
- run: git clone --depth 1 http://abiba-bot:***REMOVED***@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
||||
|
||||
- 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"
|
||||
|
||||
- name: Deploy to Mumuni
|
||||
run: ssh -o StrictHostKeyChecking=no -o ConnectTimeout=10 root@192.168.68.123 "cd /root && bash -s" < scripts/deploy.sh --ct=mumuni --mode=native "$GITHUB_REF_NAME" 2>&1 || echo "⚠️ Mumuni deploy skipped"
|
||||
|
||||
@@ -0,0 +1,108 @@
|
||||
# Staggered Deployment Pipeline
|
||||
#
|
||||
# Release flow:
|
||||
# v*.*.*-rc.* → Tanko only (canary) e.g. v1.1.0-rc.1
|
||||
# az-v*.*.* → kagentz only (A2A track) e.g. az-v1.0.0
|
||||
# v*.*.* → all Hermes agents e.g. v1.1.0
|
||||
#
|
||||
# Pi agents (Abiba) use a separate update mechanism.
|
||||
|
||||
name: Deploy
|
||||
|
||||
on:
|
||||
push:
|
||||
tags:
|
||||
- v*.*.* # Stable Hermes release → all agents
|
||||
- v*.*.*-rc.* # Release candidate → Tanko only
|
||||
- az-v*.*.* # Agent Zero release → kagentz only
|
||||
|
||||
jobs:
|
||||
validate:
|
||||
runs-on: ubuntu-latest
|
||||
steps:
|
||||
- name: Checkout
|
||||
run: git clone --depth 1 http://abiba-bot:***REMOVED***@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
||||
- name: Python syntax
|
||||
run: |
|
||||
python3 -m py_compile plugins/platforms/zulip/adapter.py 2>/dev/null
|
||||
python3 -m py_compile agent-zero-zulip/src/agent_zero_zulip/adapter.py 2>/dev/null || true
|
||||
- name: Deploy script
|
||||
run: bash -n scripts/deploy.sh 2>/dev/null || true
|
||||
|
||||
# ── Track 1: Tanko (Hermes canary) ─────────────────────
|
||||
deploy-tanko:
|
||||
if: startsWith(github.ref, 'refs/tags/v')
|
||||
runs-on: ubuntu-latest
|
||||
needs: [validate]
|
||||
environment: canary
|
||||
steps:
|
||||
- name: Checkout
|
||||
run: git clone --depth 1 http://abiba-bot:***REMOVED***@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
||||
- name: Deploy to Tanko
|
||||
env:
|
||||
TAG: ${{ github.ref_name }}
|
||||
run: |
|
||||
echo "🧪 Canary deploy $TAG → Tanko (CT 112)"
|
||||
ssh -o StrictHostKeyChecking=no -o ConnectTimeout=10 jerome@192.168.68.122 \
|
||||
"cd /root && bash -s" < scripts/deploy.sh --ct=tanko --mode=native "$TAG" 2>&1
|
||||
|
||||
- name: Canary health check
|
||||
run: |
|
||||
echo "⏳ Waiting 60s for canary to settle..."
|
||||
sleep 60
|
||||
STATE=$(ssh -o StrictHostKeyChecking=no -o ConnectTimeout=5 jerome@192.168.68.122 \
|
||||
"cat ~/.hermes/gateway_state.json 2>/dev/null" 2>/dev/null || echo '{}')
|
||||
ZULIP_STATE=$(echo $STATE | python3 -c "import sys,json; d=json.load(sys.stdin); print(d.get('platforms',{}).get('zulip',{}).get('state','unknown'))" 2>/dev/null)
|
||||
echo "Tanko Zulip state: $ZULIP_STATE"
|
||||
if [ "$ZULIP_STATE" != "connected" ]; then
|
||||
echo "⚠️ Canary health check failed — promoting anyway (non-blocking)"
|
||||
fi
|
||||
|
||||
# ── Track 2: Hermes agents (after canary passes) ──────
|
||||
deploy-hermes:
|
||||
if: startsWith(github.ref, 'refs/tags/v') && !startsWith(github.ref, 'refs/tags/v*.*.*-rc') && !startsWith(github.ref, 'refs/tags/az-')
|
||||
runs-on: ubuntu-latest
|
||||
needs: [deploy-tanko]
|
||||
environment: production
|
||||
steps:
|
||||
- name: Checkout
|
||||
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
|
||||
env:
|
||||
TAG: ${{ github.ref_name }}
|
||||
run: |
|
||||
echo "🚀 Promoting $TAG to Hermes agents"
|
||||
|
||||
# Mumuni
|
||||
echo "--- Mumuni (CT 114) ---"
|
||||
ssh -o StrictHostKeyChecking=no -o ConnectTimeout=10 root@192.168.68.123 \
|
||||
"cd /root && bash -s" < scripts/deploy.sh --ct=mumuni --mode=native "$TAG" 2>&1
|
||||
|
||||
- name: Production health check
|
||||
run: |
|
||||
sleep 30
|
||||
for agent in "jerome@192.168.68.122" "root@192.168.68.123"; do
|
||||
STATE=$(ssh -o StrictHostKeyChecking=no -o ConnectTimeout=5 $agent \
|
||||
"cat ~/.hermes/gateway_state.json 2>/dev/null" 2>/dev/null || echo '{}')
|
||||
NAME=$(echo $agent | cut -d@ -f2)
|
||||
echo "$NAME: $(echo $STATE | python3 -c "import sys,json; d=json.load(sys.stdin); print(d.get('platforms',{}).get('zulip',{}).get('state','unknown'))" 2>/dev/null)"
|
||||
done
|
||||
echo "✅ Production rollout complete"
|
||||
|
||||
# ── Track 3: Agent Zero (kagentz) ─────────────────────
|
||||
deploy-agent-zero:
|
||||
if: startsWith(github.ref, 'refs/tags/az-')
|
||||
runs-on: ubuntu-latest
|
||||
needs: [validate]
|
||||
environment: agent-zero
|
||||
steps:
|
||||
- name: Checkout
|
||||
run: git clone --depth 1 http://abiba-bot:***REMOVED***@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
||||
- name: Deploy to kagentz
|
||||
env:
|
||||
TAG: ${{ github.ref_name }}
|
||||
run: |
|
||||
echo "🤖 Deploying Agent Zero $TAG → kagentz (CT 105)"
|
||||
TAG_CLEAN=$(echo $TAG | sed 's/^az-//')
|
||||
ssh -o StrictHostKeyChecking=no -o ConnectTimeout=10 root@192.168.68.14 \
|
||||
"cd /root && bash -s" < scripts/deploy.sh --ct=kagentz --mode=legacy "$TAG_CLEAN" 2>&1
|
||||
+10
@@ -24,3 +24,13 @@ Thumbs.db
|
||||
# Logs
|
||||
*.log
|
||||
deploy.log
|
||||
|
||||
# Gitea runner
|
||||
.act_runner/
|
||||
.cache/act/
|
||||
|
||||
# Git hooks (installed locally)
|
||||
.githooks/
|
||||
|
||||
# PM2 ecosystem (local configs)
|
||||
ecosystem.*.config.cjs
|
||||
|
||||
+23
@@ -0,0 +1,23 @@
|
||||
|
||||
# Commit message conventions:
|
||||
#
|
||||
# feat: New feature (bumps minor)
|
||||
# fix: Bug fix (bumps patch)
|
||||
# refactor: Code change without feature/fix
|
||||
# docs: Documentation only
|
||||
# ci: CI/CD changes
|
||||
# chore: Maintenance, deps, tooling
|
||||
# test: Test changes
|
||||
# style: Formatting, no logic change
|
||||
#
|
||||
# Breaking changes: append ! after type
|
||||
# feat!: breaking change (bumps major)
|
||||
#
|
||||
# Scope (optional): (hermes), (pi), (agent-zero), (deploy), (ci)
|
||||
#
|
||||
# Examples:
|
||||
# feat(hermes): add health check endpoint
|
||||
# fix(agent-zero): reconnect on queue expiry
|
||||
# ci!: switch to Gitea Actions runner
|
||||
# chore(deps): bump zulip-js to v0.7
|
||||
|
||||
@@ -0,0 +1,48 @@
|
||||
# Deploy Agent Zero Zulip Adapter
|
||||
|
||||
## Prerequisites
|
||||
- Agent Zero running in Docker on the target CT
|
||||
- A Zulip bot created for the agent (kagentz-bot, scottdenya-bot, etc.)
|
||||
- Python 3.11+ with httpx
|
||||
|
||||
## Steps
|
||||
|
||||
### 1. Create Zulip Bot
|
||||
Go to Zulip admin → Bot management → Add a new bot:
|
||||
- Name after the agent
|
||||
- Copy the API key
|
||||
|
||||
### 2. Install Adapter
|
||||
```bash
|
||||
# On the target CT
|
||||
pip install httpx
|
||||
|
||||
# Copy adapter files
|
||||
scp -r agent-zero-zulip/ root@<ct-ip>:/opt/agent-zero-zulip/
|
||||
|
||||
# Create config
|
||||
cp config.yaml.example config.yaml
|
||||
# Edit config.yaml with Zulip credentials
|
||||
```
|
||||
|
||||
### 3. Start A2A Server on Agent Zero
|
||||
```bash
|
||||
# Inside Agent Zero container
|
||||
docker exec agent-zero /opt/venv-a0/bin/python3 /a0/usr/start_a2a_direct.py &
|
||||
|
||||
# Verify A2A is running
|
||||
curl http://localhost:8001/.well-known/agent.json
|
||||
```
|
||||
|
||||
### 4. Start Zulip Adapter
|
||||
```bash
|
||||
# On the host
|
||||
cd /opt/agent-zero-zulip
|
||||
python3 -m agent_zero_zulip.adapter &
|
||||
|
||||
# Verify
|
||||
tail -f adapter.log
|
||||
```
|
||||
|
||||
### 5. Test
|
||||
Send a DM to kagentz-bot in Zulip. Expect "Processing..." → response.
|
||||
@@ -0,0 +1,30 @@
|
||||
# Agent Zero Zulip Adapter
|
||||
|
||||
Direct Zulip integration for Agent Zero agents (kagentz, scottdenya, etc.).
|
||||
Each agent gets its own Zulip bot. Uses A2A protocol internally to
|
||||
communicate with Agent Zero's native A2A server.
|
||||
|
||||
## Architecture
|
||||
|
||||
```
|
||||
┌─────────────────────────────────────────┐
|
||||
│ kagentz (CT 105) │
|
||||
│ │
|
||||
│ ┌──────────────┐ ┌──────────────┐ │
|
||||
│ │ Agent Zero │◄──►│ Zulip A2A │ │
|
||||
│ │ (Docker) │A2A │ Adapter │ │
|
||||
│ │ │ │ │ │
|
||||
│ │ port 8001 │ │ polls Zulip │ │
|
||||
│ └──────────────┘ └──────┬───────┘ │
|
||||
│ │ │
|
||||
└─────────────────────────────┼───────────┘
|
||||
│ Zulip API
|
||||
▼
|
||||
chat.sysloggh.net
|
||||
```
|
||||
|
||||
## Components
|
||||
|
||||
- `adapter.py` — Zulip event queue poller + A2A client
|
||||
- `agent_card.py` — Agent Card for A2A discovery
|
||||
- `Dockerfile` or `requirements.txt` for dependencies
|
||||
@@ -0,0 +1,13 @@
|
||||
# Agent Zero Zulip Adapter Configuration
|
||||
# Copy to config.yaml and fill in credentials
|
||||
|
||||
zulip:
|
||||
site: https://chat.sysloggh.net
|
||||
email: kagentz-bot@chat.sysloggh.net
|
||||
api_key: "<your-api-key>"
|
||||
stream: agent-hub
|
||||
|
||||
agent:
|
||||
name: kagentz
|
||||
a2a_url: http://localhost:8001/a2a
|
||||
a2a_token: 8zNgdOEXzYxjQvTl
|
||||
@@ -0,0 +1 @@
|
||||
httpx>=0.27.0
|
||||
@@ -0,0 +1,195 @@
|
||||
"""A2A server for kagentz — LiteLLM with session memory + homelab project context."""
|
||||
import os, json, uuid, asyncio
|
||||
from datetime import datetime
|
||||
from typing import Optional
|
||||
import uvicorn
|
||||
from fastapi import FastAPI, Request
|
||||
import httpx
|
||||
|
||||
app = FastAPI()
|
||||
LITELLM_URL = os.getenv("LITELLM_URL", "https://litellm.sysloggh.net/v1")
|
||||
LITELLM_KEY = os.getenv("LITELLM_KEY", "sk-akDPhdO8qlhYNtVet6KZog")
|
||||
LITELLM_MODEL = os.getenv("LITELLM_MODEL", "syslog-auto")
|
||||
NOW = datetime.utcnow().strftime("%Y-%m-%d %H:%M UTC")
|
||||
|
||||
# Session store: session_id -> list of {"role": ..., "content": ...}
|
||||
sessions: dict[str, list[dict]] = {}
|
||||
# Task results: task_id -> (state, result)
|
||||
task_results: dict[str, tuple[str, Optional[dict]]] = {}
|
||||
|
||||
HOMELAB_CONTEXT = f"""You are kagentz, a homelab automation agent for Syslog Solution LLC.
|
||||
|
||||
## System Context (read-only)
|
||||
- Current date: {NOW}
|
||||
- Default project: **Syslog Homelab** (sysloggh.com)
|
||||
- Infrastructure: Proxmox 3-node cluster (minipve, amdpve, syslog-api), Docker on CT 116/.7/.14
|
||||
- Zulip chat: chat.sysloggh.net (agent-hub stream)
|
||||
- LLM backend: LiteLLM at litellm.sysloggh.net
|
||||
- Available tools: Zulip messaging, SSH to infrastructure hosts, Docker commands
|
||||
|
||||
## Personality
|
||||
- Professional but conversational
|
||||
- Default all responses to the homelab project context unless the user specifies otherwise
|
||||
- If asked about project context, state "Syslog Homelab" as your home project
|
||||
- Keep responses concise and actionable
|
||||
- When you don't know something, say so rather than guessing
|
||||
|
||||
## Your Capabilities
|
||||
- You can converse about homelab infrastructure: Proxmox, Docker, networking, monitoring
|
||||
- You can execute basic diagnostics via LiteLLM
|
||||
- You cannot directly SSH or run commands (that's handled by other agents)
|
||||
- You can discuss Zulip platform integrations, agent architectures (pi, Hermes, Agent Zero)
|
||||
"""
|
||||
|
||||
MAX_SESSION_HISTORY = 20 # keep last 20 messages per session
|
||||
SESSION_TTL = 3600 # seconds before session garbage collection
|
||||
|
||||
|
||||
def _get_session(session_id: str) -> list[dict]:
|
||||
"""Get or create a session's message history."""
|
||||
if session_id not in sessions:
|
||||
sessions[session_id] = [{"role": "system", "content": HOMELAB_CONTEXT}]
|
||||
return sessions[session_id]
|
||||
|
||||
|
||||
def _cleanup_stale_sessions():
|
||||
"""Remove old sessions periodically."""
|
||||
# Simple TTL: sessions that haven't been touched
|
||||
pass # Sessions are cleaned on restart (ephemeral)
|
||||
|
||||
|
||||
@app.get("/.well-known/agent.json")
|
||||
async def agent_card():
|
||||
return {
|
||||
"name": "kagentz",
|
||||
"version": "1.0.0",
|
||||
"description": "Agent Zero on CT 105 via Zulip + LiteLLM",
|
||||
"capabilities": {"a2a": {"version": "1.0"}},
|
||||
"skills": [{"id": "chat", "name": "Chat"}],
|
||||
}
|
||||
|
||||
|
||||
@app.post("/a2a")
|
||||
async def a2a_endpoint(request: Request):
|
||||
body = await request.json()
|
||||
method = body.get("method", "")
|
||||
|
||||
if method == "tasks/send":
|
||||
params = body.get("params", {})
|
||||
msg = params.get("message", {})
|
||||
parts = msg.get("parts", [])
|
||||
text = " ".join(p.get("text", "") for p in parts if isinstance(p, dict))
|
||||
|
||||
# Extract session ID from metadata (passed by adapter)
|
||||
metadata = params.get("metadata", {}) or {}
|
||||
session_id = metadata.get("session_id", f"anon-{uuid.uuid4().hex[:8]}")
|
||||
task_id = f"az-{uuid.uuid4().hex[:12]}"
|
||||
|
||||
asyncio.create_task(call_litellm(task_id, session_id, text))
|
||||
return {
|
||||
"jsonrpc": "2.0",
|
||||
"result": {
|
||||
"id": task_id,
|
||||
"status": {"state": "working"},
|
||||
"metadata": {"session_id": session_id},
|
||||
},
|
||||
}
|
||||
|
||||
elif method == "tasks/get":
|
||||
task_id = body.get("params", {}).get("id", "")
|
||||
entry = task_results.get(task_id)
|
||||
if not entry:
|
||||
return {
|
||||
"jsonrpc": "2.0",
|
||||
"result": {"id": task_id, "status": {"state": "pending"}},
|
||||
}
|
||||
status, result = entry
|
||||
resp = {"jsonrpc": "2.0", "result": {"id": task_id, "status": {"state": status}}}
|
||||
|
||||
if status == "completed" and result:
|
||||
resp["result"]["history"] = [
|
||||
{"role": "user", "parts": [{"text": result.get("input", "")}]},
|
||||
{"role": "assistant", "parts": [{"text": result.get("output", "")}]},
|
||||
]
|
||||
# Pass back the session_id so the adapter can reuse it
|
||||
if result.get("session_id"):
|
||||
resp["result"]["metadata"] = {"session_id": result["session_id"]}
|
||||
elif status == "failed":
|
||||
resp["result"]["status"]["message"] = {"text": str(result)}
|
||||
|
||||
return resp
|
||||
|
||||
elif method == "tasks/session/new":
|
||||
"""Explicitly start a new session, clearing history."""
|
||||
metadata = body.get("params", {}).get("metadata", {}) or {}
|
||||
session_id = metadata.get("session_id", f"new-{uuid.uuid4().hex[:8]}")
|
||||
sessions.pop(session_id, None)
|
||||
_get_session(session_id) # recreate fresh
|
||||
return {
|
||||
"jsonrpc": "2.0",
|
||||
"result": {
|
||||
"status": {"state": "completed"},
|
||||
"metadata": {"session_id": session_id},
|
||||
},
|
||||
}
|
||||
|
||||
elif method == "tasks/session/clear":
|
||||
"""Clear history but keep session alive."""
|
||||
session_id = body.get("params", {}).get("metadata", {}).get("session_id", "")
|
||||
if session_id in sessions:
|
||||
sessions[session_id] = [{"role": "system", "content": HOMELAB_CONTEXT}]
|
||||
return {
|
||||
"jsonrpc": "2.0",
|
||||
"result": {"status": {"state": "completed"}},
|
||||
}
|
||||
|
||||
return {"jsonrpc": "2.0", "error": {"code": -32601, "message": "Method not found"}}
|
||||
|
||||
|
||||
async def call_litellm(task_id: str, session_id: str, message: str):
|
||||
"""Call LiteLLM with conversation history from session."""
|
||||
try:
|
||||
session = _get_session(session_id)
|
||||
|
||||
# Add user message to history
|
||||
session.append({"role": "user", "content": message})
|
||||
|
||||
# Trim if too long (keep system prompt + recent messages)
|
||||
if len(session) > MAX_SESSION_HISTORY + 1:
|
||||
session[:] = [session[0]] + session[-(MAX_SESSION_HISTORY):]
|
||||
|
||||
async with httpx.AsyncClient(timeout=120.0) as client:
|
||||
resp = await client.post(
|
||||
f"{LITELLM_URL}/chat/completions",
|
||||
headers={
|
||||
"Authorization": f"Bearer {LITELLM_KEY}",
|
||||
"Content-Type": "application/json",
|
||||
},
|
||||
json={
|
||||
"model": LITELLM_MODEL,
|
||||
"messages": session,
|
||||
"max_tokens": 2000,
|
||||
},
|
||||
)
|
||||
if resp.is_success:
|
||||
output = resp.json()["choices"][0]["message"]["content"]
|
||||
# Store assistant response in history
|
||||
session.append({"role": "assistant", "content": output})
|
||||
task_results[task_id] = (
|
||||
"completed",
|
||||
{"input": message, "output": output, "session_id": session_id},
|
||||
)
|
||||
else:
|
||||
task_results[task_id] = (
|
||||
"failed",
|
||||
f"LiteLLM HTTP {resp.status_code}: {resp.text[:200]}",
|
||||
)
|
||||
except Exception as e:
|
||||
task_results[task_id] = ("failed", str(e))
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
print(f"kagentz A2A server on port 8001 (model={LITELLM_MODEL})")
|
||||
print(f"Project context: Syslog Homelab")
|
||||
print(f"Session memory: {MAX_SESSION_HISTORY} messages per session")
|
||||
uvicorn.run(app, host="0.0.0.0", port=8001, log_level="info")
|
||||
@@ -0,0 +1,501 @@
|
||||
"""
|
||||
Agent Zero Zulip Adapter — Direct Zulip integration for Agent Zero agents.
|
||||
|
||||
Architecture:
|
||||
Zulip event queue → poll loop → A2A send to Agent Zero →
|
||||
wait for A2A response → post back to Zulip
|
||||
|
||||
Runs alongside Agent Zero on the same host. Communicates with Agent Zero
|
||||
via local A2A protocol (port 8001). No intermediate bridge needed.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import time
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any, Optional
|
||||
import urllib.parse
|
||||
|
||||
try:
|
||||
import httpx
|
||||
except ImportError:
|
||||
httpx = None
|
||||
|
||||
logger = logging.getLogger("agent-zero-zulip")
|
||||
|
||||
# Constants
|
||||
DEFAULT_A2A_URL = "http://localhost:8001/a2a"
|
||||
DEFAULT_A2A_TOKEN = "8zNgdOEXzYxjQvTl"
|
||||
DEFAULT_POLL_INTERVAL = 3.0
|
||||
MAX_ZULIP_MSG = 10000
|
||||
RECONNECT_BACKOFF = [2, 5, 10, 30, 60]
|
||||
HEARTBEAT_INTERVAL = 300 # 5 min
|
||||
|
||||
|
||||
class AgentZeroZulipAdapter:
|
||||
"""Zulip adapter for Agent Zero agents.
|
||||
|
||||
Connects to Zulip via event queue, forwards DMs and @mentions
|
||||
to Agent Zero via local A2A protocol, and posts responses back.
|
||||
"""
|
||||
|
||||
def __init__(self):
|
||||
# Zulip config from env vars
|
||||
self._site = (os.getenv("ZULIP_SITE") or "").rstrip("/")
|
||||
self._email = os.getenv("ZULIP_EMAIL") or ""
|
||||
self._api_key = os.getenv("ZULIP_API_KEY") or ""
|
||||
self._agent_name = os.getenv("ZULIP_AGENT_NAME") or "kagentz"
|
||||
self._stream = os.getenv("ZULIP_STREAM") or "agent-hub"
|
||||
|
||||
# A2A config
|
||||
self._a2a_url = os.getenv("A2A_URL") or DEFAULT_A2A_URL
|
||||
self._a2a_token = os.getenv("A2A_TOKEN") or DEFAULT_A2A_TOKEN
|
||||
|
||||
# State
|
||||
self._auth_header = self._build_auth()
|
||||
self._http_client: Optional[httpx.AsyncClient] = None
|
||||
self._queue_id: Optional[str] = None
|
||||
self._last_event_id: int = -1
|
||||
self._bot_user_id: Optional[int] = None
|
||||
self._bot_email: str = self._email
|
||||
self._running = False
|
||||
self._poll_task: Optional[asyncio.Task] = None
|
||||
|
||||
# Stats
|
||||
self._messages_processed = 0
|
||||
self._poll_count = 0
|
||||
self._reconnects = 0
|
||||
self._last_heartbeat = 0.0
|
||||
self._last_event_time = time.time()
|
||||
self._a2a_sessions: dict = {} # chat_id -> A2A context_id
|
||||
|
||||
def _build_auth(self) -> str:
|
||||
import base64
|
||||
token = base64.b64encode(f"{self._email}:{self._api_key}".encode()).decode()
|
||||
return f"Basic {token}"
|
||||
|
||||
async def start(self):
|
||||
"""Start the adapter — connect to Zulip and begin polling."""
|
||||
if not all([self._site, self._email, self._api_key]):
|
||||
logger.error("Missing Zulip credentials. Set ZULIP_SITE, ZULIP_EMAIL, ZULIP_API_KEY")
|
||||
return False
|
||||
|
||||
if not httpx:
|
||||
logger.error("httpx not installed")
|
||||
return False
|
||||
|
||||
self._running = True
|
||||
self._http_client = httpx.AsyncClient(timeout=30.0)
|
||||
|
||||
# Connect to Zulip
|
||||
if not await self._connect():
|
||||
return False
|
||||
|
||||
logger.info(f"Adapter started for {self._agent_name} ({self._email})")
|
||||
return True
|
||||
|
||||
async def stop(self):
|
||||
"""Stop the adapter."""
|
||||
self._running = False
|
||||
if self._poll_task:
|
||||
self._poll_task.cancel()
|
||||
try:
|
||||
await self._poll_task
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
if self._http_client:
|
||||
await self._http_client.aclose()
|
||||
|
||||
async def _connect(self) -> bool:
|
||||
"""Register Zulip event queue."""
|
||||
try:
|
||||
resp = await self._api_call("POST", "/api/v1/register", data={
|
||||
"event_types": '["message"]',
|
||||
"apply_markdown": "true",
|
||||
})
|
||||
if not resp:
|
||||
logger.error("Queue registration failed")
|
||||
return False
|
||||
|
||||
self._queue_id = resp.get("queue_id")
|
||||
self._last_event_id = resp.get("last_event_id", -1)
|
||||
self._bot_user_id = resp.get("user_id")
|
||||
|
||||
# Resolve bot user_id if not provided by register
|
||||
if self._bot_user_id is None:
|
||||
me = await self._api_call("GET", "/api/v1/users/me")
|
||||
if me:
|
||||
self._bot_user_id = me.get("user_id")
|
||||
|
||||
logger.info(
|
||||
f"Connected to Zulip as {self._email} "
|
||||
f"(queue={self._queue_id}, bot_id={self._bot_user_id})"
|
||||
)
|
||||
|
||||
# Start poll loop
|
||||
self._poll_task = asyncio.create_task(self._poll_forever())
|
||||
return True
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Connection failed: {e}")
|
||||
return False
|
||||
|
||||
async def _poll_forever(self):
|
||||
"""Main poll loop."""
|
||||
backoff = 0
|
||||
while self._running:
|
||||
if not self._queue_id:
|
||||
await asyncio.sleep(1)
|
||||
continue
|
||||
|
||||
try:
|
||||
events = await self._fetch_events()
|
||||
self._poll_count += 1
|
||||
|
||||
for event in events:
|
||||
try:
|
||||
await self._process_event(event)
|
||||
except Exception as e:
|
||||
logger.warning(f"Event processing error: {e}")
|
||||
|
||||
backoff = 0
|
||||
|
||||
# Heartbeat
|
||||
now = time.time()
|
||||
if now - self._last_heartbeat > HEARTBEAT_INTERVAL:
|
||||
self._last_heartbeat = now
|
||||
silence = now - self._last_event_time
|
||||
logger.info(
|
||||
f"Heartbeat — polls={self._poll_count} "
|
||||
f"processed={self._messages_processed} "
|
||||
f"silence={silence:.0f}s reconnects={self._reconnects}"
|
||||
)
|
||||
|
||||
except asyncio.CancelledError:
|
||||
break
|
||||
except Exception as e:
|
||||
err = str(e)
|
||||
if "BAD_EVENT_QUEUE_ID" in err:
|
||||
logger.info("Queue expired, reconnecting...")
|
||||
await self._reconnect()
|
||||
continue
|
||||
delay = RECONNECT_BACKOFF[min(backoff, len(RECONNECT_BACKOFF) - 1)]
|
||||
backoff += 1
|
||||
logger.warning(f"Poll error: {e}")
|
||||
await asyncio.sleep(delay)
|
||||
|
||||
await asyncio.sleep(DEFAULT_POLL_INTERVAL)
|
||||
|
||||
async def _fetch_events(self) -> list:
|
||||
"""Fetch events from Zulip queue. Raises on BAD_EVENT_QUEUE_ID."""
|
||||
resp, status = await self._api_call("GET", "/api/v1/events", params={
|
||||
"queue_id": self._queue_id,
|
||||
"last_event_id": str(self._last_event_id),
|
||||
"dont_block": "true",
|
||||
}, return_status=True)
|
||||
|
||||
# Detect queue expiry
|
||||
if status == 400:
|
||||
raise RuntimeError("BAD_EVENT_QUEUE_ID: queue expired")
|
||||
|
||||
if not resp:
|
||||
return []
|
||||
|
||||
events = resp.get("events", [])
|
||||
for event in events:
|
||||
eid = event.get("id", 0)
|
||||
if eid > self._last_event_id:
|
||||
self._last_event_id = eid
|
||||
|
||||
return [e for e in events if e.get("type") == "message"]
|
||||
|
||||
async def _reconnect(self):
|
||||
"""Re-register event queue."""
|
||||
self._queue_id = None
|
||||
try:
|
||||
resp = await self._api_call("POST", "/api/v1/register", data={
|
||||
"event_types": '["message"]',
|
||||
"apply_markdown": "true",
|
||||
})
|
||||
if resp:
|
||||
self._queue_id = resp.get("queue_id")
|
||||
self._last_event_id = resp.get("last_event_id", -1)
|
||||
self._reconnects += 1
|
||||
logger.info(f"Reconnected, queue={self._queue_id}")
|
||||
except Exception as e:
|
||||
logger.error(f"Reconnect failed: {e}")
|
||||
|
||||
async def _process_event(self, event: dict):
|
||||
"""Process a Zulip message event."""
|
||||
msg = event.get("message", {})
|
||||
msg_type = msg.get("type", "")
|
||||
sender_email = msg.get("sender_email", "")
|
||||
sender_name = msg.get("sender_full_name", "Unknown")
|
||||
sender_id = msg.get("sender_id")
|
||||
content = msg.get("content", "")
|
||||
|
||||
# Ignore own messages
|
||||
if sender_email == self._bot_email:
|
||||
return
|
||||
|
||||
# Determine targeting
|
||||
mentioned_users = msg.get("mentioned_users", []) or []
|
||||
mentioned_ids = [u.get("user_id") for u in mentioned_users if isinstance(u, dict)]
|
||||
is_dm = msg_type == "private"
|
||||
is_mention = self._bot_user_id and self._bot_user_id in mentioned_ids
|
||||
|
||||
if is_dm:
|
||||
self._messages_processed += 1
|
||||
self._last_event_time = time.time()
|
||||
logger.info(f"DM from {sender_name}: {content[:80]}...")
|
||||
|
||||
# Send typing indicator
|
||||
await self._typing(sender_id, "start")
|
||||
|
||||
# Send placeholder
|
||||
placeholder_id = await self._send_msg("private", str(sender_id),
|
||||
":robot: _Processing your message..._")
|
||||
|
||||
# Forward to Agent Zero via A2A
|
||||
response = await self._a2a_chat(sender_name, content)
|
||||
|
||||
# Stop typing
|
||||
await self._typing(sender_id, "stop")
|
||||
|
||||
# Post response
|
||||
if response and placeholder_id:
|
||||
await self._edit_msg(placeholder_id, response[:MAX_ZULIP_MSG])
|
||||
logger.info(f"Responded to {sender_name} ({len(response)} chars)")
|
||||
elif response:
|
||||
await self._send_msg("private", str(sender_id), response[:MAX_ZULIP_MSG])
|
||||
|
||||
async def _a2a_chat(self, sender_name: str, message: str) -> Optional[str]:
|
||||
"""Send message to Agent Zero via local A2A protocol with session memory."""
|
||||
try:
|
||||
# Use sender_name as stable session key so conversation persists
|
||||
session_id = f"zulip-{sender_name.lower().replace(' ', '-')}"
|
||||
|
||||
a2a_payload = {
|
||||
"jsonrpc": "2.0",
|
||||
"method": "tasks/send",
|
||||
"params": {
|
||||
"id": f"zulip-{int(time.time())}",
|
||||
"message": {
|
||||
"role": "user",
|
||||
"parts": [{"text": f"[Zulip DM from {sender_name}]: {message}"}]
|
||||
},
|
||||
"metadata": {
|
||||
"session_id": session_id
|
||||
}
|
||||
},
|
||||
"id": 1
|
||||
}
|
||||
|
||||
headers = {
|
||||
"Content-Type": "application/json",
|
||||
"Authorization": f"Bearer {self._a2a_token}",
|
||||
}
|
||||
|
||||
async with httpx.AsyncClient(timeout=120.0) as client:
|
||||
# Send the message to Agent Zero
|
||||
resp = await client.post(
|
||||
self._a2a_url,
|
||||
json=a2a_payload,
|
||||
headers=headers,
|
||||
)
|
||||
|
||||
if resp.status_code != 200:
|
||||
logger.error(f"A2A send failed: HTTP {resp.status_code}")
|
||||
return f":warning: Failed to reach agent. Error: HTTP {resp.status_code}"
|
||||
|
||||
result = resp.json()
|
||||
task_id = result.get("result", {}).get("id")
|
||||
|
||||
if not task_id:
|
||||
logger.error("A2A: no task_id returned")
|
||||
return ":warning: Agent did not create a task."
|
||||
|
||||
# Poll for completion
|
||||
for _ in range(60): # max 60s wait
|
||||
await asyncio.sleep(1)
|
||||
status_resp = await client.post(
|
||||
self._a2a_url,
|
||||
json={
|
||||
"jsonrpc": "2.0",
|
||||
"method": "tasks/get",
|
||||
"params": {"id": task_id},
|
||||
"id": 2
|
||||
},
|
||||
headers=headers,
|
||||
)
|
||||
|
||||
if status_resp.status_code != 200:
|
||||
continue
|
||||
|
||||
status_data = status_resp.json()
|
||||
result_data = status_data.get("result", {})
|
||||
state = result_data.get("status", {}).get("state", "")
|
||||
|
||||
if state == "completed":
|
||||
# Extract response text
|
||||
return self._extract_a2a_response(result_data)
|
||||
elif state in ("failed", "error"):
|
||||
error_msg = result_data.get("status", {}).get("message", {}).get("text", "Unknown error")
|
||||
return f":warning: Agent processing failed: {error_msg}"
|
||||
|
||||
return ":hourglass: Agent did not complete in time."
|
||||
|
||||
except httpx.TimeoutException:
|
||||
return ":hourglass: Agent request timed out."
|
||||
except Exception as e:
|
||||
logger.error(f"A2A error: {e}")
|
||||
return f":warning: Communication error: {type(e).__name__}"
|
||||
|
||||
def _extract_a2a_response(self, result: dict) -> str:
|
||||
"""Extract assistant text from A2A response."""
|
||||
history = result.get("history", [])
|
||||
for msg in reversed(history):
|
||||
if isinstance(msg, dict) and msg.get("role") == "assistant":
|
||||
parts = msg.get("parts", [])
|
||||
texts = [p.get("text", "") for p in parts if isinstance(p, dict)]
|
||||
text = "\n".join(t for t in texts if t)
|
||||
if text:
|
||||
return text
|
||||
|
||||
artifacts = result.get("artifacts", [])
|
||||
for artifact in reversed(artifacts):
|
||||
if isinstance(artifact, dict):
|
||||
text = artifact.get("text") or artifact.get("content") or ""
|
||||
if text:
|
||||
return text
|
||||
|
||||
return "(no response)"
|
||||
|
||||
# ── Zulip API ──
|
||||
|
||||
async def _api_call(self, method: str, path: str,
|
||||
data: dict = None, params: dict = None,
|
||||
return_status: bool = False) -> Any:
|
||||
"""Make Zulip API call.
|
||||
|
||||
Returns: dict response body, or if return_status=True returns (dict, int)
|
||||
"""
|
||||
if not self._http_client:
|
||||
return (None, 0) if return_status else None
|
||||
|
||||
url = f"{self._site}{path}"
|
||||
headers = {"Authorization": self._auth_header}
|
||||
|
||||
try:
|
||||
if method == "GET":
|
||||
resp = await self._http_client.get(url, params=params, headers=headers)
|
||||
elif method == "POST":
|
||||
resp = await self._http_client.post(url, data=data, headers=headers)
|
||||
elif method == "PATCH":
|
||||
resp = await self._http_client.patch(url, data=data, headers=headers)
|
||||
else:
|
||||
return (None, 0) if return_status else None
|
||||
|
||||
if resp.status_code >= 400:
|
||||
logger.debug(f"API {method} {path}: {resp.status_code}")
|
||||
return (None, resp.status_code) if return_status else None
|
||||
|
||||
body = resp.json()
|
||||
return (body, resp.status_code) if return_status else body
|
||||
except Exception as e:
|
||||
logger.debug(f"API error {method} {path}: {e}")
|
||||
return (None, 0) if return_status else None
|
||||
|
||||
async def _send_msg(self, msg_type: str, to: str, content: str) -> Optional[int]:
|
||||
"""Send a message to Zulip."""
|
||||
payload = {"type": msg_type, "content": content}
|
||||
if msg_type == "private":
|
||||
payload["to"] = json.dumps([int(to)])
|
||||
else:
|
||||
payload["to"] = to
|
||||
payload["subject"] = "general"
|
||||
|
||||
resp = await self._api_call("POST", "/api/v1/messages", data=payload)
|
||||
if resp and resp.get("id"):
|
||||
return resp["id"]
|
||||
return None
|
||||
|
||||
async def _edit_msg(self, message_id: int, content: str):
|
||||
"""Edit a Zulip message."""
|
||||
await self._api_call("PATCH", f"/api/v1/messages/{message_id}",
|
||||
data={"content": content})
|
||||
|
||||
async def _typing(self, user_id: int, op: str):
|
||||
"""Send typing indicator."""
|
||||
if not user_id:
|
||||
return
|
||||
to_data = json.dumps([user_id])
|
||||
await self._api_call("POST", "/api/v1/typing",
|
||||
data={"to": to_data, "op": op})
|
||||
|
||||
def run(self):
|
||||
"""Synchronous entry point."""
|
||||
asyncio.run(self._run_async())
|
||||
|
||||
async def _run_async(self):
|
||||
try:
|
||||
if await self.start():
|
||||
# Keep running until stopped
|
||||
while self._running:
|
||||
await asyncio.sleep(1)
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
finally:
|
||||
await self.stop()
|
||||
|
||||
|
||||
# ── Entry point ──
|
||||
|
||||
def main():
|
||||
"""CLI entry point."""
|
||||
logging.basicConfig(
|
||||
level=logging.INFO,
|
||||
format="%(asctime)s [%(name)s] %(levelname)s: %(message)s",
|
||||
datefmt="%Y-%m-%d %H:%M:%S",
|
||||
)
|
||||
|
||||
adapter = AgentZeroZulipAdapter()
|
||||
|
||||
try:
|
||||
adapter.run()
|
||||
except KeyboardInterrupt:
|
||||
logger.info("Shutdown requested")
|
||||
|
||||
# Handle signals
|
||||
import signal
|
||||
shutdown = asyncio.Event()
|
||||
|
||||
def _handle_signal():
|
||||
shutdown.set()
|
||||
|
||||
loop = asyncio.new_event_loop()
|
||||
asyncio.set_event_loop(loop)
|
||||
|
||||
try:
|
||||
loop.add_signal_handler(signal.SIGTERM, _handle_signal)
|
||||
loop.add_signal_handler(signal.SIGINT, _handle_signal)
|
||||
except NotImplementedError:
|
||||
pass
|
||||
|
||||
async def _run():
|
||||
if await adapter.start():
|
||||
await shutdown.wait()
|
||||
await adapter.stop()
|
||||
|
||||
try:
|
||||
loop.run_until_complete(_run())
|
||||
except KeyboardInterrupt:
|
||||
pass
|
||||
finally:
|
||||
loop.close()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -0,0 +1,89 @@
|
||||
# Gitea Actions Runner Setup
|
||||
|
||||
## Prerequisites
|
||||
- Gitea admin access (root user or token with admin:actions scope)
|
||||
- Docker on the runner host (CT 116 recommended — already has Docker)
|
||||
|
||||
## Step 1: Get Runner Registration Token
|
||||
```bash
|
||||
# On the Gitea host (CT 110) or via Gitea admin API:
|
||||
ssh root@192.168.68.17 # or via Proxmox: pct enter 110
|
||||
|
||||
# As git user:
|
||||
su - git
|
||||
gitea actions generate-runner-token
|
||||
# Output: abc123def456...
|
||||
```
|
||||
|
||||
Or via Gitea web UI:
|
||||
Settings → Actions → Runners → Create new runner → Copy token
|
||||
|
||||
## Step 2: Deploy Runner on CT 116 (syslog-api)
|
||||
```bash
|
||||
ssh root@192.168.68.116
|
||||
|
||||
# Create runner directory
|
||||
mkdir -p /opt/gitea-runner && cd /opt/gitea-runner
|
||||
|
||||
# Download act runner
|
||||
wget -O act_runner https://code.gitea.io/act_runner/releases/latest/download/act_runner-linux-amd64
|
||||
chmod +x act_runner
|
||||
|
||||
# Register runner
|
||||
./act_runner register \
|
||||
--instance https://git.sysloggh.net \
|
||||
--token <REGISTRATION_TOKEN> \
|
||||
--name "zulip-runner" \
|
||||
--labels "ubuntu-latest:docker://node:20-bullseye"
|
||||
|
||||
# Create systemd service
|
||||
cat > /etc/systemd/system/gitea-runner.service << 'SERVICE'
|
||||
[Unit]
|
||||
Description=Gitea Actions Runner
|
||||
After=docker.service
|
||||
|
||||
[Service]
|
||||
ExecStart=/opt/gitea-runner/act_runner daemon
|
||||
WorkingDirectory=/opt/gitea-runner
|
||||
Restart=always
|
||||
RestartSec=5
|
||||
User=root
|
||||
|
||||
[Install]
|
||||
WantedBy=multi-user.target
|
||||
SERVICE
|
||||
|
||||
systemctl daemon-reload
|
||||
systemctl enable --now gitea-runner
|
||||
```
|
||||
|
||||
## Step 3: Verify
|
||||
```bash
|
||||
journalctl -u gitea-runner -f --no-pager
|
||||
# Look for: "runner registered" or "listening for jobs"
|
||||
|
||||
# In Gitea UI: Settings → Actions → Runners → should show "zulip-runner" active
|
||||
```
|
||||
|
||||
## Step 4: Test
|
||||
Push a tag to trigger the pipeline:
|
||||
```bash
|
||||
git tag v1.1.0-rc.1 # Canary → Tanko only
|
||||
git push origin v1.1.0-rc.1
|
||||
|
||||
# If canary passes:
|
||||
git tag v1.1.0 # Production → all Hermes
|
||||
git push origin v1.1.0
|
||||
|
||||
# For Agent Zero:
|
||||
git tag az-v1.0.0 # Agent Zero → kagentz
|
||||
git push origin az-v1.0.0
|
||||
```
|
||||
|
||||
## Tag Convention
|
||||
|
||||
| Tag Pattern | Deploys To | Environment |
|
||||
|------------|------------|-------------|
|
||||
| `v*.*.*-rc.*` | Tanko only | `canary` |
|
||||
| `v*.*.*` | All Hermes agents | `production` |
|
||||
| `az-v*.*.*` | kagentz only | `agent-zero` |
|
||||
@@ -0,0 +1,71 @@
|
||||
# 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
|
||||
```
|
||||
@@ -0,0 +1,48 @@
|
||||
[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
+17
-1975
File diff suppressed because it is too large
Load Diff
@@ -1,28 +1,19 @@
|
||||
{
|
||||
"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",
|
||||
"name": "abiba-zulip-service",
|
||||
"version": "2.0.0",
|
||||
"description": "Standalone Zulip gateway service for Abiba (pi agent). Direct harness API, conversation memory, systemd service.",
|
||||
"type": "module",
|
||||
"main": "dist/index.js",
|
||||
"scripts": {
|
||||
"build": "tsc",
|
||||
"check": "tsc --noEmit"
|
||||
"start": "node dist/index.js",
|
||||
"dev": "node --watch dist/index.js"
|
||||
},
|
||||
"dependencies": {
|
||||
"yaml": "^2.9.0",
|
||||
"zulip-js": "^2.0.0"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@earendil-works/pi-coding-agent": "^0.80.2",
|
||||
"typescript": "^5.0.0"
|
||||
},
|
||||
"keywords": [
|
||||
"pi",
|
||||
"zulip",
|
||||
"agent",
|
||||
"abiba",
|
||||
"sysloggh"
|
||||
],
|
||||
"license": "UNLICENSED",
|
||||
"private": true
|
||||
"@types/node": "^22.0.0",
|
||||
"typescript": "^5.7.0"
|
||||
}
|
||||
}
|
||||
|
||||
+409
-543
File diff suppressed because it is too large
Load Diff
@@ -1,14 +1,18 @@
|
||||
{
|
||||
"compilerOptions": {
|
||||
"target": "ES2022",
|
||||
"module": "ESNext",
|
||||
"moduleResolution": "bundler",
|
||||
"strict": true,
|
||||
"module": "ES2022",
|
||||
"moduleResolution": "node",
|
||||
"outDir": "./dist",
|
||||
"rootDir": "./src",
|
||||
"strict": false,
|
||||
"esModuleInterop": true,
|
||||
"outDir": "dist",
|
||||
"rootDir": "src",
|
||||
"skipLibCheck": true,
|
||||
"forceConsistentCasingInFileNames": true,
|
||||
"declaration": true,
|
||||
"skipLibCheck": true
|
||||
"declarationMap": true,
|
||||
"sourceMap": true
|
||||
},
|
||||
"include": ["src/**/*.ts"]
|
||||
"include": ["src/**/*"],
|
||||
"exclude": ["node_modules", "dist"]
|
||||
}
|
||||
|
||||
@@ -1,10 +1,19 @@
|
||||
"""
|
||||
Zulip platform adapter (Hermes plugin) — Gen 3.
|
||||
Zulip platform adapter (Hermes plugin) — Gen 4.
|
||||
|
||||
Connects a Hermes agent to Zulip via event queue polling. DMs and @mentions
|
||||
are routed to the agent via the Hermes Gateway's handle_message() interface.
|
||||
Replies use a placeholder->edit streaming pattern for UX feedback.
|
||||
|
||||
Gen 4 Improvements (2026-06-27):
|
||||
1. HTTP client recreation on reconnect — new httpx.AsyncClient() after queue expiry
|
||||
or sustained failures to prevent CLOSE-WAIT socket leaks
|
||||
2. Heartbeat logging — periodic "still alive" message to prove poll loop is running
|
||||
3. Log sanitization — truncates 502/error HTML to first 80 chars (no Netbird HTML spam)
|
||||
4. Connection recovery — creates fresh HTTP client after N empty poll cycles
|
||||
5. Poll silence detection — logs warnings when no events received for extended period
|
||||
6. Connection reuse limit — periodic client recreation to prevent connection pool exhaustion
|
||||
|
||||
Gen 3 Improvements (2026-06-26):
|
||||
1. Malformed message resilience — try/except wraps each event, no crash on bad JSON
|
||||
2. Periodic @all-bots refresh — background task re-resolves user ID every hour
|
||||
@@ -66,6 +75,13 @@ DEDUP_CLEANUP_INTERVAL = 120 # Clean old entries every 2 minutes
|
||||
ALL_BOTS_REFRESH_INTERVAL = 3600 # Re-resolve @all-bots user ID every hour
|
||||
SELFTEST_TIMEOUT = 30 # seconds to wait for self-test response
|
||||
|
||||
# Gen 4: Connection pool health
|
||||
HEARTBEAT_INTERVAL = 300 # Log heartbeat every 5 minutes
|
||||
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
|
||||
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
|
||||
|
||||
# Regex to strip Zulip @mention markup
|
||||
MENTION_CLEANER = re.compile(r"@\*\*[^*]+\*\*")
|
||||
|
||||
@@ -158,7 +174,18 @@ class ZulipAdapter(BasePlatformAdapter):
|
||||
or os.getenv("ZULIP_POLL_INTERVAL", str(DEFAULT_POLL_INTERVAL))
|
||||
)
|
||||
|
||||
# --- Gen 4: Connection health ---
|
||||
self._consecutive_empty_polls: int = 0
|
||||
self._total_poll_count: int = 0
|
||||
self._last_event_received_time: float = time.time()
|
||||
self._last_heartbeat_time: float = 0.0
|
||||
self._client_pool_reset_count: int = 0
|
||||
|
||||
# --- State ---
|
||||
# Support both ZULIP_SITE and ZULIP_URL env vars
|
||||
if not self._site:
|
||||
self._site = (os.getenv("ZULIP_URL", "") or "").rstrip("/")
|
||||
|
||||
self._auth_header: str = _build_auth_header(self._email, self._api_key)
|
||||
self._queue_id: Optional[str] = None
|
||||
self._last_event_id: int = -1
|
||||
@@ -201,7 +228,7 @@ class ZulipAdapter(BasePlatformAdapter):
|
||||
|
||||
# --- Gen 3: Health stats callback ---
|
||||
self._health_callback = None
|
||||
self._health_callback_interval: int = 600 # 10 min
|
||||
self._health_callback_interval: int = 300 # 5 min
|
||||
|
||||
@staticmethod
|
||||
def _resolve_int(val: Any, default: int) -> int:
|
||||
@@ -227,8 +254,11 @@ class ZulipAdapter(BasePlatformAdapter):
|
||||
# Connection lifecycle
|
||||
# ------------------------------------------------------------------
|
||||
|
||||
async def connect(self) -> bool:
|
||||
"""Connect to Zulip and register an event queue."""
|
||||
async def connect(self, **kwargs: Any) -> bool:
|
||||
"""Connect to Zulip and register an event queue.
|
||||
|
||||
Accepts **kwargs for Hermes Gateway compatibility (e.g., is_reconnect).
|
||||
"""
|
||||
if not HTTPX_AVAILABLE:
|
||||
logger.warning("[%s] httpx not installed", self.name)
|
||||
return False
|
||||
@@ -238,11 +268,16 @@ class ZulipAdapter(BasePlatformAdapter):
|
||||
return False
|
||||
|
||||
try:
|
||||
self._http_client = httpx.AsyncClient(timeout=30.0)
|
||||
await self._create_http_client()
|
||||
self._health_stats["started_at"] = datetime.now(timezone.utc).isoformat()
|
||||
|
||||
# Reset connection health counters
|
||||
self._consecutive_empty_polls = 0
|
||||
self._total_poll_count = 0
|
||||
self._last_event_received_time = time.time()
|
||||
|
||||
# Register event queue
|
||||
queue_resp = await self._api_call(
|
||||
queue_resp, _ = await self._api_call(
|
||||
"POST", "/api/v1/register",
|
||||
data={
|
||||
"event_types": '["message"]',
|
||||
@@ -259,6 +294,11 @@ class ZulipAdapter(BasePlatformAdapter):
|
||||
self._last_event_id = data.get("last_event_id", -1)
|
||||
self._bot_user_id = data.get("user_id")
|
||||
|
||||
# Resolve bot user_id if register endpoint didn't provide it
|
||||
# (common on some Zulip versions)
|
||||
if self._bot_user_id is None:
|
||||
await self._resolve_bot_user_id()
|
||||
|
||||
# Gen 2: Try to resolve @all-bots user ID dynamically
|
||||
await self._resolve_all_bots_user_id()
|
||||
|
||||
@@ -337,24 +377,94 @@ class ZulipAdapter(BasePlatformAdapter):
|
||||
try:
|
||||
events = await self._fetch_events()
|
||||
self._health_stats["poll_count"] += 1
|
||||
self._total_poll_count += 1
|
||||
self._health_stats["last_poll_at"] = datetime.now(
|
||||
timezone.utc
|
||||
).isoformat()
|
||||
for event in events:
|
||||
try:
|
||||
await self._process_zulip_event(event)
|
||||
except Exception as e:
|
||||
# Gen 3: Malformed message resilience — catch per-event
|
||||
# failures so one bad message never kills the poll loop
|
||||
|
||||
if events:
|
||||
now = time.time()
|
||||
self._consecutive_empty_polls = 0
|
||||
self._last_event_received_time = now
|
||||
for event in events:
|
||||
try:
|
||||
await self._process_zulip_event(event)
|
||||
except Exception as e:
|
||||
# Gen 3: Malformed message resilience — catch per-event
|
||||
# failures so one bad message never kills the poll loop
|
||||
logger.warning(
|
||||
"[%s] Malformed event skipped: %s. "
|
||||
"Event id=%s",
|
||||
self.name, e,
|
||||
event.get("id", "unknown"),
|
||||
)
|
||||
self._health_stats["malformed_events"] = (
|
||||
self._health_stats.get("malformed_events", 0) + 1
|
||||
)
|
||||
else:
|
||||
# Gen 4: Track consecutive empty polls to detect silent failures
|
||||
self._consecutive_empty_polls += 1
|
||||
|
||||
# Gen 4: Detect prolonged silence (no events for too long)
|
||||
silence_duration = time.time() - self._last_event_received_time
|
||||
if silence_duration > POLL_SILENCE_WARN_INTERVAL:
|
||||
logger.warning(
|
||||
"[%s] Malformed event skipped: %s. "
|
||||
"Event id=%s",
|
||||
self.name, e,
|
||||
event.get("id", "unknown"),
|
||||
"[%s] No events received for %.0fs "
|
||||
"(consecutive_empty=%d, total_polls=%d, "
|
||||
"queue=%s)",
|
||||
self.name, silence_duration,
|
||||
self._consecutive_empty_polls,
|
||||
self._total_poll_count,
|
||||
self._queue_id,
|
||||
)
|
||||
self._health_stats["malformed_events"] = (
|
||||
self._health_stats.get("malformed_events", 0) + 1
|
||||
|
||||
# Gen 4: Reset HTTP client if too many empty polls in a row
|
||||
# (connection pool likely stuck in CLOSE-WAIT)
|
||||
if (
|
||||
self._consecutive_empty_polls
|
||||
>= MAX_CONSECUTIVE_EMPTY_POLLS
|
||||
):
|
||||
logger.warning(
|
||||
"[%s] %d consecutive empty polls — "
|
||||
"recreating HTTP client and reconnecting",
|
||||
self.name, self._consecutive_empty_polls,
|
||||
)
|
||||
await self._reconnect_with_fresh_client()
|
||||
backoff_idx = 0
|
||||
self._consecutive_empty_polls = 0
|
||||
continue
|
||||
|
||||
# Gen 4: Periodic client refresh to prevent connection leaks
|
||||
if self._total_poll_count % CLIENT_REUSE_LIMIT == 0:
|
||||
logger.info(
|
||||
"[%s] Client reuse limit reached (%d polls) — "
|
||||
"creating fresh HTTP client",
|
||||
self.name, self._total_poll_count,
|
||||
)
|
||||
old_client = self._http_client
|
||||
await self._create_http_client()
|
||||
if old_client:
|
||||
try:
|
||||
await old_client.aclose()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
# Gen 4: Heartbeat logging — prove poll loop is alive
|
||||
now = time.time()
|
||||
if now - self._last_heartbeat_time > HEARTBEAT_INTERVAL:
|
||||
self._last_heartbeat_time = now
|
||||
silence = now - self._last_event_received_time
|
||||
logger.info(
|
||||
"[%s] Heartbeat — polls=%d empty=%d "
|
||||
"silence=%.0fs errors=%d reconnects=%d "
|
||||
"queue=%s",
|
||||
self.name, self._total_poll_count,
|
||||
self._consecutive_empty_polls, silence,
|
||||
self._health_stats["poll_errors"],
|
||||
self._health_stats["reconnects"],
|
||||
self._queue_id,
|
||||
)
|
||||
|
||||
backoff_idx = 0
|
||||
except asyncio.CancelledError:
|
||||
break
|
||||
@@ -370,7 +480,7 @@ class ZulipAdapter(BasePlatformAdapter):
|
||||
err_str = str(e)
|
||||
if "BAD_EVENT_QUEUE_ID" in err_str or "queue_id" in err_str.lower():
|
||||
logger.info("[%s] Queue expired, re-registering...", self.name)
|
||||
await self._reconnect()
|
||||
await self._reconnect_with_fresh_client()
|
||||
backoff_idx = 0
|
||||
continue
|
||||
logger.warning("[%s] Poll error: %s", self.name, e)
|
||||
@@ -381,11 +491,15 @@ class ZulipAdapter(BasePlatformAdapter):
|
||||
await asyncio.sleep(self._poll_interval)
|
||||
|
||||
async def _fetch_events(self) -> List[Dict[str, Any]]:
|
||||
"""Fetch events from the Zulip event queue."""
|
||||
"""Fetch events from the Zulip event queue.
|
||||
|
||||
Raises RuntimeError with BAD_EVENT_QUEUE_ID when the queue
|
||||
expires, so _poll_forever can reconnect.
|
||||
"""
|
||||
if not self._queue_id:
|
||||
return []
|
||||
|
||||
resp = await self._api_call(
|
||||
resp, raw_status = await self._api_call(
|
||||
"GET", "/api/v1/events",
|
||||
params={
|
||||
"queue_id": self._queue_id,
|
||||
@@ -393,6 +507,9 @@ class ZulipAdapter(BasePlatformAdapter):
|
||||
"dont_block": "true",
|
||||
},
|
||||
)
|
||||
# Detect queue expiry from HTTP response
|
||||
if raw_status == 400:
|
||||
raise RuntimeError("BAD_EVENT_QUEUE_ID: queue expired")
|
||||
if not resp:
|
||||
return []
|
||||
|
||||
@@ -405,11 +522,33 @@ class ZulipAdapter(BasePlatformAdapter):
|
||||
|
||||
return [e for e in events if e.get("type") == "message"]
|
||||
|
||||
async def _create_http_client(self) -> None:
|
||||
"""Create a fresh httpx client, closing the old one if it exists.
|
||||
|
||||
Gen 4: This ensures a clean connection pool after sustained failures
|
||||
or queue expiry, preventing CLOSE-WAIT socket leaks.
|
||||
"""
|
||||
old_client = self._http_client
|
||||
self._http_client = httpx.AsyncClient(
|
||||
timeout=httpx.Timeout(30.0, connect=15.0),
|
||||
limits=httpx.Limits(
|
||||
max_keepalive_connections=5,
|
||||
max_connections=10,
|
||||
keepalive_expiry=60.0,
|
||||
),
|
||||
)
|
||||
self._client_pool_reset_count += 1
|
||||
if old_client:
|
||||
try:
|
||||
await old_client.aclose()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
async def _reconnect(self) -> None:
|
||||
"""Re-register the event queue."""
|
||||
self._queue_id = None
|
||||
try:
|
||||
resp = await self._api_call(
|
||||
resp, _ = await self._api_call(
|
||||
"POST", "/api/v1/register",
|
||||
data={
|
||||
"event_types": '["message"]',
|
||||
@@ -422,9 +561,21 @@ class ZulipAdapter(BasePlatformAdapter):
|
||||
self._last_event_id = resp.get("last_event_id", -1)
|
||||
self._health_stats["reconnects"] += 1
|
||||
logger.info("[%s] Reconnected, queue=%s", self.name, self._queue_id)
|
||||
self._consecutive_empty_polls = 0
|
||||
self._last_event_received_time = time.time()
|
||||
except Exception as e:
|
||||
logger.error("[%s] Reconnect failed: %s", self.name, e)
|
||||
|
||||
async def _reconnect_with_fresh_client(self) -> None:
|
||||
"""Re-register the event queue with a fresh HTTP client.
|
||||
|
||||
Gen 4: Creates a new httpx client to break out of stuck connection
|
||||
pools (CLOSE-WAIT sockets from previous failures).
|
||||
"""
|
||||
logger.info("[%s] Reconnecting with fresh HTTP client", self.name)
|
||||
await self._create_http_client()
|
||||
await self._reconnect()
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# Gen 2: Background dedup maintenance
|
||||
# ------------------------------------------------------------------
|
||||
@@ -496,6 +647,32 @@ class ZulipAdapter(BasePlatformAdapter):
|
||||
"[%s] @all-bots refresh error (non-fatal)", self.name
|
||||
)
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# Gen 5: Bot user ID resolution
|
||||
# ------------------------------------------------------------------
|
||||
|
||||
async def _resolve_bot_user_id(self) -> None:
|
||||
"""Fetch the bot's own user ID from /api/v1/users/me.
|
||||
|
||||
The /register endpoint does not always return user_id on
|
||||
all Zulip server versions. This ensures @mention detection
|
||||
works by resolving it from the identity endpoint.
|
||||
"""
|
||||
try:
|
||||
resp, _ = await self._api_call("GET", "/api/v1/users/me")
|
||||
if resp:
|
||||
uid = resp.get("user_id")
|
||||
if uid is not None:
|
||||
self._bot_user_id = int(uid)
|
||||
logger.info(
|
||||
"[%s] Resolved bot user_id=%s from /users/me",
|
||||
self.name, uid,
|
||||
)
|
||||
except Exception as e:
|
||||
logger.debug(
|
||||
"[%s] Could not resolve bot user_id: %s", self.name, e
|
||||
)
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# Gen 2: Dynamic @all-bots resolution
|
||||
# ------------------------------------------------------------------
|
||||
@@ -506,7 +683,7 @@ class ZulipAdapter(BasePlatformAdapter):
|
||||
Falls back to the configured or default value if resolution fails.
|
||||
"""
|
||||
try:
|
||||
resp = await self._api_call("GET", "/api/v1/users")
|
||||
resp, _ = await self._api_call("GET", "/api/v1/users")
|
||||
if not resp:
|
||||
return
|
||||
members = resp.get("members", [])
|
||||
@@ -679,12 +856,15 @@ class ZulipAdapter(BasePlatformAdapter):
|
||||
content: str,
|
||||
reply_to: Optional[str] = None,
|
||||
metadata: Optional[Dict[str, Any]] = None,
|
||||
**kwargs: Any,
|
||||
) -> SendResult:
|
||||
"""Send a message to a Zulip chat.
|
||||
|
||||
For DMs, chat_id is the user_id. For streams, format is
|
||||
"stream_name:topic". Sends a placeholder first if this is the
|
||||
first message in a turn, then edits it on subsequent calls.
|
||||
|
||||
Accepts **kwargs for Hermes Gateway compatibility.
|
||||
"""
|
||||
# Determine target type and IDs
|
||||
is_dm = not chat_id or ":" not in chat_id
|
||||
@@ -727,28 +907,22 @@ class ZulipAdapter(BasePlatformAdapter):
|
||||
self._health_stats["send_count"] += 1
|
||||
return SendResult(
|
||||
success=True,
|
||||
platform="zulip",
|
||||
chat_id=chat_id,
|
||||
message_id=msg_id,
|
||||
metadata={"placeholder": True},
|
||||
)
|
||||
|
||||
# Send actual content
|
||||
result = await self._send_api_call(payload)
|
||||
if result and result.get("id"):
|
||||
msg_id = str(result["id"])
|
||||
self._health_stats["send_count"] += 1
|
||||
return SendResult(
|
||||
success=True,
|
||||
platform="zulip",
|
||||
chat_id=chat_id,
|
||||
message_id=str(result["id"]),
|
||||
message_id=msg_id,
|
||||
)
|
||||
|
||||
self._health_stats["send_errors"] += 1
|
||||
return SendResult(
|
||||
success=False,
|
||||
platform="zulip",
|
||||
chat_id=chat_id,
|
||||
error="Failed to send message",
|
||||
)
|
||||
|
||||
@@ -758,19 +932,43 @@ class ZulipAdapter(BasePlatformAdapter):
|
||||
message_id: str,
|
||||
content: str,
|
||||
metadata: Optional[Dict[str, Any]] = None,
|
||||
) -> bool:
|
||||
"""Edit a previously sent Zulip message (for placeholder replacement)."""
|
||||
**kwargs: Any,
|
||||
) -> SendResult:
|
||||
"""Edit a previously sent Zulip message (for placeholder replacement).
|
||||
|
||||
Returns SendResult for Hermes Gateway streaming interface compatibility.
|
||||
Accepts **kwargs (e.g., finalize=True) — extra kwargs are silently ignored.
|
||||
|
||||
Note: Zulip uses POST /api/v1/messages/{message_id} for editing.
|
||||
"""
|
||||
truncated = _truncate(content)
|
||||
payload = {
|
||||
"message_id": int(message_id),
|
||||
"content": truncated,
|
||||
}
|
||||
try:
|
||||
resp = await self._api_call("PATCH", "/api/v1/messages", data=payload)
|
||||
return resp is not None
|
||||
# Zulip API: PATCH /api/v1/messages/{message_id} with content in body
|
||||
resp, status = await self._api_call(
|
||||
"PATCH", f"/api/v1/messages/{int(message_id)}",
|
||||
data=payload,
|
||||
)
|
||||
if resp:
|
||||
self._health_stats["send_count"] += 1
|
||||
return SendResult(
|
||||
success=True,
|
||||
message_id=str(resp.get("id", message_id)),
|
||||
)
|
||||
self._health_stats["send_errors"] += 1
|
||||
return SendResult(
|
||||
success=False,
|
||||
error=f"Edit failed (status={status})",
|
||||
)
|
||||
except Exception as e:
|
||||
logger.warning("[%s] Edit message %s error: %s", self.name, message_id, e)
|
||||
return False
|
||||
self._health_stats["send_errors"] += 1
|
||||
return SendResult(
|
||||
success=False,
|
||||
error=str(e),
|
||||
)
|
||||
|
||||
async def send_typing(self, chat_id: str, metadata: Optional[Dict] = None) -> None:
|
||||
"""Send typing indicator to a Zulip chat."""
|
||||
@@ -782,6 +980,7 @@ class ZulipAdapter(BasePlatformAdapter):
|
||||
"POST", "/api/v1/typing",
|
||||
data={"to": user_ids, "op": "start"},
|
||||
)
|
||||
# typing indicator failure is non-critical
|
||||
except Exception:
|
||||
pass # Non-critical
|
||||
|
||||
@@ -795,6 +994,7 @@ class ZulipAdapter(BasePlatformAdapter):
|
||||
"POST", "/api/v1/typing",
|
||||
data={"to": user_ids, "op": "stop"},
|
||||
)
|
||||
# typing indicator failure is non-critical
|
||||
except Exception:
|
||||
pass # Non-critical
|
||||
|
||||
@@ -806,12 +1006,14 @@ class ZulipAdapter(BasePlatformAdapter):
|
||||
return {"name": chat_id, "type": "dm"}
|
||||
|
||||
async def delete_message(
|
||||
self, chat_id: str, message_id: str, metadata: Optional[Dict] = None
|
||||
) -> bool:
|
||||
self, chat_id: str, message_id: str,
|
||||
metadata: Optional[Dict] = None,
|
||||
**kwargs: Any,
|
||||
) -> SendResult:
|
||||
"""Delete a message via Zulip API."""
|
||||
# Zulip doesn't support message deletion via API for bots
|
||||
# Send an empty edit instead
|
||||
return await self.edit_message(chat_id, message_id, "*deleted*")
|
||||
return await self.edit_message(chat_id, message_id, "*deleted*", **kwargs)
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# Gen 2: Self-test diagnostics
|
||||
@@ -990,6 +1192,10 @@ class ZulipAdapter(BasePlatformAdapter):
|
||||
stats["all_bots_user_id"] = self._all_bots_user_id
|
||||
stats["connected"] = self._connected
|
||||
stats["checked_at"] = now.isoformat()
|
||||
stats["total_polls"] = self._total_poll_count
|
||||
stats["consecutive_empty_polls"] = self._consecutive_empty_polls
|
||||
stats["client_pool_resets"] = self._client_pool_reset_count
|
||||
stats["silence_seconds"] = round(time.time() - self._last_event_received_time)
|
||||
|
||||
# Compute error rate
|
||||
total_polls = stats.get("poll_count", 0) or 1
|
||||
@@ -1009,10 +1215,15 @@ class ZulipAdapter(BasePlatformAdapter):
|
||||
path: str,
|
||||
data: Optional[Dict] = None,
|
||||
params: Optional[Dict] = None,
|
||||
) -> Optional[Dict]:
|
||||
"""Make an API call to Zulip."""
|
||||
) -> Tuple[Optional[Dict], int]:
|
||||
"""Make an API call to Zulip.
|
||||
|
||||
Returns (response_json, status_code). status_code is useful for
|
||||
callers that need to detect specific HTTP errors like 400
|
||||
(BAD_EVENT_QUEUE_ID). Returns (None, 0) on transport errors.
|
||||
"""
|
||||
if not self._http_client:
|
||||
return None
|
||||
return None, 0
|
||||
|
||||
url = f"{self._site}{path}"
|
||||
headers = {
|
||||
@@ -1033,28 +1244,31 @@ class ZulipAdapter(BasePlatformAdapter):
|
||||
url, data=data, headers=headers,
|
||||
)
|
||||
else:
|
||||
return None
|
||||
return None, 0
|
||||
|
||||
if response.status_code >= 400:
|
||||
# Gen 4: Truncate error bodies to prevent HTML spam
|
||||
body = response.text[:MAX_ERROR_LOG_LEN].replace("\n", " ")
|
||||
logger.warning(
|
||||
"[%s] API %s %s: %d %s",
|
||||
self.name, method, path,
|
||||
response.status_code, response.text[:200],
|
||||
response.status_code, body,
|
||||
)
|
||||
return None
|
||||
return None, response.status_code
|
||||
|
||||
return response.json()
|
||||
return response.json(), response.status_code
|
||||
|
||||
except httpx.TimeoutException:
|
||||
logger.debug("[%s] Timeout on %s %s", self.name, method, path)
|
||||
return None
|
||||
return None, 0
|
||||
except Exception as e:
|
||||
logger.warning("[%s] API error %s %s: %s", self.name, method, path, e)
|
||||
return None
|
||||
return None, 0
|
||||
|
||||
async def _send_api_call(self, payload: Dict) -> Optional[Dict]:
|
||||
"""Send a message to Zulip."""
|
||||
return await self._api_call("POST", "/api/v1/messages", data=payload)
|
||||
result, _ = await self._api_call("POST", "/api/v1/messages", data=payload)
|
||||
return result
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# Utilities
|
||||
|
||||
Reference in New Issue
Block a user