Compare commits
20
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
0dd3c5868d | ||
|
|
2d7acdb94a | ||
|
|
8b588f027d | ||
|
|
20271f0eb3 | ||
|
|
f205230e84 | ||
|
|
c8accb7135 | ||
|
|
cb83baf299 | ||
|
|
e7b658e3f7 | ||
|
|
2f6c4283d2 | ||
|
|
033a5c3042 | ||
|
|
08d0b7371b | ||
|
|
28dabe8834 | ||
|
|
92f281c43a | ||
|
|
7474ab1dc4 | ||
|
|
7047749ef0 | ||
|
|
553d469174 | ||
|
|
080d2ad756 | ||
|
|
9611f935bc | ||
|
|
84b80179ec | ||
|
|
5300c0f998 |
+28
-29
@@ -14,43 +14,42 @@ jobs:
|
|||||||
validate:
|
validate:
|
||||||
runs-on: ubuntu-latest
|
runs-on: ubuntu-latest
|
||||||
steps:
|
steps:
|
||||||
- uses: actions/checkout@v4
|
- run: python3 --version
|
||||||
- name: Ensure python3 + PyYAML (act container is alpine-based, no python preinstalled)
|
|
||||||
run: |
|
|
||||||
. /etc/os-release
|
|
||||||
echo "job OS: $ID"
|
|
||||||
if command -v python3 >/dev/null 2>&1; then
|
|
||||||
python3 --version
|
|
||||||
elif command -v apk >/dev/null 2>&1; then
|
|
||||||
echo "installing via apk"
|
|
||||||
apk add --no-cache python3 py3-yaml
|
|
||||||
python3 --version
|
|
||||||
python3 -c "import yaml" || { echo "yaml import failed after py3-yaml; trying pip"; python3 -m pip install --break-system-packages pyyaml; }
|
|
||||||
elif command -v apt-get >/dev/null 2>&1; then
|
|
||||||
echo "installing via apt"
|
|
||||||
apt-get update -qq
|
|
||||||
apt-get install -y -qq python3 python3-yaml
|
|
||||||
python3 --version
|
|
||||||
else
|
|
||||||
echo "no python3 and no apk/apt in job image — cannot run CI checks"
|
|
||||||
exit 1
|
|
||||||
fi
|
|
||||||
- run: node --version
|
- run: node --version
|
||||||
|
- run: echo "Runner works!"
|
||||||
|
- run: git clone --depth 1 http://abiba-bot:mipjoq-tybbox-2ciHru@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
||||||
|
|
||||||
- name: Python syntax check
|
- name: Python syntax check
|
||||||
run: |
|
run: |
|
||||||
python3 -m py_compile plugins/platforms/zulip/adapter.py && echo "✅ adapter.py OK"
|
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 && echo "✅ a2a adapter OK"
|
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: Workflow YAML parse check
|
|
||||||
run: python3 ci_check.py workflows
|
|
||||||
|
|
||||||
- name: Config validation
|
- name: Config validation
|
||||||
run: python3 ci_check.py config
|
run: |
|
||||||
|
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: No secrets check
|
- name: No secrets check
|
||||||
run: python3 ci_check.py secrets
|
run: |
|
||||||
|
! grep -r "api_key.*[A-Za-z0-9]\{20,\}" --include="*.py" --include="*.ts" --include="*.yaml" . 2>/dev/null || echo "⚠️ Possible API key found!"
|
||||||
|
|
||||||
- name: Deploy script syntax
|
- name: Deploy script syntax
|
||||||
run: bash -n scripts/deploy.sh && echo "✅ deploy.sh syntax OK"
|
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:mipjoq-tybbox-2ciHru@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"
|
||||||
|
|||||||
@@ -21,7 +21,7 @@ jobs:
|
|||||||
runs-on: ubuntu-latest
|
runs-on: ubuntu-latest
|
||||||
steps:
|
steps:
|
||||||
- name: Checkout
|
- name: Checkout
|
||||||
uses: actions/checkout@v4
|
run: git clone --depth 1 http://abiba-bot:mipjoq-tybbox-2ciHru@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
||||||
- name: Python syntax
|
- name: Python syntax
|
||||||
run: |
|
run: |
|
||||||
python3 -m py_compile plugins/platforms/zulip/adapter.py 2>/dev/null
|
python3 -m py_compile plugins/platforms/zulip/adapter.py 2>/dev/null
|
||||||
@@ -37,7 +37,7 @@ jobs:
|
|||||||
environment: canary
|
environment: canary
|
||||||
steps:
|
steps:
|
||||||
- name: Checkout
|
- name: Checkout
|
||||||
uses: actions/checkout@v4
|
run: git clone --depth 1 http://abiba-bot:mipjoq-tybbox-2ciHru@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
||||||
- name: Deploy to Tanko
|
- name: Deploy to Tanko
|
||||||
env:
|
env:
|
||||||
TAG: ${{ github.ref_name }}
|
TAG: ${{ github.ref_name }}
|
||||||
@@ -66,7 +66,7 @@ jobs:
|
|||||||
environment: production
|
environment: production
|
||||||
steps:
|
steps:
|
||||||
- name: Checkout
|
- name: Checkout
|
||||||
uses: actions/checkout@v4
|
run: git clone --depth 1 http://abiba-bot:mipjoq-tybbox-2ciHru@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
||||||
- name: Deploy to all Hermes agents
|
- name: Deploy to all Hermes agents
|
||||||
env:
|
env:
|
||||||
TAG: ${{ github.ref_name }}
|
TAG: ${{ github.ref_name }}
|
||||||
@@ -97,7 +97,7 @@ jobs:
|
|||||||
environment: agent-zero
|
environment: agent-zero
|
||||||
steps:
|
steps:
|
||||||
- name: Checkout
|
- name: Checkout
|
||||||
uses: actions/checkout@v4
|
run: git clone --depth 1 http://abiba-bot:mipjoq-tybbox-2ciHru@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
||||||
- name: Deploy to kagentz
|
- name: Deploy to kagentz
|
||||||
env:
|
env:
|
||||||
TAG: ${{ github.ref_name }}
|
TAG: ${{ github.ref_name }}
|
||||||
|
|||||||
+1
-26
@@ -1,27 +1,2 @@
|
|||||||
# CI Pipeline Status
|
# CI Pipeline Status
|
||||||
|
Status: Active
|
||||||
**Last verified: 2026-09-25** — see PR "main security + truncate + CI repair".
|
|
||||||
|
|
||||||
## Reality check (before that PR)
|
|
||||||
`ci.yml` was **invalid YAML** — the inline `run: |` blocks for the config and
|
|
||||||
secrets checks dedented out of their block scalar, so Gitea Actions could never
|
|
||||||
parse the workflow. CI never ran on any PR, and this file's "Active" status was
|
|
||||||
fiction. The old "No secrets check" also swallowed hits with `|| echo` (always
|
|
||||||
green) and never scanned `.yml` files — which is where six embedded
|
|
||||||
credentials were living.
|
|
||||||
|
|
||||||
## Current pipeline
|
|
||||||
| Trigger | Job | Checks |
|
|
||||||
|---|---|---|
|
|
||||||
| PR → main, push main, tag `v*` | `validate` | adapter + a2a `py_compile`; workflow YAML parse; `config.yaml.example`; secret scan; `bash -n scripts/deploy.sh` |
|
|
||||||
| tag `v*` | `deploy` | Tanko canary + Mumuni via `scripts/deploy.sh` (native mode) |
|
|
||||||
|
|
||||||
All validation logic lives in **`ci_check.py`** — run `python3 ci_check.py all`
|
|
||||||
locally before pushing; it fails the check on real hits instead of echoing
|
|
||||||
warnings.
|
|
||||||
|
|
||||||
## Known caveats
|
|
||||||
- The removed credentials still exist in git history (pre-July commits) —
|
|
||||||
rotating `abiba-bot`'s password is a **server-side** task, out of repo scope.
|
|
||||||
- Runner availability (`zulip-runner` on CT 116) is not verifiable from agents
|
|
||||||
— the first green run after merge is the proof.
|
|
||||||
|
|||||||
-105
@@ -1,105 +0,0 @@
|
|||||||
#!/usr/bin/env python3
|
|
||||||
"""CI validation checks for zulip-platform-plugins.
|
|
||||||
|
|
||||||
The inline validation steps this script replaced were never valid YAML in the
|
|
||||||
first place (the `run: |` blocks dedented out of their block scalar), so the
|
|
||||||
whole `validate` job silently never ran. Keeping the logic in a real file
|
|
||||||
means it is testable locally — run `python3 ci_check.py all` before pushing.
|
|
||||||
|
|
||||||
Usage: python3 ci_check.py {workflows|config|secrets|all}
|
|
||||||
"""
|
|
||||||
import re
|
|
||||||
import sys
|
|
||||||
from pathlib import Path
|
|
||||||
|
|
||||||
ROOT = Path(__file__).resolve().parent
|
|
||||||
CODE_GLOBS = ("*.py", "*.ts", "*.yaml", "*.yml", "*.cjs", "*.sh")
|
|
||||||
# scheme://user:pass@host — embedded HTTP Basic credentials
|
|
||||||
CRED_RE = re.compile(r"://[^/ \"']+:[^/ \"'@]+@")
|
|
||||||
API_KEY_RE = re.compile(r"""api_key["']?\s*[:=]\s*["']?[A-Za-z0-9]{20,}""")
|
|
||||||
# Allowlist: placeholder patterns, not real secrets
|
|
||||||
PLACEHOLDER = ("${", "{{", "user:pass", "USER:TOKEN", "user:token", "example")
|
|
||||||
|
|
||||||
|
|
||||||
def iter_code_files():
|
|
||||||
for path in sorted(ROOT.rglob("*")):
|
|
||||||
if not path.is_file():
|
|
||||||
continue
|
|
||||||
if any(part in {".git", "__pycache__", "node_modules"} for part in path.parts):
|
|
||||||
continue
|
|
||||||
if path.name == Path(__file__).name: # this file documents the patterns
|
|
||||||
continue
|
|
||||||
if any(path.match(glob) for glob in CODE_GLOBS):
|
|
||||||
yield path
|
|
||||||
|
|
||||||
|
|
||||||
def check_workflows() -> bool:
|
|
||||||
try:
|
|
||||||
import yaml
|
|
||||||
except ModuleNotFoundError:
|
|
||||||
print("⚠️ pyyaml not installed — workflow parse check skipped")
|
|
||||||
return True
|
|
||||||
ok = True
|
|
||||||
for f in sorted((ROOT / ".gitea" / "workflows").glob("*.yml")):
|
|
||||||
try:
|
|
||||||
doc = yaml.safe_load(f.read_text())
|
|
||||||
assert isinstance(doc, dict) and "jobs" in doc, "missing 'jobs' mapping"
|
|
||||||
print(f"✅ {f.relative_to(ROOT)} valid — jobs: {list(doc['jobs'])}")
|
|
||||||
except Exception as e:
|
|
||||||
print(f"❌ {f.relative_to(ROOT)}: {e}")
|
|
||||||
ok = False
|
|
||||||
return ok
|
|
||||||
|
|
||||||
|
|
||||||
def check_config() -> bool:
|
|
||||||
try:
|
|
||||||
import yaml
|
|
||||||
except ModuleNotFoundError:
|
|
||||||
print("⚠️ pyyaml not installed — config check skipped")
|
|
||||||
return True
|
|
||||||
try:
|
|
||||||
cfg = yaml.safe_load((ROOT / "config.yaml.example").read_text())
|
|
||||||
assert isinstance(cfg, dict) and "zulip" in cfg, "missing 'zulip' section"
|
|
||||||
print("✅ config.yaml.example valid")
|
|
||||||
return True
|
|
||||||
except Exception as e:
|
|
||||||
print(f"❌ config.yaml.example: {e}")
|
|
||||||
return False
|
|
||||||
|
|
||||||
|
|
||||||
def check_secrets() -> bool:
|
|
||||||
hits = []
|
|
||||||
for path in iter_code_files():
|
|
||||||
try:
|
|
||||||
text = path.read_text(errors="ignore")
|
|
||||||
except OSError:
|
|
||||||
continue
|
|
||||||
for lineno, line in enumerate(text.splitlines(), 1):
|
|
||||||
for rx in (CRED_RE, API_KEY_RE):
|
|
||||||
if rx.search(line) and not any(p in line for p in PLACEHOLDER):
|
|
||||||
hits.append(f"{path.relative_to(ROOT)}:{lineno}: {line.strip()[:100]}")
|
|
||||||
break
|
|
||||||
if hits:
|
|
||||||
print("❌ Embedded credentials detected:")
|
|
||||||
for h in hits:
|
|
||||||
print(f" {h}")
|
|
||||||
return False
|
|
||||||
print("✅ No embedded credentials in tracked code/workflow files")
|
|
||||||
return True
|
|
||||||
|
|
||||||
|
|
||||||
CHECKS = {"workflows": check_workflows, "config": check_config, "secrets": check_secrets}
|
|
||||||
|
|
||||||
|
|
||||||
def main() -> int:
|
|
||||||
args = sys.argv[1:] or ["all"]
|
|
||||||
names = list(CHECKS) if "all" in args else args
|
|
||||||
failed = [n for n in names if n not in CHECKS or not CHECKS[n]()]
|
|
||||||
if failed:
|
|
||||||
print(f"❌ CI checks failed: {', '.join(failed)}")
|
|
||||||
return 1
|
|
||||||
return 0
|
|
||||||
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
|
||||||
sys.exit(main())
|
|
||||||
@@ -203,22 +203,6 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
logger.error("Zulip POST %s network error: %s", path, exc)
|
logger.error("Zulip POST %s network error: %s", path, exc)
|
||||||
return {"result": "error", "msg": str(exc)}
|
return {"result": "error", "msg": str(exc)}
|
||||||
|
|
||||||
async def _api_patch(self, path: str, payload: Dict[str, Any]) -> Dict[str, Any]:
|
|
||||||
import aiohttp
|
|
||||||
url = f"{self._site}/api/v1/{path.lstrip('/')}"
|
|
||||||
try:
|
|
||||||
async with self._session.patch(
|
|
||||||
url, data=payload, auth=self._auth(),
|
|
||||||
timeout=aiohttp.ClientTimeout(total=15),
|
|
||||||
) as resp:
|
|
||||||
data = await resp.json()
|
|
||||||
if resp.status >= 400:
|
|
||||||
logger.debug("Zulip PATCH %s -> %s: %s", path, resp.status, str(data)[:200])
|
|
||||||
return data
|
|
||||||
except (aiohttp.ClientError, asyncio.TimeoutError) as exc:
|
|
||||||
logger.debug("Zulip PATCH %s network error: %s", path, exc)
|
|
||||||
return {"result": "error", "msg": str(exc)}
|
|
||||||
|
|
||||||
async def _api_delete(self, path: str, params: Optional[Dict[str, Any]] = None) -> Dict[str, Any]:
|
async def _api_delete(self, path: str, params: Optional[Dict[str, Any]] = None) -> Dict[str, Any]:
|
||||||
import aiohttp
|
import aiohttp
|
||||||
url = f"{self._site}/api/v1/{path.lstrip('/')}"
|
url = f"{self._site}/api/v1/{path.lstrip('/')}"
|
||||||
@@ -342,23 +326,6 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
|
|
||||||
return SendResult(success=True, message_id=str(last_id) if last_id else None)
|
return SendResult(success=True, message_id=str(last_id) if last_id else None)
|
||||||
|
|
||||||
async def edit_message(
|
|
||||||
self,
|
|
||||||
chat_id: str,
|
|
||||||
message_id: str,
|
|
||||||
content: str,
|
|
||||||
**kwargs,
|
|
||||||
) -> SendResult:
|
|
||||||
"""Edit a previously-sent message. Used by the gateway stream consumer
|
|
||||||
for progressive message updates during agent streaming."""
|
|
||||||
formatted = self.format_message(content)
|
|
||||||
data = await self._api_patch(f"messages/{message_id}", {"content": formatted})
|
|
||||||
if data.get("result") != "success":
|
|
||||||
msg = str(data.get("msg", "edit failed"))
|
|
||||||
logger.debug("Zulip edit_message(%s) -> %s", message_id, msg)
|
|
||||||
return SendResult(success=False, error=msg)
|
|
||||||
return SendResult(success=True, message_id=message_id)
|
|
||||||
|
|
||||||
async def get_chat_info(self, chat_id: str) -> Dict[str, Any]:
|
async def get_chat_info(self, chat_id: str) -> Dict[str, Any]:
|
||||||
kind, to, topic = _parse_target(chat_id, self._default_topic)
|
kind, to, topic = _parse_target(chat_id, self._default_topic)
|
||||||
if kind == "direct":
|
if kind == "direct":
|
||||||
|
|||||||
@@ -1,76 +1,26 @@
|
|||||||
/**
|
|
||||||
* PM2 Ecosystem Config — Zulip Gateway v3 (Resilience)
|
|
||||||
*
|
|
||||||
* Deploy: pm2 start /root/.pm2/ecosystem.config.cjs
|
|
||||||
* Status: pm2 status
|
|
||||||
* Logs: pm2 logs abiba-zulip
|
|
||||||
*/
|
|
||||||
|
|
||||||
module.exports = {
|
module.exports = {
|
||||||
apps: [
|
apps: [
|
||||||
{
|
{
|
||||||
// ── Router (main Zulip gateway) ──
|
|
||||||
name: "abiba-zulip",
|
name: "abiba-zulip",
|
||||||
script: "/bin/pi",
|
script: "/usr/bin/pi",
|
||||||
args: "--mode rpc --session-id zulip-service",
|
args: "--mode rpc --session-id zulip-service",
|
||||||
cwd: "/root",
|
cwd: "/root",
|
||||||
|
|
||||||
// Resilience hardening — up from pi defaults
|
|
||||||
max_restarts: 100, // Crash loops won't exhaust PM2 (was default 10)
|
|
||||||
min_uptime: "10s", // Must survive 10s to count as "alive"
|
|
||||||
max_memory_restart: "500M", // OOM protection — restart before swap thrash
|
|
||||||
restart_delay: 5000, // 5s cooldown between restarts
|
|
||||||
kill_timeout: 15000, // 15s SIGTERM grace before SIGKILL
|
|
||||||
listen_timeout: 30000, // 30s to bind health port
|
|
||||||
|
|
||||||
// Logging
|
|
||||||
log_date_format: "YYYY-MM-DD HH:mm:ss Z",
|
|
||||||
error_file: "/root/.pm2/logs/abiba-zulip-error.log",
|
|
||||||
out_file: "/root/.pm2/logs/abiba-zulip-out.log",
|
|
||||||
merge_logs: true,
|
|
||||||
log_type: "json",
|
|
||||||
|
|
||||||
// Process management
|
|
||||||
autorestart: true,
|
|
||||||
watch: false,
|
|
||||||
instances: 1,
|
|
||||||
exec_mode: "fork",
|
|
||||||
|
|
||||||
// Environment
|
|
||||||
env: {
|
env: {
|
||||||
ZULIP_ROLE: "router",
|
ZULIP_ROLE: "router",
|
||||||
ZULIP_SITE: "https://chat.sysloggh.net",
|
ZULIP_EXTENSION_ACTIVE: "true",
|
||||||
ZULIP_EMAIL: "abiba-bot@chat.sysloggh.net",
|
NODE_OPTIONS: "--max-old-space-size=512",
|
||||||
ZULIP_API_KEY: "cKTDMZAPW08dk3zl05sStzO7HRztzyn8",
|
|
||||||
AGENT_NAME: "abiba",
|
|
||||||
AGENT_OWNER_EMAIL: "jerome@sysloggh.com",
|
|
||||||
NODE_ENV: "production",
|
|
||||||
},
|
},
|
||||||
|
// Prevent rapid crash-looping: restart with backoff, limit retries
|
||||||
|
min_uptime: "30s",
|
||||||
|
max_restarts: 15,
|
||||||
|
restart_delay: 10000,
|
||||||
|
kill_timeout: 5000,
|
||||||
|
autorestart: true,
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
// ── Supervisor (external watchdog) ──
|
name: "abiba-telegram",
|
||||||
name: "zulip-watchdog",
|
script: "/usr/bin/pitg",
|
||||||
script: "/root/.pi/agent/extensions/zulip/watchdog.js",
|
|
||||||
cwd: "/root",
|
cwd: "/root",
|
||||||
|
|
||||||
max_restarts: 10,
|
|
||||||
min_uptime: "3s",
|
|
||||||
restart_delay: 3000,
|
|
||||||
kill_timeout: 5000,
|
|
||||||
|
|
||||||
log_date_format: "YYYY-MM-DD HH:mm:ss Z",
|
|
||||||
error_file: "/root/.pm2/logs/zulip-watchdog-error.log",
|
|
||||||
out_file: "/root/.pm2/logs/zulip-watchdog-out.log",
|
|
||||||
merge_logs: true,
|
|
||||||
|
|
||||||
autorestart: true,
|
|
||||||
watch: false,
|
|
||||||
instances: 1,
|
|
||||||
exec_mode: "fork",
|
|
||||||
|
|
||||||
env: {
|
|
||||||
NODE_ENV: "production",
|
|
||||||
},
|
|
||||||
},
|
},
|
||||||
],
|
],
|
||||||
};
|
};
|
||||||
|
|||||||
@@ -1,11 +0,0 @@
|
|||||||
{
|
|
||||||
"lastChangelogVersion": "0.80.6",
|
|
||||||
"defaultProvider": "syslog-harness",
|
|
||||||
"defaultModel": "syslog-auto",
|
|
||||||
"defaultThinkingLevel": "high",
|
|
||||||
"extensions": [
|
|
||||||
"+extensions/mcp/index.ts",
|
|
||||||
"+extensions/zulip/index.js"
|
|
||||||
],
|
|
||||||
"theme": "dark"
|
|
||||||
}
|
|
||||||
@@ -1,84 +0,0 @@
|
|||||||
/**
|
|
||||||
* Zulip Gateway Supervisor — external watchdog process.
|
|
||||||
*
|
|
||||||
* Monitors the router health endpoint every 30s. If 3 consecutive checks fail,
|
|
||||||
* restarts the abiba-zulip PM2 process gracefully.
|
|
||||||
*
|
|
||||||
* This is the pattern Hermes uses: an external supervisor that can recover
|
|
||||||
* the gateway even when the gateway process itself is hung (not just crashed).
|
|
||||||
*
|
|
||||||
* Deployed via PM2 as a separate process in ecosystem.config.cjs.
|
|
||||||
*/
|
|
||||||
|
|
||||||
const HEALTH_URL = "http://127.0.0.1:9200/health";
|
|
||||||
const CHECK_INTERVAL_MS = 30_000;
|
|
||||||
const MAX_FAILURES = 3;
|
|
||||||
const RESTART_GRACE_MS = 15_000;
|
|
||||||
|
|
||||||
let failures = 0;
|
|
||||||
|
|
||||||
async function check() {
|
|
||||||
try {
|
|
||||||
const res = await fetch(HEALTH_URL, {
|
|
||||||
signal: AbortSignal.timeout(5000),
|
|
||||||
});
|
|
||||||
|
|
||||||
if (res.ok) {
|
|
||||||
const data = await res.json();
|
|
||||||
|
|
||||||
if (data.status === "ok" && data.zulip?.connected) {
|
|
||||||
if (failures > 0) {
|
|
||||||
console.log(`[watchdog] Router recovered after ${failures} failure(s)`);
|
|
||||||
}
|
|
||||||
failures = 0;
|
|
||||||
// Silent health log every 10 checks (~5 min) for monitoring
|
|
||||||
if (Math.random() < 0.1) {
|
|
||||||
console.log(`[watchdog] Router healthy (uptime: ${Math.round(process.uptime())}s)`);
|
|
||||||
}
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
|
|
||||||
// Connected but degraded
|
|
||||||
console.warn(`[watchdog] Router degraded: status=${data.status}, connected=${data.zulip?.connected}`);
|
|
||||||
failures++;
|
|
||||||
} else {
|
|
||||||
console.warn(`[watchdog] Health check returned ${res.status}`);
|
|
||||||
failures++;
|
|
||||||
}
|
|
||||||
} catch (err) {
|
|
||||||
failures++;
|
|
||||||
console.warn(`[watchdog] Health check ${failures}/${MAX_FAILURES}: ${err.message}`);
|
|
||||||
}
|
|
||||||
|
|
||||||
if (failures >= MAX_FAILURES) {
|
|
||||||
console.error(`[watchdog] ${MAX_FAILURES} consecutive failures — restarting abiba-zulip`);
|
|
||||||
|
|
||||||
const { execSync } = await import("node:child_process");
|
|
||||||
try {
|
|
||||||
execSync("pm2 restart abiba-zulip", { timeout: 30000, encoding: "utf-8" });
|
|
||||||
console.log("[watchdog] Restart command sent successfully");
|
|
||||||
} catch (e) {
|
|
||||||
console.error(`[watchdog] Restart failed: ${e.message}`);
|
|
||||||
|
|
||||||
// Fallback: try resurrect if restart fails (process may be deleted)
|
|
||||||
try {
|
|
||||||
execSync("pm2 resurrect", { timeout: 30000 });
|
|
||||||
console.log("[watchdog] PM2 resurrected (fallback)");
|
|
||||||
} catch (e2) {
|
|
||||||
console.error(`[watchdog] Resurrect also failed: ${e2.message}`);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
failures = 0;
|
|
||||||
// Wait for restart to fully initialize before checking again
|
|
||||||
await new Promise((r) => setTimeout(r, RESTART_GRACE_MS));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
console.log("[watchdog] Zulip gateway supervisor started");
|
|
||||||
console.log(`[watchdog] Monitoring ${HEALTH_URL} every ${CHECK_INTERVAL_MS / 1000}s`);
|
|
||||||
console.log(`[watchdog] Max failures before restart: ${MAX_FAILURES}`);
|
|
||||||
|
|
||||||
// Immediate first check, then periodic
|
|
||||||
check();
|
|
||||||
setInterval(check, CHECK_INTERVAL_MS);
|
|
||||||
File diff suppressed because it is too large
Load Diff
@@ -43,6 +43,7 @@ import time
|
|||||||
import uuid
|
import uuid
|
||||||
from collections import deque
|
from collections import deque
|
||||||
from datetime import datetime, timezone
|
from datetime import datetime, timezone
|
||||||
|
from pathlib import Path
|
||||||
from typing import Any, Dict, List, Optional, Tuple
|
from typing import Any, Dict, List, Optional, Tuple
|
||||||
|
|
||||||
try:
|
try:
|
||||||
@@ -64,14 +65,9 @@ logger = logging.getLogger(__name__)
|
|||||||
|
|
||||||
# Constants
|
# Constants
|
||||||
DEFAULT_STREAM = "agent-hub"
|
DEFAULT_STREAM = "agent-hub"
|
||||||
# SyslogGH realm: the @all-bots user has user_id=20 (verified via
|
DEFAULT_ALL_BOTS_USER_ID = 1
|
||||||
# /api/v1/users and CONTRACT_VERIFICATION_2026-06-29). Dynamic resolution
|
|
||||||
# on connect overrides this; this value is the fallback when that fails —
|
|
||||||
# user_id=1 was never valid in this realm and silently dropped @all-bots.
|
|
||||||
DEFAULT_ALL_BOTS_USER_ID = 20
|
|
||||||
DEFAULT_POLL_INTERVAL = 3.0
|
DEFAULT_POLL_INTERVAL = 3.0
|
||||||
MAX_ZULIP_MESSAGE = 10000
|
MAX_ZULIP_MESSAGE = 10000
|
||||||
TRUNCATION_NOTICE = "\n\n[...truncated at Zulip limit]"
|
|
||||||
ECHO_TAG_PREFIX = "hermes-zulip-"
|
ECHO_TAG_PREFIX = "hermes-zulip-"
|
||||||
RECONNECT_BACKOFF = [2, 5, 10, 30, 60]
|
RECONNECT_BACKOFF = [2, 5, 10, 30, 60]
|
||||||
DEDUP_WINDOW = 300 # 5 minutes
|
DEDUP_WINDOW = 300 # 5 minutes
|
||||||
@@ -125,7 +121,7 @@ def _truncate(text: str, limit: int = MAX_ZULIP_MESSAGE) -> str:
|
|||||||
"""Truncate to Zulip's message limit with notice."""
|
"""Truncate to Zulip's message limit with notice."""
|
||||||
if len(text) <= limit:
|
if len(text) <= limit:
|
||||||
return text
|
return text
|
||||||
return text[:limit - len(TRUNCATION_NOTICE)] + TRUNCATION_NOTICE
|
return text[:limit] + "\n\n[...truncated at Zulip limit]"
|
||||||
|
|
||||||
|
|
||||||
def _parse_zulip_timestamp(ts: Any) -> datetime:
|
def _parse_zulip_timestamp(ts: Any) -> datetime:
|
||||||
@@ -886,6 +882,116 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
reply_to_topic=topic,
|
reply_to_topic=topic,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
# ------------------------------------------------------------------
|
||||||
|
# Media / attachment extraction (vision support)
|
||||||
|
# ------------------------------------------------------------------
|
||||||
|
|
||||||
|
async def _extract_inline_media(
|
||||||
|
self, content: str, media_urls: List[str], media_types: List[str],
|
||||||
|
) -> None:
|
||||||
|
"""Download Zulip ``/user_uploads/...`` attachments referenced in a
|
||||||
|
message so vision/transcription tools can read them locally.
|
||||||
|
|
||||||
|
Zulip renders uploaded files as ``<img src="/user_uploads/...">``
|
||||||
|
in HTML. After ``_strip_html()`` the markdown-style references
|
||||||
|
survive in the plain text. This method finds those references,
|
||||||
|
downloads the files, and caches them into Hermes media storage.
|
||||||
|
"""
|
||||||
|
for rel in re.findall(r"\]\((/user_uploads/[^)]+)\)", content):
|
||||||
|
dl_url = f"{self._site}{rel}"
|
||||||
|
try:
|
||||||
|
resp = await self._http_client.get(
|
||||||
|
dl_url,
|
||||||
|
headers={"Authorization": self._auth_header},
|
||||||
|
timeout=30.0,
|
||||||
|
)
|
||||||
|
if resp.status_code >= 400:
|
||||||
|
logger.debug(
|
||||||
|
"[%s] Inline media %s -> %d",
|
||||||
|
self.name, rel, resp.status_code,
|
||||||
|
)
|
||||||
|
continue
|
||||||
|
data = resp.content
|
||||||
|
mime = (
|
||||||
|
resp.headers.get("content-type", "application/octet-stream")
|
||||||
|
.split(";")[0]
|
||||||
|
.strip()
|
||||||
|
)
|
||||||
|
except Exception:
|
||||||
|
logger.debug("[%s] Failed to download inline media %s", self.name, rel)
|
||||||
|
continue
|
||||||
|
fname = rel.rsplit("/", 1)[-1] or "file"
|
||||||
|
await self._cache_attachment(data, fname, mime, media_urls, media_types)
|
||||||
|
|
||||||
|
async def _extract_event_attachments(
|
||||||
|
self, message: Dict[str, Any], media_urls: List[str], media_types: List[str],
|
||||||
|
) -> None:
|
||||||
|
"""Download files from Zulip's ``message.attachments`` field (UI file
|
||||||
|
uploads that Zulip sends as structured event data rather than inline
|
||||||
|
markdown links)."""
|
||||||
|
atts = message.get("attachments") or []
|
||||||
|
for att in atts:
|
||||||
|
path_id = att.get("path_id")
|
||||||
|
fname = att.get("name") or "file"
|
||||||
|
if not path_id:
|
||||||
|
continue
|
||||||
|
dl_url = f"{self._site}/api/v1/user_uploads/{path_id}"
|
||||||
|
try:
|
||||||
|
resp = await self._http_client.get(
|
||||||
|
dl_url,
|
||||||
|
headers={"Authorization": self._auth_header},
|
||||||
|
timeout=30.0,
|
||||||
|
)
|
||||||
|
if resp.status_code >= 400:
|
||||||
|
logger.debug(
|
||||||
|
"[%s] Event attachment %s -> %d",
|
||||||
|
self.name, fname, resp.status_code,
|
||||||
|
)
|
||||||
|
continue
|
||||||
|
data = resp.content
|
||||||
|
mime = (
|
||||||
|
resp.headers.get("content-type", "application/octet-stream")
|
||||||
|
.split(";")[0]
|
||||||
|
.strip()
|
||||||
|
)
|
||||||
|
except Exception:
|
||||||
|
logger.debug("[%s] Failed to download event attachment %s", self.name, fname)
|
||||||
|
continue
|
||||||
|
await self._cache_attachment(data, fname, mime, media_urls, media_types)
|
||||||
|
|
||||||
|
async def _cache_attachment(
|
||||||
|
self, data: bytes, fname: str, mime: str,
|
||||||
|
media_urls: List[str], media_types: List[str],
|
||||||
|
) -> None:
|
||||||
|
"""Cache a downloaded attachment into Hermes media storage so the
|
||||||
|
vision_analyze / ocr tools can access it via a local file path."""
|
||||||
|
try:
|
||||||
|
from gateway.platforms.base import (
|
||||||
|
cache_image_from_bytes,
|
||||||
|
cache_audio_from_bytes,
|
||||||
|
cache_document_from_bytes,
|
||||||
|
)
|
||||||
|
except ImportError:
|
||||||
|
logger.warning("[%s] cache_*_from_bytes not available in gateway", self.name)
|
||||||
|
return
|
||||||
|
ext = Path(fname).suffix
|
||||||
|
try:
|
||||||
|
if mime.startswith("image/"):
|
||||||
|
media_urls.append(cache_image_from_bytes(data, ext or ".png"))
|
||||||
|
media_types.append(mime)
|
||||||
|
elif mime.startswith("audio/"):
|
||||||
|
media_urls.append(cache_audio_from_bytes(data, ext or ".ogg"))
|
||||||
|
media_types.append(mime)
|
||||||
|
else:
|
||||||
|
media_urls.append(cache_document_from_bytes(data, fname))
|
||||||
|
media_types.append(mime)
|
||||||
|
except Exception as exc:
|
||||||
|
logger.warning("[%s] Failed to cache attachment %s: %s", self.name, fname, exc)
|
||||||
|
|
||||||
|
# ------------------------------------------------------------------
|
||||||
|
# Message routing
|
||||||
|
# ------------------------------------------------------------------
|
||||||
|
|
||||||
async def _route_message(
|
async def _route_message(
|
||||||
self,
|
self,
|
||||||
text: str,
|
text: str,
|
||||||
@@ -916,7 +1022,21 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
chat_id_alt=str(reply_to_id) if reply_to_id else None,
|
chat_id_alt=str(reply_to_id) if reply_to_id else None,
|
||||||
)
|
)
|
||||||
|
|
||||||
# Derive message type
|
# Gen 5: Extract media/attachments so vision tools can read them
|
||||||
|
media_urls: List[str] = []
|
||||||
|
media_types: List[str] = []
|
||||||
|
await self._extract_inline_media(text, media_urls, media_types)
|
||||||
|
await self._extract_event_attachments(raw, media_urls, media_types)
|
||||||
|
|
||||||
|
# Derive message type — upgrade to PHOTO/VOICe/DOCUMENT when media present
|
||||||
|
if media_types:
|
||||||
|
if any(m.startswith("image/") for m in media_types):
|
||||||
|
message_type = MessageType.PHOTO
|
||||||
|
elif any(m.startswith("audio/") for m in media_types):
|
||||||
|
message_type = MessageType.VOICE
|
||||||
|
else:
|
||||||
|
message_type = MessageType.DOCUMENT
|
||||||
|
else:
|
||||||
message_type = MessageType.TEXT
|
message_type = MessageType.TEXT
|
||||||
|
|
||||||
message_event = MessageEvent(
|
message_event = MessageEvent(
|
||||||
@@ -926,6 +1046,8 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
message_id=msg_id,
|
message_id=msg_id,
|
||||||
raw_message=raw,
|
raw_message=raw,
|
||||||
timestamp=_parse_zulip_timestamp(timestamp),
|
timestamp=_parse_zulip_timestamp(timestamp),
|
||||||
|
media_urls=media_urls or None,
|
||||||
|
media_types=media_types or None,
|
||||||
)
|
)
|
||||||
|
|
||||||
# Store reply metadata so send() can respond to the right thread
|
# Store reply metadata so send() can respond to the right thread
|
||||||
|
|||||||
Reference in New Issue
Block a user