Files
prose-contracts/zulip-resilience-v3.prose.md
root 23f3f378c5
PR Pipeline — Authorize → Validate → Review → Merge / auth (pull_request) Successful in 4s
PR Pipeline — Authorize → Validate → Review → Merge / validate (pull_request) Successful in 3s
PR Pipeline — Authorize → Validate → Review → Merge / lint (pull_request) Successful in 6s
PR Pipeline — Authorize → Validate → Review → Merge / ai-review (pull_request) Successful in 5s
PR Pipeline — Authorize → Validate → Review → Merge / gate (pull_request) Successful in 0s
zulip-resilience: add fleet-wide audit findings and fixes from 2026-07-18
- Added Incident Log section documenting fleet-wide Zulip audit
- Abiba: poll timeout AbortError fix (returns [] instead of error)
- Abiba: credential fallback .env file for Infisical outages
- Tanko: full gateway restart to recover Zulip connection
- Fleet health metrics summary table
- Hermes agent improvement recommendations (circuit breaker, credential fallback, queue re-registration, watchdog)
- Abiba pi extension v2.1 resilience feature matrix
2026-07-18 16:33:12 +00:00

477 lines
18 KiB
Markdown

---
kind: responsibility
name: zulip-resilience-v3
description: >
Rewrite the pi Zulip gateway with production-grade resilience patterns drawn from
Zulip's own event system docs (queue lifecycle, heartbeat monitoring, BAD_EVENT_QUEUE_ID
handling, idle_queue_timeout) and battle-tested Node.js resilience patterns
(circuit breaker, exponential backoff with jitter, bulkhead isolation, supervisor watchdog).
replaces: zulip-self-heal (retired)
agent: abiba
triggers:
- "/zulip self-heal v3"
- "zulip stopped responding"
- "PM2 abiba-zulip crashed"
---
# Zulip Gateway v3 — Production Resilience
## Architecture Overview
The current v2 gateway (`/root/.pi/agent/extensions/zulip/index.js`) has three structural
weaknesses that cause repeated deaths:
1. **No crash recovery** — uncaught errors kill the Node process, PM2 exhausts max_restarts
2. **No circuit breaker** — 502/fetch-failed errors escalate to process death with no fallback
3. **No queue lifecycle management** — doesn't use Zulip's documented heartbeat protocol or
idle_queue_timeout, so BAD_EVENT_QUEUE_ID errors cascade into crashes
The v3 rewrite addresses all three, following patterns from:
- [Zulip Events System docs](https://zulip.readthedocs.io/en/11.6/subsystems/events-system.html) —
queue registration, heartbeat, BAD_EVENT_QUEUE_ID recovery, call_on_each_event loop
- [Zulip API: Get Events](https://zulip.com/api/get-events) — long-poll timeout, dont_block, event ack
- [Circuit Breaker & Retry Patterns in Node.js 2026](https://1xapi.com/blog/resilient-api-circuit-breaker-bulkhead-retry-nodejs-2026) —
Opossum-based circuit breaker with fallback, retry with jitter, bulkhead isolation
---
## Maintains
- `zulip-gateway`: { status: "healthy" | "degraded" | "down" }
- `circuit-breaker`: { state: "CLOSED" | "OPEN" | "HALF_OPEN", failures, successes }
- `queue-lifecycle`: { queue_id, last_event_id, idle_timeout, heartbeat_age }
- `workers`: { count, busy, idle, stuck }
- `supervisor`: { pid, last_check, health_failures }
---
## Detection Rules
### Rule 1: Queue Expired (BAD_EVENT_QUEUE_ID)
- **Detect**: Events API returns error with BAD_EVENT_QUEUE_ID in body
- **Fix**: Call `POST /register` to create new queue, update queue_id and last_event_id
- **Debounce**: If 3 re-registrations fail within 60s, escalate (server may be down)
- **Ref**: Zulip docs: "Your software will need to handle that error condition by re-initializing itself"
### Rule 2: Network Degradation (502/ECONNREFUSED/fetch failed)
- **Detect**: Events API returns 502 or network error
- **Circuit breaker**: Track failure rate over 10s rolling window
- CLOSED → OPEN: 50% failure rate with ≥5 requests
- OPEN → HALF_OPEN: After 30s reset timeout
- HALF_OPEN → CLOSED: Probe succeeds
- HALF_OPEN → OPEN: Probe fails
- **While OPEN**: Log errors, skip events, notify user via DM: "⚠️ Zulip connection degraded — will retry in 30s"
### Rule 3: Long-Poll Timeout (natural)
- **Detect**: Events API response takes > `event_queue_longpoll_timeout_seconds`
- **Not an error**: Server sends heartbeat events when no real events. Simply re-poll.
### Rule 4: Worker Busy Timeout (>5 min)
- **Detect**: Worker `busySince` exceeds 5 minutes
- **Fix**: SIGKILL worker, send error DM, clean up pending replies
### Rule 5: Process Crash (uncaught)
- **Detect**: `uncaughtException` / `unhandledRejection` fires
- **Fix**: Log → clear poll timer → attempt reconnect with backoff → if reconnect fails 3x, exit(1) and let PM2 restart
### Rule 6: Supervisor Detects Router Stall
- **Detect**: External supervisor (`zulip-watchdog`) polls `/health` every 30s. If 3 consecutive failures:
- **Fix**: `pm2 restart abiba-zulip` gracefully (SIGTERM, drain workers, restart)
---
## Implementation Plan
### Phase 1: Rewrite Router Core (circuit-breaker + queue lifecycle)
Replace the poll loop in index.js with a resilience-first event loop:
```js
// Queue lifecycle (Zulip docs pattern)
async function createOrRefreshQueue() {
// POST /register with event_types=["message"]
// Store: queueId, lastEventId, eventQueueLongpollTimeoutSeconds
// NEW: pass idle_queue_timeout parameter (Zulip 12.0+)
}
// Circuit breaker (Opossum pattern, implemented inline to avoid dependency)
class ZulipCircuitBreaker {
constructor({ failureThreshold=0.5, resetTimeout=30000, volumeThreshold=5, windowMs=10000 }) {
this.state = "CLOSED"; // CLOSED | OPEN | HALF_OPEN
this.failures = 0;
this.successes = 0;
this.totalRequests = 0;
this.lastFailureTime = null;
this.openedAt = null;
this.failureThreshold = failureThreshold;
this.resetTimeout = resetTimeout;
this.volumeThreshold = volumeThreshold;
this.windowMs = windowMs;
}
async fire(fn) {
if (this.state === "OPEN") {
if (Date.now() - this.openedAt > this.resetTimeout) {
this.state = "HALF_OPEN";
} else {
throw new CircuitOpenError("Circuit is OPEN");
}
}
try {
const result = await fn();
this.onSuccess();
return result;
} catch (err) {
this.onFailure();
throw err;
}
}
onSuccess() {
this.successes++;
this.totalRequests++;
if (this.state === "HALF_OPEN") {
this.state = "CLOSED";
this.failures = 0;
}
// Reset counters periodically
if (this.totalRequests > this.volumeThreshold * 2) {
this.failures = Math.floor(this.failures / 2);
this.successes = Math.floor(this.successes / 2);
this.totalRequests = Math.floor(this.totalRequests / 2);
}
}
onFailure() {
this.failures++;
this.totalRequests++;
this.lastFailureTime = Date.now();
if (this.totalRequests >= this.volumeThreshold &&
this.failures / this.totalRequests >= this.failureThreshold) {
if (this.state !== "OPEN") {
this.state = "OPEN";
this.openedAt = Date.now();
console.error(`[zulip-ext] CIRCUIT BREAKER OPEN — ${this.failures}/${this.totalRequests} failures`);
}
}
}
}
// Retry with exponential backoff + jitter (from resilience patterns)
async function withRetry(fn, { maxAttempts=3, baseDelay=200, maxDelay=10000, shouldRetry=()=>true }={}) {
let lastError;
for (let attempt = 1; attempt <= maxAttempts; attempt++) {
try {
return await fn();
} catch (err) {
lastError = err;
if (attempt === maxAttempts || !shouldRetry(err)) throw err;
const delay = Math.min(baseDelay * Math.pow(2, attempt - 1), maxDelay);
const jitter = delay * (0.5 + Math.random() * 0.5); // 50-100% of delay
console.warn(`[zulip-ext] Retry ${attempt}/${maxAttempts} after ${Math.round(jitter)}ms: ${err.message.slice(0,80)}`);
await new Promise(r => setTimeout(r, jitter));
}
}
throw lastError;
}
// Resilience-first event loop (Zulip call_on_each_event pattern)
async function resilientPollLoop() {
while (connected) {
try {
const events = await circuitBreaker.fire(() =>
withRetry(() => zulipQueue.poll(), {
maxAttempts: 2,
baseDelay: 1000,
shouldRetry: (err) => {
const msg = err.message || "";
return msg.includes("fetch failed") || msg.includes("ECONN") || msg.includes("network");
}
})
);
lastError = null;
retryCount = 0;
for (const ev of events) {
await processEvent(ev);
}
heartbeat();
} catch (err) {
const msg = err instanceof Error ? err.message : String(err);
if (msg.includes("BAD_EVENT_QUEUE_ID") || msg.includes("deregistered")) {
// Queue expired — re-register (Zulip docs pattern)
console.log(`[zulip-ext] Queue expired, re-registering… (${msg.slice(0,80)})`);
try {
zulipQueue = await createZulipQueue();
console.log(`[zulip-ext] Re-registered, new queue=${zulipQueue.queueId}`);
} catch (reRegErr) {
console.error(`[zulip-ext] Re-registration failed: ${reRegErr.message}`);
connected = false;
retryCount++;
const backoff = Math.min(5000 * Math.pow(2, retryCount), 300000);
console.log(`[zulip-ext] Full reconnect in ${Math.round(backoff/1000)}s`);
await new Promise(r => setTimeout(r, backoff));
await startPolling();
return;
}
} else if (err.name === "CircuitOpenError") {
// Circuit is open — skip this cycle, wait for HALF_OPEN
lastError = "circuit_open";
await new Promise(r => setTimeout(r, POLL_INTERVAL_MS));
} else {
lastError = msg;
retryCount++;
const backoff = Math.min(POLL_INTERVAL_MS * Math.pow(1.5, Math.min(retryCount, 8)), 60000);
console.error(`[zulip-ext] Poll error (retry ${retryCount}, backoff ${backoff}ms): ${msg}`);
await new Promise(r => setTimeout(r, backoff));
}
}
}
}
```
### Phase 2: PM2 Hardening
Create `/root/.pm2/ecosystem.config.cjs`:
```js
module.exports = {
apps: [
{
name: "abiba-zulip",
script: "/bin/pi",
args: "--mode rpc --session-id zulip-service",
env: {
ZULIP_ROLE: "router",
ZULIP_SITE: "https://chat.sysloggh.net",
ZULIP_EMAIL: "abiba-bot@chat.sysloggh.net",
ZULIP_API_KEY: process.env.ZULIP_API_KEY,
AGENT_NAME: "abiba",
AGENT_OWNER_EMAIL: "jerome@sysloggh.com",
},
max_restarts: 100, // Up from default 10 — crash loops won't exhaust
min_uptime: "10s", // Must survive 10s to count as "alive"
max_memory_restart: "500M", // OOM protection
restart_delay: 5000, // 5s between restarts
kill_timeout: 15000, // 15s SIGTERM grace before SIGKILL
listen_timeout: 30000, // 30s to bind health port
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,
autorestart: true,
watch: false,
instances: 1,
exec_mode: "fork",
},
{
name: "zulip-watchdog",
script: "/root/.pi/agent/extensions/zulip/watchdog.js",
max_restarts: 10,
min_uptime: "3s",
restart_delay: 3000,
autorestart: true,
},
],
};
```
### Phase 3: Supervisor Watchdog
Create `/root/.pi/agent/extensions/zulip/watchdog.js`:
```js
// External supervisor — monitors router health and restarts if stalled.
// This is the pattern Hermes uses: an external process that can recover
// the gateway even if the gateway process itself is hung (not just crashed).
const HEALTH_URL = "http://127.0.0.1:9200/health";
const CHECK_INTERVAL_MS = 30_000;
const MAX_FAILURES = 3;
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} failures`);
}
failures = 0;
return;
}
}
failures++;
console.warn(`[watchdog] Health check ${failures}/${MAX_FAILURES}: status not ok`);
} 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 } = require("child_process");
try {
execSync("pm2 restart abiba-zulip", { timeout: 30000 });
console.log("[watchdog] Restart command sent");
} catch (e) {
console.error(`[watchdog] Restart failed: ${e.message}`);
}
failures = 0;
// Wait for restart to complete before checking again
await new Promise(r => setTimeout(r, 15000));
}
}
console.log("[watchdog] Zulip gateway supervisor started");
setInterval(check, CHECK_INTERVAL_MS);
check(); // Immediate first check
```
### Phase 4: Health Endpoint Enhancement
Add circuit breaker stats to the existing health endpoint:
```js
// In /health response, add:
"circuit_breaker": {
"state": circuitBreaker.state,
"failures": circuitBreaker.failures,
"successes": circuitBreaker.successes,
"total_requests": circuitBreaker.totalRequests,
"failure_rate": circuitBreaker.totalRequests > 0
? (circuitBreaker.failures / circuitBreaker.totalRequests).toFixed(2)
: "0.00"
}
```
---
## Test Plan
### Test 1: Queue Re-registration
1. Manually delete the Zulip event queue via API
2. Next poll should detect BAD_EVENT_QUEUE_ID
3. Router should auto re-register within 1 poll cycle
4. Verify: `/health` shows new queue_id, connected=true
### Test 2: Circuit Breaker Trip
1. Block Zulip server with iptables: `iptables -A OUTPUT -d 192.168.68.19 -j DROP`
2. Router should detect failures, trip circuit after 5 failures
3. `/health` should show circuit_breaker.state = "OPEN"
4. Remove iptables rule
5. Circuit should transition to HALF_OPEN → CLOSED within 60s
6. Verify: messages processed after recovery
### Test 3: Supervisor Recovery
1. Kill the router process: `kill -STOP $(pm2 pid abiba-zulip)` (freeze, don't kill)
2. Watchdog should detect 3 failed health checks in 90s
3. Watchdog should execute `pm2 restart abiba-zulip`
4. Verify: router back online, connected=true
### Test 4: Worker Busy Timeout
1. Send a message that triggers a long-running operation
2. If worker stays busy >5 minutes, should receive SIGKILL
3. User should receive error DM: "Response timed out"
### Test 5: End-to-End Message
1. Send DM "What time is it?" from Jerome
2. Should receive response within 30s
3. `/health` should show messages_processed incremented
---
## Rollback Plan
If v3 causes issues:
1. `pm2 delete abiba-zulip; pm2 delete zulip-watchdog`
2. Restore v2 from git: `cd /root/.pi/agent/extensions/zulip && git checkout index.js`
3. `pm2 resurrect` to reload previous process list
4. Verify: `/health` returns ok
Backup v2 before starting: `cp index.js index.js.v2-backup-$(date +%Y%m%d-%H%M%S)`
---
## Success Metrics
| Metric | Current (v2) | Target (v3) |
|--------|-------------|-------------|
| Uptime between manual interventions | 1-3 days | 30+ days |
| Crash recovery | Manual (PM2 resurrect) | Automatic (circuit breaker + supervisor) |
| Queue expiry handling | Crash | Auto re-register |
| Busy worker deadlock | Router death | Worker SIGKILL + error DM |
| PM2 restart exhaustion | Yes (max_restarts=10) | No (max_restarts=100 + watchdog) |
---
## Incident Log — 2026-07-18 Fleet-Wide Audit
### Fleet State After Audit
| Agent | Platform | Zulip State | Issues Found | Fix Applied |
|-------|----------|-------------|--------------|-------------|
| **Abiba** | pi (CT 100) | ✅ Connected | API key missing from Infisical injection; poll timeout noise | Added .env fallback; AbortError treated as empty poll (no retry); poll timeout 65s→90s |
| **Tanko** | Hermes (CT 112) | ✅ Connected | Gateway disconnected since Jul 11; watchdog restart didn't re-establish Zulip | Full gateway restart (kill wrapper, let infisical-gateway.sh respawn) |
| **Mumuni** | Hermes (CT 114) | ✅ Connected | No issues found | None needed |
### Key Fixes Applied
**1. Abiba — Credential Fallback (L4 Pattern)**
- Root cause: `zulip.api_key` in config.yaml is `""` (expected from Infisical). Infisical vault `ABIBA_ZULIP_API_KEY` wasn't being injected into the process environment.
- Fix: Added `.env` file fallback at `/root/.pi/agent/extensions/zulip/.env` with known-working key, sourced before the Infisical `exec`.
- Lesson: Per L4 from gpu-self-heal, Infisical is not always available — always keep a local `.env` fallback.
**2. Abiba — Poll Timeout Handling**
- Root cause: Zulip long-poll uses `AbortSignal.timeout(65000)`. Zulip's default `event_queue_longpoll_timeout_seconds` can exceed 65s. When the signal fires, an `AbortError` is thrown and caught by the circuit breaker as a failure.
- Fix: Caught `AbortError` inside `poll()` and return empty array (no events) instead of throwing. Extended timeout to 90s to match Zulip server default.
- Reference: [Zulip Events System — long-poll timeout](https://zulip.readthedocs.io/en/11.6/subsystems/events-system.html)
**3. Tanko — Gateway Restart**
- Root cause: Gateway process was running but Zulip platform stayed in "disconnected" state since Jul 11, 2026. The wrapper script (`infisical-gateway.sh`) restarts on crash but the gateway wasn't re-establishing Zulip on restart.
- Fix: Killed gateway PID to trigger wrapper restart. New gateway (PID 331991) established Zulip connection successfully.
### Fleet-Wide Zulip Health Metrics (as of 2026-07-18)
| Metric | Value |
|--------|-------|
| Zulip server | ✅ HTTP 200 |
| Agents connected | 3/3 (Abiba, Tanko, Mumuni) |
| Abiba circuit breaker | CLOSED (0 failures) |
| Abiba uptime | 2D (post-restart) |
| Tanko gateway uptime | Ongoing |
| Mumuni gateway uptime | Ongoing |
| Watchdog status | ✅ Online (2D uptime) |
### Hermes Agent Zulip Plugin Improvements
Based on the audit, improvements that should be ported to all Hermes Zulip adapters:
1. **Circuit breaker pattern** — Already in Abiba's pi extension. Hermes adapters should add the same CLOSED→OPEN→HALF_OPEN state machine with exponential backoff.
2. **Credential fallback** — All Hermes agents use Infisical for credentials. Add `.env` local fallback per L4 pattern for `ZULIP_API_KEY`.
3. **Queue re-registration** — Handle `BAD_EVENT_QUEUE_ID` with automatic re-registration instead of gateway restart.
4. **Supervisor watchdog** — Hermes uses PM2 which auto-restarts on crash, but has no health-check watchdog. Add lightweight external health checks.
5. **Streaming** — All agents have `streaming: true` in their zulip config. Verify `edit_message()` is implemented in each adapter.
### Abiba pi Zulip Extension v2 — Implemented Resilience Summary
| Feature | Status | Notes |
|---------|--------|-------|
| Circuit breaker | ✅ | CLOSED→OPEN→HALF_OPEN; 50% failure threshold; 30s reset timeout |
| Retry with jitter | ✅ | 2 attempts, 200ms base, 50-100% jitter |
| Queue lifecycle | ✅ | 10min idle_queue_timeout; BAD_EVENT_QUEUE_ID handling |
| Crash prevention | ✅ | uncaughtException + unhandledRejection recovery |
| Worker timeout | ✅ | 5min busy timeout → SIGKILL + error DM |
| Health endpoint | ✅ | :9200 with circuit breaker metrics |
| Echo prevention | ✅ | Dynamic bot user resolution |
| Poll timeout (AbortError) | ✅ v2.1 | Normal timeout returns [] instead of error |
| Credential fallback | ✅ v2.1 | .env file before Infisical exec |
| Provider auto-fix | ✅ | Detects reasoning_content models, switches to compatible |