/** * pi-zulip-extension v2 — Router-worker architecture for per-sender session isolation * * Deploy to: ~/.pi/agent/extensions/zulip/ * Config: config.yaml (alongside extension) * * Architecture: * Router (this process, PM2 abiba-zulip, ZULIP_ROLE=router): * Polls Zulip for events, maintains a pool of per-sender pi RPC workers. * Each sender gets their own pi --mode rpc process with a dedicated session file. * Messages are routed to the correct worker; worker responses are relayed * back to Zulip (with streaming edits). * Worker (child pi --mode rpc --session-dir ...): * Extension loads but is NO-OP (ZULIP_ROLE != router). Just runs the agent * with per-sender persistent sessions. No Zulip logic in the worker. */ import fs from "node:fs"; import path from "node:path"; import url from "node:url"; import http from "node:http"; import { spawn } from "node:child_process"; import { execSync } from "node:child_process"; import { parse } from "yaml"; // --------------------------------------------------------------------------- // Constants // --------------------------------------------------------------------------- const EXT_DIR = path.dirname(url.fileURLToPath(import.meta.url)); const SESSION_DIR = path.join(process.env.HOME, ".pi", "agent", "sessions", "zulip"); const SESSIONS_FILE = path.join(SESSION_DIR, "worker-sessions.json"); const POLL_INTERVAL_MS = 3000; const WORKER_IDLE_TIMEOUT_MS = 30 * 60 * 1000; // kill idle workers after 30 min const WORKER_BUSY_TIMEOUT_MS = 5 * 60 * 1000; // kill stuck busy workers after 5 min const STREAMING_EDIT_MS = 8000; // throttle Zulip message edits during streaming const MAX_ZULIP_MSG = 10000; const HEARTBEAT_INTERVAL_MS = 5 * 60 * 1000; // --------------------------------------------------------------------------- // Config loading (cleaned — flat yaml format expected) // --------------------------------------------------------------------------- function loadConfig() { const configPath = path.join(EXT_DIR, "config.yaml"); if (!fs.existsSync(configPath)) { throw new Error(`No config.yaml at ${configPath}`); } const raw = fs.readFileSync(configPath, "utf-8"); const cfg = parse(raw); function get(key) { if (cfg[key] !== undefined) return cfg[key]; const parts = key.split("."); let cur = cfg; for (const p of parts) { if (cur && typeof cur === "object" && p in cur) cur = cur[p]; else return undefined; } return cur; } return { zulip: { email: String(get("zulip.email") ?? ""), api_key: process.env.ZULIP_API_KEY ?? String(get("zulip.api_key") ?? ""), site: String(get("zulip.site") ?? ""), all_bots_user_id: Number(get("zulip.all_bots_user_id") ?? 1), }, agent: { name: String(get("agent.name") ?? "abiba"), owner_email: String(get("agent.owner_email") ?? "jerome@sysloggh.com"), }, health_port: parseInt(String(get("health_port") ?? "9200"), 10), }; } const config = loadConfig(); // --------------------------------------------------------------------------- // Crash prevention — the router must never die from an unhandled error. // Without these, any uncaught throw or rejected promise kills the process and // PM2 eventually exhausts max_restarts, taking the Zulip gateway offline. // --------------------------------------------------------------------------- process.on("uncaughtException", (err) => { console.error(`[zulip-ext] UNCAUGHT EXCEPTION (recovering): ${err.message}`); console.error(err.stack); lastError = `uncaught: ${err.message}`; // Attempt reconnect instead of dying connected = false; retryDelay = 10000; setTimeout(startPolling, retryDelay); }); process.on("unhandledRejection", (reason, promise) => { const msg = reason instanceof Error ? reason.message : String(reason); console.error(`[zulip-ext] UNHANDLED REJECTION (recovering): ${msg}`); if (reason?.stack) console.error(reason.stack); lastError = `unhandledRejection: ${msg}`; }); // --------------------------------------------------------------------------- // Resilience patterns — circuit breaker + retry with exponential backoff // Based on Zulip event system docs and battle-tested Node.js resilience patterns. // --------------------------------------------------------------------------- /** * Circuit breaker with CLOSED → OPEN → HALF_OPEN states. * Implements the same pattern as Opossum but inline (no dependency). * Prevents cascading failures when Zulip server is degraded. */ class ZulipCircuitBreaker { constructor(opts = {}) { this.state = "CLOSED"; this.failures = 0; this.successes = 0; this.totalRequests = 0; this.lastFailureTime = null; this.openedAt = null; this.failureThreshold = opts.failureThreshold ?? 0.5; this.resetTimeout = opts.resetTimeout ?? 30000; this.volumeThreshold = opts.volumeThreshold ?? 5; this.onStateChange = opts.onStateChange ?? (() => {}); } async fire(fn) { if (this.state === "OPEN") { if (Date.now() - this.openedAt > this.resetTimeout) { this._transition("HALF_OPEN"); } else { const err = new Error("Circuit is OPEN"); err.name = "CircuitOpenError"; err.circuitState = this.state; throw err; } } try { const result = await fn(); this._onSuccess(); return result; } catch (err) { this._onFailure(err); throw err; } } _transition(newState) { const old = this.state; this.state = newState; if (newState === "OPEN") this.openedAt = Date.now(); if (newState !== old) { console.error(`[zulip-ext] CIRCUIT BREAKER: ${old} → ${newState} (${this.failures}/${this.totalRequests} failures)`); this.onStateChange(newState, old); } } _onSuccess() { this.successes++; this.totalRequests++; if (this.state === "HALF_OPEN") { this._transition("CLOSED"); this.failures = 0; this.successes = 0; this.totalRequests = 0; } // Periodic counter decay to avoid stale stats if (this.totalRequests >= this.volumeThreshold * 4) { this.failures = Math.floor(this.failures / 2); this.successes = Math.floor(this.successes / 2); this.totalRequests = Math.floor(this.totalRequests / 2); } } _onFailure(err) { this.failures++; this.totalRequests++; this.lastFailureTime = Date.now(); if (this.totalRequests >= this.volumeThreshold && this.failures / this.totalRequests >= this.failureThreshold) { this._transition("OPEN"); } } getStats() { return { state: this.state, failures: this.failures, successes: this.successes, totalRequests: this.totalRequests, failureRate: this.totalRequests > 0 ? (this.failures / this.totalRequests).toFixed(3) : "0.000", openedAt: this.openedAt, }; } } /** * Retry with exponential backoff + jitter. * Only retries transient errors (network, fetch failed). * Respects the circuit breaker — if circuit is open, don't retry. */ async function withRetry(fn, opts = {}) { const { maxAttempts = 3, baseDelay = 200, maxDelay = 10000, shouldRetry = () => true, } = opts; let lastError; for (let attempt = 1; attempt <= maxAttempts; attempt++) { try { return await fn(); } catch (err) { lastError = err; if (attempt === maxAttempts) throw err; if (err.name === "CircuitOpenError") throw err; // Don't retry open circuit if (!shouldRetry(err)) throw err; const delay = Math.min(baseDelay * Math.pow(2, attempt - 1), maxDelay); const jitter = delay * (0.5 + Math.random() * 0.5); 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; } function isTransientError(err) { const msg = (err instanceof Error ? err.message : String(err)).toLowerCase(); return msg.includes("fetch failed") || msg.includes("econn") || msg.includes("network") || msg.includes("timeout") || msg.includes("abort") || msg.includes("502") || msg.includes("503") || msg.includes("504"); } // Global circuit breaker for Zulip API calls const zulipBreaker = new ZulipCircuitBreaker({ failureThreshold: 0.5, resetTimeout: 30000, volumeThreshold: 2, // Low threshold — polls are ~10s apart, 5 would need 50s onStateChange: (newState, oldState) => { if (newState === "OPEN" && zulipQueue) { // Notify owner that Zulip connection is degraded zulipQueue.sendMessage({ type: "private", to: config.agent.owner_email, content: `⚠️ **Zulip connection degraded** — circuit breaker open. Will retry in 30s.`, }).catch(() => {}); } else if (newState === "CLOSED" && oldState === "HALF_OPEN" && zulipQueue) { zulipQueue.sendMessage({ type: "private", to: config.agent.owner_email, content: `✅ **Zulip connection restored** — circuit breaker closed.`, }).catch(() => {}); } }, }); // --------------------------------------------------------------------------- // Zulip client (raw fetch — no zulip-js dependency) // --------------------------------------------------------------------------- function authHeader() { return "Basic " + Buffer.from(`${config.zulip.email}:${config.zulip.api_key}`).toString("base64"); } /** * Create a Zulip event queue and resolve bot + all-bots user IDs. */ async function createZulipQueue() { const auth = authHeader(); // Register event queue (with idle_timeout for Zulip 12.0+ queue persistence) const regForm = new URLSearchParams(); regForm.append("event_types", JSON.stringify(["message"])); regForm.append("idle_queue_timeout", "600"); // 10 min — server keeps queue alive longer const regRes = await fetch(`${config.zulip.site}/api/v1/register`, { method: "POST", headers: { Authorization: auth, "Content-Type": "application/x-www-form-urlencoded" }, body: regForm, }); if (!regRes.ok) throw new Error(`Queue registration failed: ${regRes.status} ${await regRes.text().catch(() => '').then(t => t.slice(0,100))}`); const qr = await regRes.json(); let lastEventId = qr.last_event_id ?? -1; let botUserId = null; let allBotsUserId = config.zulip.all_bots_user_id ?? 1; // Resolve bot's own user_id try { const meRes = await fetch(`${config.zulip.site}/api/v1/users/me`, { headers: { Authorization: auth } }); if (meRes.ok) { const md = await meRes.json(); if (md.user_id) botUserId = md.user_id; } } catch { /* non-critical */ } // Resolve @all-bots user_id dynamically try { const usersRes = await fetch(`${config.zulip.site}/api/v1/users`, { headers: { Authorization: auth } }); if (usersRes.ok) { const ud = await usersRes.json(); for (const m of ud.members || []) { if ((m.email || "").toLowerCase().includes("all-bots")) { allBotsUserId = m.user_id; break; } } } } catch { /* fall back to config */ } console.log(`[zulip-ext] Bot user_id=${botUserId}, all-bots user_id=${allBotsUserId}`); return { queueId: qr.queue_id, lastEventId, botUserId, allBotsUserId, async poll() { const res = await fetch( `${config.zulip.site}/api/v1/events?queue_id=${encodeURIComponent(qr.queue_id)}&last_event_id=${lastEventId}`, { headers: { Authorization: authHeader() }, signal: AbortSignal.timeout(65000) } ); if (!res.ok) { if (res.status === 429) { const ra = parseFloat(res.headers.get("retry-after") || "1"); await new Promise(r => setTimeout(r, ra * 1000 + 100)); return []; } const errBody = await res.text().catch(() => ''); throw new Error(`Events API ${res.status}: ${errBody.slice(0, 200)}`); } const data = await res.json(); if (data.events) { for (const e of data.events) { if (e.id > lastEventId) lastEventId = e.id; } return data.events.filter(e => e.type === "message"); } return []; }, async sendMessage(params) { const fb = new URLSearchParams(); fb.append("type", params.type); fb.append("to", params.type === "private" ? JSON.stringify([params.to]) : String(params.to)); if (params.subject) fb.append("subject", params.subject); fb.append("content", params.content); const res = await fetch(`${config.zulip.site}/api/v1/messages`, { method: "POST", headers: { Authorization: authHeader(), "Content-Type": "application/x-www-form-urlencoded" }, body: fb.toString(), }); if (!res.ok) { const errBody = await res.text().catch(() => ''); throw new Error(`Send failed ${res.status}: ${(errBody || '').slice(0, 200)}`); } return (await res.json()).id; }, async editMessage(messageId, content) { const fb = new URLSearchParams(); fb.append("content", content); const res = await fetch(`${config.zulip.site}/api/v1/messages/${messageId}`, { method: "PATCH", headers: { Authorization: authHeader(), "Content-Type": "application/x-www-form-urlencoded" }, body: fb.toString(), }); if (!res.ok) { console.warn(`[zulip-ext] Edit ${messageId} failed: ${res.status}`); } }, async sendTyping(userIds, op) { const fb = new URLSearchParams(); fb.append("to", JSON.stringify(userIds)); fb.append("op", op); await fetch(`${config.zulip.site}/api/v1/typing`, { method: "POST", headers: { Authorization: authHeader(), "Content-Type": "application/x-www-form-urlencoded" }, body: fb.toString(), }).catch(() => {}); }, }; } // --------------------------------------------------------------------------- // Dynamic echo prevention — fetch bot users from Zulip API // --------------------------------------------------------------------------- let BOT_EMAILS = new Set([config.zulip.email]); // always include self async function loadBotEmails(queue) { try { const res = await fetch(`${config.zulip.site}/api/v1/users`, { headers: { Authorization: authHeader() }, }); if (res.ok) { const data = await res.json(); const count = BOT_EMAILS.size; for (const m of data.members || []) { if (m.is_bot && m.email) BOT_EMAILS.add(m.email); } console.log(`[zulip-ext] Echo prevention: ${BOT_EMAILS.size} bot emails (${BOT_EMAILS.size - count} dynamic)`); } } catch (e) { console.warn(`[zulip-ext] Could not fetch bot list: ${e.message}`); } } function isBotSender(email) { return BOT_EMAILS.has(email); } const ALL_BOTS_RE = /@\*\*(?:all-bots|All Bots)\*\*/i; function isAllBotsMention(content) { return ALL_BOTS_RE.test(content); } // --------------------------------------------------------------------------- // Worker pool manager // --------------------------------------------------------------------------- /** @type {Map} */ const workers = new Map(); let workerSessions = {}; // senderId → sessionFilePath (persisted to disk) let cmdIdCounter = 0; let idleCleanupTimer = null; function loadWorkerSessions() { try { if (fs.existsSync(SESSIONS_FILE)) { return JSON.parse(fs.readFileSync(SESSIONS_FILE, "utf-8")); } } catch (e) { console.warn(`[zulip-ext] Failed to load worker sessions: ${e.message}`); } return {}; } function saveWorkerSessions() { try { if (!fs.existsSync(SESSION_DIR)) fs.mkdirSync(SESSION_DIR, { recursive: true }); fs.writeFileSync(SESSIONS_FILE, JSON.stringify(workerSessions, null, 2)); } catch (e) { console.warn(`[zulip-ext] Failed to save worker sessions: ${e.message}`); } } /** * Attach a JSONL line reader to a stream. */ function attachJsonlReader(stream, onLine, onError) { let buf = ""; const td = new TextDecoder(); stream.on("data", (chunk) => { buf += typeof chunk === "string" ? chunk : td.decode(chunk, { stream: true }); while (true) { const idx = buf.indexOf("\n"); if (idx === -1) break; const line = buf.slice(0, idx).replace(/\r$/, ""); buf = buf.slice(idx + 1); if (!line) continue; try { onLine(JSON.parse(line)); } catch (e) { if (onError) onError(e, line); } } }); } /** * Spawn a pi RPC worker process for a sender. */ function spawnWorker(senderId) { const sessionDir = SESSION_DIR; if (!fs.existsSync(sessionDir)) fs.mkdirSync(sessionDir, { recursive: true }); const proc = spawn("pi", [ "--mode", "rpc", "--session-dir", sessionDir, ], { env: { ...process.env, ZULIP_ROLE: "worker", ZULIP_EXTENSION_ACTIVE: "", // ensure NOT active }, stdio: ["pipe", "pipe", "pipe"], cwd: process.env.HOME, }); const worker = { process: proc, senderId, busy: false, streamingBuffer: "", lastStreamEdit: 0, pendingReplies: [], // [{ zulipMsgId, senderId, senderName, replyInfo, createdAt }] lastMessageContent: null, lastReplyInfo: null, lastActivity: Date.now(), busySince: null, createdAt: Date.now(), sessionFile: workerSessions[senderId] || null, pendingCommands: new Map(), rpcCommandQueue: [], }; // Read stdout events (JSONL) attachJsonlReader( proc.stdout, (event) => handleWorkerEvent(worker, event), (err, line) => console.warn(`[zulip-ext] Worker ${senderId} parse error: ${err.message} | ${line.slice(0, 100)}`) ); // Pipe stderr for diagnostics proc.stderr.on("data", (d) => { const text = d.toString().trim(); if (text) process.stderr.write(`[worker-${senderId}] ${text}\n`); }); // Handle exit proc.on("exit", (code, signal) => { console.log(`[zulip-ext] Worker ${senderId} exited: code=${code} signal=${signal || "none"}`); // Fail any pending replies while (worker.pendingReplies.length > 0) { const reply = worker.pendingReplies.shift(); if (reply.zulipMsgId && zulipQueue) { zulipQueue.editMessage(reply.zulipMsgId, ":warning: Response interrupted — worker process exited. Please try again." ).catch(() => {}); } } workers.delete(senderId); }); // Resolve session file (async () => { // Guard: verify session file exists before trying to resume. // A stale reference to a deleted file causes the worker to get stuck // in a spawn-resume-fail loop — see WAL node #657 for the incident. if (worker.sessionFile && !fs.existsSync(worker.sessionFile)) { console.log(`[zulip-ext] Worker ${senderId}: stale session ref (file missing), clearing`); delete workerSessions[senderId]; saveWorkerSessions(); worker.sessionFile = null; } if (worker.sessionFile) { // Resume existing session const cmdId = String(++cmdIdCounter); sendRpc(worker, { id: cmdId, type: "switch_session", sessionPath: worker.sessionFile }); console.log(`[zulip-ext] Worker ${senderId}: resuming session ${worker.sessionFile}`); } // If no session file, the auto-created session is fine. Discover it. if (!worker.sessionFile) { const cmdId = String(++cmdIdCounter); sendRpc(worker, { id: cmdId, type: "get_state" }); } })(); return worker; } function sendRpc(worker, cmd) { worker.process.stdin.write(JSON.stringify(cmd) + "\n"); } /** * Handle all RPC events from a worker. */ function handleWorkerEvent(worker, event) { // Debug: log all worker events to diagnose processing if (event.type !== 'message_update' || !event.assistantMessageEvent || event.assistantMessageEvent.type === 'text_start' || event.assistantMessageEvent.type === 'text_end') { console.log(`[zulip-ext] Worker ${worker.senderId} event: ${event.type}${event.assistantMessageEvent ? ':' + event.assistantMessageEvent.type : ''}${event.command ? ' cmd=' + event.command : ''}`); } switch (event.type) { case "response": handleWorkerResponse(worker, event); break; case "message_update": handleWorkerStreaming(worker, event); break; case "agent_end": handleWorkerAgentEnd(worker, event); break; case "agent_start": break; case "turn_start": break; case "turn_end": break; case "compaction_start": console.log(`[zulip-ext] Worker ${worker.senderId}: compaction started`); break; case "compaction_end": console.log(`[zulip-ext] Worker ${worker.senderId}: compaction ended`); // Track new session file path if it changed break; case "extension_error": console.error(`[zulip-ext] Worker ${worker.senderId} ext error: ${event.error?.slice(0, 200)}`); break; case "message_start": case "message_end": case "tool_execution_start": case "tool_execution_update": case "tool_execution_end": case "queue_update": case "auto_retry_start": case "auto_retry_end": break; // informational, not needed by router } } function handleWorkerResponse(worker, event) { // Resolve pending command callbacks if (event.id) { const cb = worker.pendingCommands.get(event.id); if (cb) { worker.pendingCommands.delete(event.id); cb.resolve(event); return; } } // Handle responses without explicit callbacks if (event.command === "get_state" && event.success) { const sessionFile = event.data?.sessionFile; if (sessionFile && !worker.sessionFile) { worker.sessionFile = sessionFile; workerSessions[worker.senderId] = sessionFile; saveWorkerSessions(); console.log(`[zulip-ext] Worker ${worker.senderId}: session file = ${sessionFile}`); } } if (event.command === "switch_session") { if (event.success && !event.data?.cancelled) { console.log(`[zulip-ext] Worker ${worker.senderId}: session switched OK`); } else { // Session restore failed (missing file, corrupted, etc). // Clear stale reference and discover the new auto-created session. console.log(`[zulip-ext] Worker ${worker.senderId}: session switch failed (cancelled=${event.data?.cancelled}, err=${event.data?.error}), falling back to new session`); delete workerSessions[worker.senderId]; saveWorkerSessions(); worker.sessionFile = null; const cmdId = String(++cmdIdCounter); sendRpc(worker, { id: cmdId, type: "get_state" }); } } if (event.command === "new_session" && event.success) { // Track the new session file (get_state will follow if needed) if (!event.data?.cancelled) { const cmdId = String(++cmdIdCounter); worker.pendingCommands.set(cmdId, { resolve: (ev) => { const sf = ev.data?.sessionFile; if (sf) { worker.sessionFile = sf; workerSessions[worker.senderId] = sf; saveWorkerSessions(); console.log(`[zulip-ext] Worker ${worker.senderId}: new session = ${sf}`); } }, }); sendRpc(worker, { id: cmdId, type: "get_state" }); } } } function handleWorkerStreaming(worker, event) { const delta = event.assistantMessageEvent; if (!delta || delta.type !== "text_delta") return; worker.streamingBuffer += delta.delta; const now = Date.now(); if (worker.streamingBuffer.length > 10 && (!worker.lastStreamEdit || now - worker.lastStreamEdit >= STREAMING_EDIT_MS)) { worker.lastStreamEdit = now; const reply = worker.pendingReplies[0]; if (reply?.zulipMsgId && zulipQueue) { const preview = worker.streamingBuffer.length > 9500 ? worker.streamingBuffer.slice(0, 9500) + "\n\n_… still generating…_" : worker.streamingBuffer; zulipQueue.editMessage(reply.zulipMsgId, preview).catch(() => {}); } } } async function handleWorkerAgentEnd(worker, event) { worker.busy = false; worker.busySince = null; worker.streamingBuffer = ""; worker.lastStreamEdit = 0; worker.lastActivity = Date.now(); if (worker.pendingReplies.length === 0) return; const reply = worker.pendingReplies.shift(); // Extract final assistant text — walk backwards to find the last message with text content // (the final turn may only have tool calls with brief inline text like "On it. Let me") const assistantMsgs = (event.messages || []).filter(m => m.role === "assistant"); let responseText = ""; for (let i = assistantMsgs.length - 1; i >= 0; i--) { const content = assistantMsgs[i].content; if (!content) continue; if (typeof content === "string") { responseText = content; break; } if (Array.isArray(content)) { const texts = content.filter(c => c.type === "text").map(c => c.text).join("\n"); if (texts.trim()) { responseText = texts; break; } } } // Stop typing if (reply.senderId && zulipQueue) { zulipQueue.sendTyping([reply.senderId], "stop").catch(() => {}); } if (!responseText.trim()) { // All assistant messages are tool-call-only. Send a summary. console.log(`[zulip-ext] Tool-only response for ${reply.senderName} — sending activity note`); responseText = "⚙️ _Running tools…_\n\nUse `/continue` to see results or send another message."; } const truncated = responseText.length > MAX_ZULIP_MSG ? responseText.slice(0, MAX_ZULIP_MSG) + "\n\n[...truncated at Zulip limit]" : responseText; const target = reply.replyInfo.type === "private" ? { type: "private", to: reply.senderId } : { type: "stream", to: reply.replyInfo.streamId ?? reply.replyInfo.streamName, subject: reply.replyInfo.topic }; try { if (reply.zulipMsgId && zulipQueue) { await zulipQueue.editMessage(reply.zulipMsgId, truncated); console.log(`[zulip-ext] Finalized reply for ${reply.senderName} on worker ${worker.senderId}`); } else if (zulipQueue) { await zulipQueue.sendMessage({ ...target, content: truncated }); console.log(`[zulip-ext] Sent reply for ${reply.senderName}`); } } catch (err) { console.error(`[zulip-ext] Failed to send reply: ${err.message}`); } } /** * Kill a worker and clean up. */ function killWorker(senderId) { const w = workers.get(senderId); if (!w) return; console.log(`[zulip-ext] Killing idle worker ${senderId}`); workers.delete(senderId); try { w.process.kill("SIGTERM"); } catch { /* already dead */ } // Note: workerSessions entry is kept so session resumes on re-connect } /** * Periodically kill idle workers. */ function cleanupIdleWorkers() { const now = Date.now(); for (const [senderId, w] of workers) { if (w.busy) continue; if (now - w.lastActivity > WORKER_IDLE_TIMEOUT_MS) { killWorker(senderId); } } } /** * Kill workers stuck in busy state > WORKER_BUSY_TIMEOUT_MS. * This is the critical fix — without it, a hung LLM call deadlocks the router forever. */ function cleanupStuckWorkers() { const now = Date.now(); for (const [senderId, w] of workers) { if (!w.busy || !w.busySince) continue; if (now - w.busySince > WORKER_BUSY_TIMEOUT_MS) { console.log(`[zulip-ext] Worker ${senderId} stuck busy for ${Math.round((now - w.busySince) / 1000)}s — force-killing`); // Send error to pending replies while (w.pendingReplies.length > 0) { const reply = w.pendingReplies.shift(); if (reply.zulipMsgId && zulipQueue) { zulipQueue.editMessage(reply.zulipMsgId, ":warning: Response timed out — worker was stuck for over 5 minutes. Please try again." ).catch(() => {}); } } // Kill and remove the stuck worker (session ref kept for re-connect) workers.delete(senderId); try { w.process.kill("SIGKILL"); } catch { /* already dead */ } console.log(`[zulip-ext] Worker ${senderId} force-killed. Will re-spawn on next message.`); } } } // --------------------------------------------------------------------------- // Attachment download // --------------------------------------------------------------------------- /** Download as raw buffer (for text extraction / base64 image encoding). */ async function downloadAttachmentRaw(pathId) { try { const res = await fetch(`${config.zulip.site}/api/v1/user_uploads/${pathId}`, { headers: { Authorization: authHeader() }, signal: AbortSignal.timeout(30000), }); if (!res.ok) return null; const buf = Buffer.from(await res.arrayBuffer()); return { contentType: res.headers.get("content-type") || "application/octet-stream", buf, size: buf.length }; } catch (e) { console.error(`[zulip-ext] Download failed: ${e.message}`); return null; } } /** Download and base64-encode (for pi image API). */ async function downloadAttachment(pathId) { const raw = await downloadAttachmentRaw(pathId); if (!raw) return null; return { contentType: raw.contentType, data: raw.buf.toString("base64"), size: raw.size }; } // --------------------------------------------------------------------------- // Attachment classification & text extraction // --------------------------------------------------------------------------- const IMAGE_EXTS = new Set(["png", "jpg", "jpeg", "gif", "webp", "bmp", "svg"]); const PDF_EXTS = new Set(["pdf"]); const TEXT_EXTS = new Set([ "txt", "md", "py", "js", "ts", "jsx", "tsx", "json", "yaml", "yml", "log", "csv", "sh", "bash", "env", "cfg", "conf", "ini", "toml", "xml", "html", "htm", "css", "scss", "sql", "rs", "go", "java", "c", "cpp", "cc", "h", "hpp", "rb", "php", "swift", "kt", "scala", "r", "lua", "pl", "ps1", "bat", "makefile", "dockerfile", "gitignore", "editorconfig", "diff", "patch", "tex", "rst", "org", "nix", ]); const MAX_TEXT_ATTACHMENT_CHARS = 50_000; function classifyAttachment(att) { const name = att.name || ""; const ext = name.split(".").pop()?.toLowerCase() || ""; if (IMAGE_EXTS.has(ext)) return "image"; if (PDF_EXTS.has(ext)) return "pdf"; if (TEXT_EXTS.has(ext)) return "text"; // MIME-based fallback return "binary"; } function formatSize(bytes) { if (bytes < 1024) return `${bytes} B`; if (bytes < 1024 * 1024) return `${(bytes / 1024).toFixed(1)} KB`; return `${(bytes / (1024 * 1024)).toFixed(1)} MB`; } function tryDecodeText(buf) { try { return buf.toString("utf-8"); } catch { return null; } } /** Try to extract text from a PDF using pdftotext (poppler-utils). */ function extractPdfText(buf) { try { const tmpIn = `/tmp/zulip-pdf-${Date.now()}.pdf`; fs.writeFileSync(tmpIn, buf); const text = execSync(`pdftotext -layout -nopgbrk "${tmpIn}" -`, { timeout: 15000, encoding: "utf-8", maxBuffer: 5 * 1024 * 1024, }); fs.unlinkSync(tmpIn); return text.trim() || null; } catch (e) { console.error(`[zulip-ext] PDF extraction failed: ${e.message}`); return null; } } // --------------------------------------------------------------------------- // Message routing // --------------------------------------------------------------------------- /** * Send a message to a worker and track the pending reply. */ async function routeToWorker(senderId, senderName, replyInfo, content, attachments) { let worker = workers.get(senderId); if (!worker) { worker = spawnWorker(senderId); workers.set(senderId, worker); } // Send placeholder to Zulip const placeholders = [":robot: _Processing…_", ":hourglass_flowing_sand: _Thinking…_", ":brain: _Generating…_"]; const placeholder = placeholders[Date.now() % placeholders.length]; const target = replyInfo.type === "private" ? { type: "private", to: senderId } : { type: "stream", to: replyInfo.streamId ?? replyInfo.streamName, subject: replyInfo.topic }; let zulipMsgId = null; if (zulipQueue) { try { zulipMsgId = await zulipQueue.sendMessage({ ...target, content: placeholder }); } catch (e) { console.warn(`[zulip-ext] Placeholder send failed: ${e.message}`); } } worker.pendingReplies.push({ zulipMsgId, senderId, senderName, replyInfo, createdAt: Date.now(), }); // Build prompt text const label = replyInfo.type === "stream" ? `@all-bots in #${replyInfo.streamName}` : `DM from @${senderName}`; const promptText = `[Zulip ${label}]: ${content}`; // Handle attachments — download, classify, extract content let images = []; let attachmentTexts = []; if (attachments?.length) { for (const att of attachments) { const category = classifyAttachment(att); const raw = await downloadAttachmentRaw(att.path_id); if (!raw) { attachmentTexts.push(`[Attachment: ${att.name} — download failed]`); continue; } switch (category) { case "image": { const b64 = raw.buf.toString("base64"); images.push({ type: "image", data: b64, mimeType: raw.contentType }); break; } case "text": { const text = tryDecodeText(raw.buf); if (text) { const truncated = text.length > MAX_TEXT_ATTACHMENT_CHARS ? text.slice(0, MAX_TEXT_ATTACHMENT_CHARS) + `\n\n[...truncated at ${MAX_TEXT_ATTACHMENT_CHARS.toLocaleString()} chars, ${(text.length - MAX_TEXT_ATTACHMENT_CHARS).toLocaleString()} chars removed]` : text; attachmentTexts.push( `[Attachment: ${att.name} (${formatSize(raw.size)})]\n\`\`\`${truncated}\n\`\`\`` ); } else { attachmentTexts.push(`[Attachment: ${att.name} (${formatSize(raw.size)}, ${raw.contentType}) — binary content, could not decode as text]`); } break; } case "pdf": { const pdfText = extractPdfText(raw.buf); if (pdfText) { const truncated = pdfText.length > MAX_TEXT_ATTACHMENT_CHARS ? pdfText.slice(0, MAX_TEXT_ATTACHMENT_CHARS) + `\n\n[...truncated at ${MAX_TEXT_ATTACHMENT_CHARS.toLocaleString()} chars]` : pdfText; attachmentTexts.push( `[Attachment: ${att.name} (PDF, ${formatSize(raw.size)})]\n\`\`\`${truncated}\n\`\`\`` ); } else { attachmentTexts.push(`[Attachment: ${att.name} (PDF, ${formatSize(raw.size)}) — text extraction failed (pdftotext unavailable or unreadable PDF)]`); } break; } default: { attachmentTexts.push(`[Attachment: ${att.name} (${formatSize(raw.size)}, ${raw.contentType}) — binary file, content not extracted]`); break; } } } } // Merge attachment text into prompt const fullText = attachmentTexts.length > 0 ? promptText + "\n\n---\n" + attachmentTexts.join("\n\n") : promptText; // Send to worker if (worker.busy) { sendRpc(worker, { type: "steer", message: fullText }); console.log(`[zulip-ext] Steered to worker ${senderId} for ${senderName}${attachmentTexts.length ? ' with ' + attachmentTexts.length + ' attachment(s)' : ''}`); } else { const cmd = { type: "prompt", message: fullText }; if (images.length > 0) cmd.images = images; sendRpc(worker, cmd); const attInfo = []; if (images.length > 0) attInfo.push(`${images.length} image(s)`); if (attachmentTexts.length > 0) attInfo.push(`${attachmentTexts.length} file(s)`); console.log(`[zulip-ext] Routed to worker ${senderId} for ${senderName}${attInfo.length ? ' with ' + attInfo.join(' + ') : ''}`); } worker.busy = true; worker.busySince = Date.now(); worker.lastMessageContent = content; worker.lastReplyInfo = replyInfo; worker.lastActivity = Date.now(); // Send typing indicator if (zulipQueue) { zulipQueue.sendTyping([senderId], "start").catch(() => {}); } } // --------------------------------------------------------------------------- // Zulip command handlers // --------------------------------------------------------------------------- const ZULIP_COMMANDS = ["status", "retry", "continue", "new", "help", "stop", "model", "compact", "yolo", "verbose", "approve", "deny"]; function parseCommand(content) { const m = content.match(/^\/(\w+)(?:\s+(.+))?$/); return m ? { cmd: m[1].toLowerCase(), args: m[2] || "" } : null; } /** * Handle a Zulip command. Returns true if handled, false to pass through as normal message. */ async function handleZulipCommand(cmd, args, senderId, senderName, replyInfo) { const target = replyInfo.type === "private" ? { type: "private", to: senderId } : { type: "stream", to: replyInfo.streamId ?? replyInfo.streamName, subject: replyInfo.topic }; const send = (content) => { if (zulipQueue) zulipQueue.sendMessage({ ...target, content }).catch(() => {}); }; switch (cmd) { // --- Global commands (no worker needed) --- case "help": { const cmds = { status: "Show system status — PM2 processes, containers, uptime", retry: "Retry the last response", continue: "Continue the last response if cut off", new: "Start fresh — clears conversation context", help: "Show this help message", stop: "Stop the current response", model: "Show or change LLM model", compact: "Compact conversation context", yolo: "Acknowledge YOLO mode (no confirmations)", verbose: "Toggle verbose output", approve: "Approve pending actions (Hermes compat)", deny: "Deny pending actions (Hermes compat)", }; if (args && cmds[args]) { send(`**/${args}** — ${cmds[args]}\n\nUsage: \`/${args}\``); } else { let helpText = "**Abiba Zulip Commands**\n\n"; for (const [name, desc] of Object.entries(cmds)) { helpText += `• \`/${name}\` — ${desc}\n`; } helpText += `\nSend any message to chat.`; send(helpText); } return true; } case "status": { let resp = "**Abiba System Status**\n\n**PM2 Processes:**\n"; try { const pm2 = execSync("pm2 status --no-color 2>/dev/null", { timeout: 10000, encoding: "utf-8" }); for (const line of pm2.split("\n")) { if (line.includes("abiba-")) { const cells = line.split("│").filter(c => c.trim()); if (cells.length >= 6) { const clean = (s) => s.replace(/\x1b\[[0-9;]*m/g, ""); resp += `• ${clean(cells[1])}: ${clean(cells[5])} (uptime: ${clean(cells[3])})\n`; } } } } catch (e) { resp += `_Error: ${e.message}_\n`; } resp += `\n**Worker Sessions:** ${workers.size} active\n`; for (const [sid, w] of workers) { resp += `• sender_${sid}: ${w.busy ? "busy" : "idle"}, replies pending: ${w.pendingReplies.length}\n`; } resp += `\n**Zulip:** connected=${!!zulipQueue} | processed=${msgCount} | echo prevented bots=${BOT_EMAILS.size}`; send(resp); return true; } // --- Per-sender commands (need worker) --- case "retry": { const worker = workers.get(senderId); if (!worker?.lastMessageContent) { send(":warning: No previous message to retry."); return true; } send(":repeat: Retrying last message…"); routeToWorker(senderId, senderName, worker.lastReplyInfo, worker.lastMessageContent); return true; } case "continue": { send(":arrow_forward: Continuing…"); routeToWorker(senderId, senderName, replyInfo, "Continue your response from where you left off."); return true; } case "new": { const worker = workers.get(senderId); if (worker) { worker.lastMessageContent = null; sendRpc(worker, { type: "new_session" }); } const topicMsg = args ? `Start fresh, focusing on: ${args}` : "Start fresh."; send(`:broom: ${topicMsg}\n\n_Previous context cleared._`); routeToWorker(senderId, senderName, replyInfo, `[System: Session reset. ${topicMsg} Treat this as a fresh conversation.]`); return true; } case "stop": { const worker = workers.get(senderId); if (worker?.busy) { sendRpc(worker, { type: "abort" }); send("⏹️ **Stop signal sent.** Current processing interrupted."); } else { send(":information_source: No active processing to stop."); } return true; } case "model": { const worker = workers.get(senderId); if (!worker) { send("`/model` requires an active conversation. Send a message first."); return true; } sendRpc(worker, { type: "get_available_models" }); // We'll need to async-read the response. For now, just acknowledge. // In a follow-up, the response can be captured and sent to Zulip. send("🔍 Checking available models…"); return true; } case "compact": { const worker = workers.get(senderId); if (worker) sendRpc(worker, { type: "compact" }); send("🗜️ **Compaction triggered.** Context will be summarized."); return true; } case "yolo": send("🤘 YOLO mode acknowledged. All tools execute immediately."); return true; case "verbose": send("📝 For detailed output, include 'be thorough' or 'show your work' in your message."); return true; case "approve": case "deny": send("ℹ️ Pi doesn't have pending approvals — commands execute immediately."); return true; default: send(`:question: Unknown command \`/${cmd}\`. Try \`/help\`.`); return true; } } // --------------------------------------------------------------------------- // Polling loop // --------------------------------------------------------------------------- let zulipQueue = null; let connected = false; let retryCount = 0; let retryDelay = 5000; let msgCount = 0; let skippedCount = 0; let lastHeartbeatTime = 0; let lastError = null; async function startPolling() { console.log(`[zulip-ext] Connecting to ${config.zulip.site} as ${config.zulip.email}…`); try { zulipQueue = await zulipBreaker.fire(() => createZulipQueue()); connected = true; retryCount = 0; retryDelay = 5000; lastError = null; await loadBotEmails(zulipQueue); console.log(`[zulip-ext] Connected, queue=${zulipQueue.queueId}`); // Resilience-first event loop (Zulip call_on_each_event pattern) // Uses circuit breaker + retry with jitter instead of raw polling while (connected && zulipQueue) { try { const events = await zulipBreaker.fire(() => withRetry(() => zulipQueue.poll(), { maxAttempts: 2, baseDelay: 1000, shouldRetry: isTransientError, }) ); lastError = null; retryCount = 0; if (events.length > 0) { for (const ev of events) { await processEvent(ev); } } // Heartbeat const now = Date.now(); if (now - lastHeartbeatTime > HEARTBEAT_INTERVAL_MS) { lastHeartbeatTime = now; const workerInfo = Array.from(workers.values()).map(w => `${w.senderId}:${w.busy ? "busy" : "idle"}:${w.pendingReplies.length}`).join(","); console.log(`[zulip-ext] Heartbeat — msgs=${msgCount} skipped=${skippedCount} workers=[${workerInfo}] queue=${zulipQueue.queueId}`); } // Normal poll interval between cycles await new Promise(r => setTimeout(r, POLL_INTERVAL_MS)); } 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}`); retryCount = 0; lastError = null; } catch (reRegErr) { console.error(`[zulip-ext] Re-registration failed: ${reRegErr.message}`); connected = false; break; // Exit loop, full reconnect below } } else if (err.name === "CircuitOpenError") { // Circuit is open — pause and retry lastError = "circuit_open"; console.warn(`[zulip-ext] Circuit OPEN — pausing ${Math.round(zulipBreaker.resetTimeout / 1000)}s`); await new Promise(r => setTimeout(r, zulipBreaker.resetTimeout)); } else if (isTransientError(err)) { 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)); } else { lastError = msg; console.error(`[zulip-ext] Unrecoverable poll error: ${msg}`); connected = false; break; } } } // If we exited the loop (disconnected), attempt full reconnect if (!connected) { retryDelay = Math.min(5000 * Math.pow(2, retryCount), 300000); console.log(`[zulip-ext] Full reconnect in ${Math.round(retryDelay / 1000)}s (attempt ${retryCount + 1})…`); retryCount++; setTimeout(startPolling, retryDelay); } } catch (err) { const msg = err instanceof Error ? err.message : String(err); lastError = msg; console.error(`[zulip-ext] Connection failed: ${msg}`); // Exponential backoff (capped at ~5 min) retryDelay = Math.min(5000 * Math.pow(2, retryCount), 300000); console.log(`[zulip-ext] Reconnecting in ${Math.round(retryDelay / 1000)}s (attempt ${retryCount + 1})…`); retryCount++; setTimeout(startPolling, retryDelay); } } async function processEvent(event) { if (event.type !== "message") return; const msg = event.message; if (!msg) return; // Ignore own messages if (msg.sender_email === config.zulip.email) return; // Echo prevention if (isBotSender(msg.sender_email)) { skippedCount++; if (skippedCount % 10 === 1) { console.log(`[zulip-ext] Echo skip: ${msg.sender_email} (total=${skippedCount})`); } return; } const senderName = msg.sender_full_name || msg.sender_email; const senderId = msg.sender_id; // keep as number for correct API format // Private messages (DM-first, ADR-001) if (msg.type === "private") { msgCount++; console.log(`[zulip-ext] DM from ${senderName} (id=${senderId}): ${(msg.content || "").slice(0, 80)}`); // Check for commands const parsed = parseCommand(msg.content); if (parsed && ZULIP_COMMANDS.includes(parsed.cmd)) { await handleZulipCommand(parsed.cmd, parsed.args, senderId, senderName, { type: "private", senderId, senderName, }); return; } routeToWorker(senderId, senderName, { type: "private", senderId, senderName }, msg.content, msg.attachments); return; } // Stream messages: @mention or @all-bots (ADR-005, ADR-006) const mentionedUserIds = (msg.mentioned_users || []).map(u => u.user_id); const isDirectMention = zulipQueue?.botUserId && mentionedUserIds.includes(zulipQueue.botUserId); const isAllBots = isAllBotsMention(msg.content || ""); if (msg.type === "stream" && (isDirectMention || isAllBots)) { msgCount++; const streamName = typeof msg.display_recipient === "string" ? msg.display_recipient : "unknown"; const topic = msg.subject || "general"; console.log(`[zulip-ext] ${isDirectMention ? "@mention" : "@all-bots"} in #${streamName} > ${topic}`); const replyInfo = { type: "stream", streamName, topic, streamId: msg.stream_id, senderId, senderName }; const parsed = parseCommand(msg.content); if (parsed && ZULIP_COMMANDS.includes(parsed.cmd)) { await handleZulipCommand(parsed.cmd, parsed.args, senderId, senderName, replyInfo); return; } routeToWorker(senderId, senderName, replyInfo, msg.content, msg.attachments); } } // --------------------------------------------------------------------------- // Health endpoint // --------------------------------------------------------------------------- let healthServer = null; function startHealthServer(port) { const server = http.createServer((req, res) => { if (req.url === "/health" || req.url === "/") { const workerList = []; for (const [sid, w] of workers) { workerList.push({ sender_id: sid, busy: w.busy, pending_replies: w.pendingReplies.length, last_activity: Math.floor((Date.now() - w.lastActivity) / 1000), }); } res.writeHead(200, { "Content-Type": "application/json" }); res.end(JSON.stringify({ status: connected && zulipQueue ? "ok" : "down", platform: "pi", agent: config.agent.name, zulip: { connected, site: config.zulip.site, email: config.zulip.email, queue_id: zulipQueue?.queueId ?? null, bot_user_id: zulipQueue?.botUserId ?? null, messages_processed: msgCount, skipped: skippedCount, last_error: lastError, }, circuit_breaker: zulipBreaker.getStats(), workers: workerList, worker_count: workers.size, })); } else { res.writeHead(404); res.end("Not found"); } }); server.on("error", (err) => { if (err.code === "EADDRINUSE") { console.warn(`[zulip-ext] Port :${port} in use — skipping health`); } }); server.listen(port, "127.0.0.1", () => { console.log(`[zulip-ext] Health endpoint on :${port}`); }); return server; } // --------------------------------------------------------------------------- // Extension lifecycle // --------------------------------------------------------------------------- export default function (pi) { // Guard: only activate in the designated PM2 router process if (process.env.ZULIP_ROLE !== "router") { // Worker processes and TUI sessions: NO-OP return; } pi.on("session_start", async () => { console.log(`[zulip-ext] Router starting for ${config.agent.name}`); // Load persisted worker session map if (!fs.existsSync(SESSION_DIR)) fs.mkdirSync(SESSION_DIR, { recursive: true }); workerSessions = loadWorkerSessions(); // Start health server if (!healthServer) { healthServer = startHealthServer(config.health_port); } // Start Zulip poller startPolling(); // Start worker cleanup (idle + stuck, every 60s) idleCleanupTimer = setInterval(() => { cleanupIdleWorkers(); cleanupStuckWorkers(); }, 60 * 1000); console.log(`[zulip-ext] Router ready. Known sessions: ${Object.keys(workerSessions).length}`); }); pi.on("session_shutdown", async () => { console.log(`[zulip-ext] Router shutting down…`); connected = false; if (idleCleanupTimer) { clearInterval(idleCleanupTimer); idleCleanupTimer = null; } // Kill all workers gracefully for (const [sid, w] of workers) { try { w.process.kill("SIGTERM"); } catch {} } workers.clear(); // Delete Zulip event queue if (zulipQueue?.queueId) { try { await fetch(`${config.zulip.site}/api/v1/queues/${zulipQueue.queueId}`, { method: "DELETE", headers: { Authorization: authHeader() }, }); console.log(`[zulip-ext] Queue ${zulipQueue.queueId} deregistered`); } catch (e) { console.warn(`[zulip-ext] Queue deregister failed: ${e.message}`); } } if (healthServer) { healthServer.close(); healthServer = null; } connected = false; zulipQueue = null; saveWorkerSessions(); }); // Register diagnostics commands pi.registerCommand("zulip-status", { description: "Show Zulip extension status and worker pool", handler: async (_args, ctx) => { const lines = [ "=== Zulip Extension v2 (router-worker) ===", `Agent: ${config.agent.name}`, `Server: ${config.zulip.site}`, `Connected: ${connected ? "✅" : "❌"} Queue: ${zulipQueue?.queueId ?? "N/A"}`, `Messages: ${msgCount} processed, ${skippedCount} echo-skipped`, `Workers: ${workers.size} active`, `Retries: ${retryCount} Last error: ${lastError ?? "none"}`, "", ]; for (const [sid, w] of workers) { lines.push(`• sender_${sid}: ${w.busy ? "busy" : "idle"}, ${w.pendingReplies.length} pending, session=${w.sessionFile ?? "new"}`); } ctx.ui.notify(lines.join("\n"), "info"); }, }); pi.registerCommand("zulip-send-test", { description: "Send test DM: /zulip-send-test ", handler: async (args, ctx) => { const parts = args.trim().match(/^([^\s]+)\s+(.+)$/); const to = parts?.[1] ?? config.agent.owner_email; const msg = parts?.[2] ?? "Test from pi Zulip extension v2"; if (!zulipQueue) { ctx.ui.notify("Not connected. Check /zulip-status.", "error"); return; } try { await zulipQueue.sendMessage({ type: "private", to, content: msg }); ctx.ui.notify(`Test sent to ${to}`, "info"); } catch (err) { ctx.ui.notify(`Failed: ${err.message}`, "error"); } }, }); pi.registerCommand("zulip-test", { description: "Run comprehensive Zulip self-test", handler: async (_args, ctx) => { const results = []; const passed = (n, d) => results.push({ name: n, ok: true, detail: d }); const failed = (n, d) => results.push({ name: n, ok: false, detail: d }); passed("queue_connected", `Queue: ${zulipQueue?.queueId ?? "N/A"}`); zulipQueue?.botUserId ? passed("bot_identity", `Bot ID: ${zulipQueue.botUserId}`) : failed("bot_identity", "Bot ID not resolved"); zulipQueue && zulipQueue.allBotsUserId > 1 ? passed("all_bots", `ID: ${zulipQueue.allBotsUserId}`) : failed("all_bots", `ID: ${zulipQueue?.allBotsUserId ?? 1}`); BOT_EMAILS.size > 1 ? passed("echo_prevention", `${BOT_EMAILS.size} bot emails`) : failed("echo_prevention", "Only self excluded"); workers.size >= 0 ? passed("worker_pool", `${workers.size} active workers`) : failed("worker_pool", "Pool broken"); // Loopback send if (zulipQueue) { try { const id = await zulipQueue.sendMessage({ type: "private", to: config.agent.owner_email, content: `🔬 Zulip self-test at ${new Date().toISOString()}` }); passed("api_send", `Loopback DM sent (id=${id})`); } catch (e) { failed("api_send", e.message.slice(0, 80)); } } const ok = results.filter(r => r.ok).length; const total = results.length; const verdict = ok === total ? "HEALTHY" : ok >= total - 1 ? "DEGRADED" : "CRITICAL"; let report = `=== Zulip Self-Test: ${verdict} ===\n${ok}/${total} checks passed\n\n`; for (const r of results) { report += `${r.ok ? "✅" : "❌"} ${r.name}: ${r.detail}\n`; } ctx.ui.notify(report, verdict === "HEALTHY" ? "info" : "error"); }, }); }