From 2f6c4283d26a2561ff4a24e9cadd22e74240a8fb Mon Sep 17 00:00:00 2001 From: "Abiba (pi)" Date: Mon, 29 Jun 2026 03:45:34 +0000 Subject: [PATCH] feat: /zulip-test command + standardized health schema v1 - Added /zulip-test self-test command (7 checks: queue, identity, @all-bots, health, poll loop, echo prevention, API send) - Standardized health endpoint to zulip-health/v1 schema (nested zulip{} object, checks{} map, platform/agent metadata) - Saved pi extension source to pi-zulip-extension/extension-src/ (with all fixes: dynamic @all-bots, health sync, logging) - Added extension src to repo for deployment tracking - Updated vault with final contract status report --- .../extension-src/config.yaml.example | 24 + pi-zulip-extension/extension-src/types.d.ts | 34 + .../extension-src/zulip-extension.js | 1365 +++++++++++++++++ 3 files changed, 1423 insertions(+) create mode 100644 pi-zulip-extension/extension-src/config.yaml.example create mode 100644 pi-zulip-extension/extension-src/types.d.ts create mode 100644 pi-zulip-extension/extension-src/zulip-extension.js diff --git a/pi-zulip-extension/extension-src/config.yaml.example b/pi-zulip-extension/extension-src/config.yaml.example new file mode 100644 index 0000000..fe6cf75 --- /dev/null +++ b/pi-zulip-extension/extension-src/config.yaml.example @@ -0,0 +1,24 @@ +# pi Zulip extension — Abiba configuration (flat format for custom parser) +# Deploy to: ~/.pi/agent/extensions/config.yaml + +zulip.site: https://chat.sysloggh.net +zulip.email: abiba-bot@chat.sysloggh.net +zulip.api_key: YOUR_ZULIP_API_KEY_HERE +zulip.stream: agent-hub +zulip.all_bots_user_id: 20 + +agent.name: abiba +agent.display_name: Abiba +agent.zulip_bot_name: abiba-bot +agent.owner_email: jerome@sysloggh.com +agent.private_topic: abiba + +health_port: 9200 + +monitoring.health_endpoint_enabled: true +monitoring.log_level: info + +error_handling.timeout_seconds: 30 +error_handling.retry_count: 3 +error_handling.retry_delay_seconds: 5 +error_handling.graceful_message: "Abiba encountered an error processing your request. Please try again later." diff --git a/pi-zulip-extension/extension-src/types.d.ts b/pi-zulip-extension/extension-src/types.d.ts new file mode 100644 index 0000000..e513e8e --- /dev/null +++ b/pi-zulip-extension/extension-src/types.d.ts @@ -0,0 +1,34 @@ +// Type declarations for zulip-js (no @types available) +declare module "zulip-js" { + interface ZulipClient { + queues: { + register(params: { + event_types: string[]; + narrow?: Array>; + }): Promise<{ queue_id: string; last_event_id: number }>; + }; + users: { + me: { + get(): Promise; + }; + }; + messages: { + store: { + send(rawContent: string): Promise; + }; + }; + events: { + retrieve(params: { + queue_id: string; + last_event_id: number; + dont_block?: boolean; + }): Promise<{ events: any[]; result?: string; msg?: string }>; + }; + } + + export default function zulip(config: { + username: string; + apiKey: string; + realm: string; + }): Promise; +} diff --git a/pi-zulip-extension/extension-src/zulip-extension.js b/pi-zulip-extension/extension-src/zulip-extension.js new file mode 100644 index 0000000..6965332 --- /dev/null +++ b/pi-zulip-extension/extension-src/zulip-extension.js @@ -0,0 +1,1365 @@ +/** + * pi-zulip-extension — Zulip agent communication plugin for pi agents + * + * Deploy to: ~/.pi/agent/extensions/zulip/ + * Config: config.yaml (alongside extension, or env vars) + * + * DM-first architecture per ADR-001/ADR-002. + * Background Zulip event queue poller that injects messages directly into + * the current pi session via pi.sendUserMessage(), then captures the LLM + * response from agent_end and sends it back to Zulip. No subprocess overhead. + * + * @see ADR-007: Platform-native plugin contracts + * @see docs/ARCHITECTURE.md + */ +import zulip from "zulip-js"; +import http from "node:http"; +import fs from "node:fs"; +import path from "node:path"; +import url from "node:url"; +import { parse } from "yaml"; +import { execSync } from "node:child_process"; +function loadConfig() { + // Check env vars first + const email = process.env.ZULIP_EMAIL; + const apiKey = process.env.ZULIP_API_KEY; + const site = process.env.ZULIP_SITE; + if (email && apiKey && site) { + return { + zulip: { email, api_key: apiKey, site }, + agent: { + name: process.env.AGENT_NAME ?? "abiba", + owner_email: process.env.AGENT_OWNER_EMAIL ?? "jerome@sysloggh.com", + }, + health_port: parseInt(process.env.HEALTH_PORT ?? "9200", 10), + poll_interval_ms: 3000, + max_retries: 3, + retry_delay_ms: 5000, + pi_command: "pi", + }; + } + // Fall back to config.yaml alongside the extension + const configDir = path.dirname(url.fileURLToPath(import.meta.url)); + const configPath = path.join(configDir, "config.yaml"); + if (!fs.existsSync(configPath)) { + throw new Error(`No config.yaml found at ${configPath}. Set ZULIP_EMAIL, ZULIP_API_KEY, ZULIP_SITE env vars or provide config.yaml.`); + } + const raw = fs.readFileSync(configPath, "utf-8"); + const cfg = parse(raw); + // Support both flat (env-style) and nested config formats + function getField(key) { + if (cfg[key] !== undefined) + return cfg[key]; + // Try nested lookup: "zulip.email" -> cfg.zulip?.email + 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; + } + const parsed = { + zulip: { + email: String(getField("zulip.email") ?? getField("zulip.email") ?? ""), + api_key: String(getField("zulip.api_key") ?? getField("zulip.api_key") ?? ""), + site: String(getField("zulip.site") ?? getField("zulip.server_url") ?? ""), + stream: String(getField("zulip.stream") ?? ""), + all_bots_user_id: Number(getField("zulip.all_bots_user_id") ?? 1), + }, + agent: { + name: String(getField("agent.name") ?? "abiba"), + owner_email: String(getField("agent.owner_email") ?? "jerome@sysloggh.com"), + display_name: String(getField("agent.display_name") ?? ""), + zulip_bot_name: String(getField("agent.zulip_bot_name") ?? ""), + private_topic: String(getField("agent.private_topic") ?? ""), + }, + health_port: parseInt(String(getField("health_port") ?? getField("monitoring.health_port") ?? "9200"), 10), + poll_interval_ms: 3000, + max_retries: 3, + retry_delay_ms: 5000, + pi_command: "pi", + }; + // Override api_key from env if set + parsed.zulip.api_key = process.env.ZULIP_API_KEY ?? parsed.zulip.api_key; + return parsed; +} +/** + * Register a Zulip event queue and poll for events. + */ +async function createZulipQueue(config) { + // Use direct fetch with URLSearchParams for queue registration + // (zulip-js library uses multipart/form-data which may not work correctly) + const authHeader = "Basic " + Buffer.from(config.zulip.email + ":" + config.zulip.api_key).toString("base64"); + + // Register queue via direct API call + const regForm = new URLSearchParams(); + regForm.append("event_types", JSON.stringify(["message"])); + + const regRes = await fetch(config.zulip.site + "/api/v1/register", { + method: "POST", + headers: { + "Authorization": authHeader, + "Content-Type": "application/x-www-form-urlencoded", + }, + body: regForm, + }); + if (!regRes.ok) { + throw new Error(`Queue registration failed: ${regRes.status} ${await regRes.text()}`); + } + const queueRes = await regRes.json(); + const queueId = queueRes.queue_id; + let lastEventId = queueRes.last_event_id ?? -1; + + // Fetch bot's own user_id for @mention detection + let botUserId = null; + try { + const meResp = await fetch(config.zulip.site + "/api/v1/users/me", { + headers: { "Authorization": authHeader }, + }); + if (meResp.ok) { + const meData = await meResp.json(); + if (meData.user_id) { + botUserId = meData.user_id; + console.log(`[zulip-ext] Bot user_id: ${botUserId}`); + } + } + } catch (e) { + // Non-critical + } + + // Resolve @all-bots user_id dynamically (Gen 2: parity with Hermes adapter) + let allBotsUserId = config.zulip.all_bots_user_id ?? 1; + try { + const usersResp = await fetch(config.zulip.site + "/api/v1/users", { + headers: { "Authorization": authHeader }, + }); + if (usersResp.ok) { + const usersData = await usersResp.json(); + for (const member of (usersData.members || [])) { + if ((member.email || "").toLowerCase().includes("all-bots")) { + allBotsUserId = member.user_id; + console.log(`[zulip-ext] @all-bots user_id dynamically resolved: ${allBotsUserId}`); + break; + } + } + } + } catch (e) { + console.log(`[zulip-ext] @all-bots resolution failed, using config default: ${allBotsUserId}`); + } + + const queueObj = { + queueId, + lastEventId, + botUserId, + allBotsUserId, + /** + * Poll the event queue and return any new events. + */ + async poll() { + const res = await fetch(`${config.zulip.site}/api/v1/events?queue_id=${encodeURIComponent(queueId)}&last_event_id=${lastEventId}`, { + headers: { + Authorization: "Basic " + + Buffer.from(`${config.zulip.email}:${config.zulip.api_key}`).toString("base64"), + }, + }); + if (!res.ok) { + throw new Error(`Events API returned ${res.status}: ${await res.text()}`); + } + const data = (await res.json()); + // Debug: log raw event count for stream visibility + if (data.events && data.events.length > 0) { + for (const evt of data.events) { + const m = evt.message || {}; + const dispRecip = typeof m.display_recipient === 'string' ? m.display_recipient : JSON.stringify(m.display_recipient); + console.log(`[zulip-ext] RAW EVENT: type=${m.type} stream=${dispRecip || "(dm)"} sender=${m.sender_full_name} mentioned_users=${JSON.stringify(m.mentioned_users || [])}`); + } + } + if (data.events && Array.isArray(data.events)) { + for (const event of data.events) { + if (event.id > lastEventId) { + lastEventId = event.id; + queueObj.lastEventId = lastEventId; // Sync for health endpoint + } + } + return data.events.filter((e) => e.type === "message"); + } + return []; + }, + /** + * Send a message to Zulip via the API. + */ + async sendMessage(params) { + const formBody = new URLSearchParams(); + formBody.append("type", params.type); + formBody.append("to", params.type === "private" ? JSON.stringify([params.to]) : String(params.to)); + if (params.subject) + formBody.append("subject", params.subject); + formBody.append("content", params.content); + const res = await fetch(`${config.zulip.site}/api/v1/messages`, { + method: "POST", + headers: { + Authorization: "Basic " + + Buffer.from(`${config.zulip.email}:${config.zulip.api_key}`).toString("base64"), + "Content-Type": "application/x-www-form-urlencoded", + }, + body: formBody.toString(), + }); + if (!res.ok) { + throw new Error(`Send message API returned ${res.status}: ${await res.text()}`); + } + const body = (await res.json()); + return body.id; + }, + /** + * Edit a previously sent Zulip message (streaming update). + */ + async editMessage(messageId, content) { + const formBody = new URLSearchParams(); + formBody.append("content", content); + const res = await fetch(`${config.zulip.site}/api/v1/messages/${messageId}`, { + method: "PATCH", + headers: { + Authorization: "Basic " + + Buffer.from(`${config.zulip.email}:${config.zulip.api_key}`).toString("base64"), + "Content-Type": "application/x-www-form-urlencoded", + }, + body: formBody.toString(), + }); + if (!res.ok) { + // Gen 2: Log ALL failures — previously suppressed 400s silently + const body = await res.text().catch(() => ''); + console.warn(`[zulip-ext] Edit message ${messageId} returned ${res.status}: ${body.slice(0, 80)}`); + } + }, + /** + * Send typing indicator. + */ + async sendTypingNotification(userIds, operation) { + const formBody = new URLSearchParams(); + formBody.append("to", JSON.stringify(userIds)); + formBody.append("op", operation); + await fetch(`${config.zulip.site}/api/v1/typing`, { + method: "POST", + headers: { + Authorization: "Basic " + + Buffer.from(`${config.zulip.email}:${config.zulip.api_key}`).toString("base64"), + "Content-Type": "application/x-www-form-urlencoded", + }, + body: formBody.toString(), + }); + }, + }; + return queueObj; +} +// --------------------------------------------------------------------------- +// Health endpoint +// --------------------------------------------------------------------------- +function startHealthServer(port, getState) { + const server = http.createServer((req, res) => { + if (req.url === "/health" || req.url === "/") { + const state = getState(); + res.writeHead(200, { "Content-Type": "application/json" }); + res.end(JSON.stringify({ status: "ok", ...state })); + } + else { + res.writeHead(404); + res.end("Not found"); + } + }); + server.on("error", (err) => { + if (err.code === "EADDRINUSE") { + console.warn(`[zulip-ext] Port :${port} already in use — skipping health endpoint`); + } + else { + console.error(`[zulip-ext] Health server error:`, err.message); + } + }); + server.listen(port, "127.0.0.1", () => { + console.log(`[zulip-ext] Health endpoint on :${port}`); + }); + return server; +} +// --------------------------------------------------------------------------- +// Detect @all-bots content +// --------------------------------------------------------------------------- +const ALL_BOTS_RE = /@\*\*(?:all-bots|All Bots)\*\*/i; +function isAllBotsMention(content) { + return ALL_BOTS_RE.test(content); +} +// --------------------------------------------------------------------------- +// Main extension +// --------------------------------------------------------------------------- +export default function (pi) { + const config = loadConfig(); + + // Exit diagnostics — log why the process is terminating + const _exitSignals = ["SIGINT", "SIGTERM", "SIGHUP", "SIGQUIT"]; + for (const sig of _exitSignals) { + process.on(sig, () => { + console.error(`[zulip-ext] EXIT-DIAG: received ${sig} — PM2 or external stop signal`); + }); + } + process.on("beforeExit", (code) => { + console.error(`[zulip-ext] EXIT-DIAG: beforeExit with code=${code}`); + }); + process.on("exit", (code) => { + // Synchronous-only here — last chance to log + fs.writeSync(2, `[zulip-ext] EXIT-DIAG: exit with code=${code}\n`); + }); + process.on("uncaughtException", (err) => { + console.error(`[zulip-ext] EXIT-DIAG: UNCAUGHT EXCEPTION: ${err.message}\n${err.stack}`); + }); + process.on("unhandledRejection", (reason) => { + console.error(`[zulip-ext] EXIT-DIAG: UNHANDLED REJECTION: ${reason instanceof Error ? reason.message : String(reason)}`); + }); + // Track if stdin end was the cause + if (process.stdin.isTTY) { + process.stdin.on("end", () => { + console.error(`[zulip-ext] EXIT-DIAG: stdin ended — RPC mode will shutdown`); + }); + } + + let queue = null; + let pollTimer = null; + let retryCount = 0; + let connected = false; + let lastError = null; + let processedCount = 0; + let healthServer = null; + // Stuck detection: track last message activity time + let lastActivityTime = Date.now(); + const STUCK_THRESHOLD_MS = 30 * 60 * 1000; // 30 min of inactivity = stuck + let stuckRecoveryTimer = null; + const STUCK_CHECK_INTERVAL_MS = 60 * 1000; // check every 60s + // Track state for health endpoint + function getState() { + const now = Date.now(); + const idleSeconds = Math.floor((now - lastActivityTime) / 1000); + const isStuck = connected && idleSeconds > (STUCK_THRESHOLD_MS / 1000); + const passedChecks = [ + { name: "queue", status: connected && queue !== null }, + { name: "bot_identity", status: (queue?.botUserId ?? 0) > 0 }, + { name: "all_bots", status: (queue?.allBotsUserId ?? 0) > 1 }, + { name: "poll_loop", status: connected && !isStuck }, + { name: "echo_prevention", status: BOT_EMAILS.length > 0 }, + ]; + const checkResults = {}; + for (const c of passedChecks) checkResults[c.name] = c.status ? "ok" : "fail"; + const okChecks = passedChecks.filter(c => c.status).length; + return { + // Standardized: zulip-health/v1 + status: connected ? (okChecks >= passedChecks.length - 1 ? "ok" : "degraded") : "down", + platform: "pi", + agent: config.agent.name, + zulip: { + connected, + site: config.zulip.site, + email: config.zulip.email, + queue_id: queue?.queueId ?? null, + last_event_id: queue?.lastEventId ?? null, + bot_user_id: queue?.botUserId ?? null, + all_bots_user_id: queue?.allBotsUserId ?? null, + messages_processed: processedCount, + silence_seconds: idleSeconds, + stuck: isStuck, + }, + uptime: {}, + checks: checkResults, + // Legacy flat fields (backward compat) + connected, + stuck: isStuck, + idle_seconds: idleSeconds, + email: config.zulip.email, + site: config.zulip.site, + agent: config.agent.name, + is_owner: config.agent.owner_email, + queue_id: queue?.queueId ?? null, + last_event_id: queue?.lastEventId ?? null, + bot_user_id: queue?.botUserId ?? null, + all_bots_user_id: queue?.allBotsUserId ?? null, + messages_processed: processedCount, + last_error: lastError, + retry_count: retryCount, + last_activity_time: lastActivityTime, + }; + } + // ---- Gen 2: Known bot emails to prevent echo loops ---- + const BOT_EMAILS = [ + config.zulip.email, // own email + "tanko-bot@chat.sysloggh.net", + "mumuni-bot@chat.sysloggh.net", + "koby-bot@chat.sysloggh.net", + "koonimo-bot@chat.sysloggh.net", + "kagentz-bot@chat.sysloggh.net", + ]; + const OWNER_EMAIL = config.agent.owner_email || "jerome@sysloggh.com"; + let lastNonBotMessageTime = 0; + let lastHeartbeatTime = 0; + const HEARTBEAT_INTERVAL_MS = 300000; // 5 min + + // Pending Zulip replies awaiting the LLM response from agent_end. + // zulipMessageId is set after the placeholder is sent. + const pendingZulipReplies = []; + let zulipReplyIdCounter = 0; + let isAgentBusy = false; + let streamingTimer = null; + let streamingAccumulator = ""; + let streamingReplyId = null; + // Token tracking for /status context display + let sessionTokens = { input: 0, output: 0, cacheRead: 0, totalCost: 0 }; + let totalMessagesSkipped = 0; + // Track last user message for /retry command + let lastUserMessage = null; + /** + * Zulip commands available via DM. + */ + const ZULIP_COMMANDS = { + "status": { + description: "Show system status — PM2 processes, containers, uptime", + usage: "/status", + }, + "retry": { + description: "Retry the last response — re-injects your previous message", + usage: "/retry", + }, + "continue": { + description: "Continue the last response if it was cut off", + usage: "/continue", + }, + "new": { + description: "Start fresh — clears conversation context for a new topic", + usage: "/new [topic]", + }, + "help": { + description: "Show this help message", + usage: "/help [command]", + }, + }; + /** + * Handle a Zulip command — sends response directly back to Zulip without LLM. + * Returns true if the message was handled as a command, false to pass through. + */ + function handleZulipCommand(content, senderName, senderId, replyInfo) { + const cmdMatch = content.match(/^\/(\w+)(?:\s+(.+))?$/); + if (!cmdMatch) + return false; + const cmd = cmdMatch[1].toLowerCase(); + const args = cmdMatch[2] || ""; + let response = ""; + const target = replyInfo.type === "private" + ? { type: "private", to: senderId } + : { type: "stream", to: replyInfo.streamId ?? replyInfo.streamName, subject: replyInfo.topic }; + switch (cmd) { + case "status": { + response = `**Abiba System Status**\n\n`; + try { + // PM2 status + response += `**PM2 Processes:**\n`; + const pm2 = execSync(`pm2 status --no-color 2>/dev/null`, { timeout: 10000, encoding: "utf-8" }); + pm2.split("\n").forEach(line => { + if (line.includes("abiba-")) { + const cells = line.split("│").filter(c => c.trim()); + const name = cells[1] || "?"; + const status = cells[5] || "?"; + const uptime = cells[3] || "?"; + const restarts = cells[4] || "0"; + response += `• ${name.replace(/\[\d+m|\[\d{2}m|\[m|\[7m|\[1m|\[36m|\[32m|\[39m|\[22m|\[90m/g,"")}: ${status.replace(/\[\d+m|\[\d{2}m|\[m|\[7m|\[1m|\[36m|\[32m|\[39m|\[22m|\[90m/g,"")} (uptime: ${uptime}, restarts: ${restarts})\n`; + } + }); + + // LiteLLM containers + response += `\n**LiteLLM Containers:**\n`; + const containers = execSync( + `/usr/bin/ssh -o StrictHostKeyChecking=no -o ConnectTimeout=5 -i /root/.ssh/id_ed25519 192.168.68.116 "docker ps --format '{{.Names}}|{{.Status}}'" 2>&1`, + { timeout: 15000, encoding: "utf-8" } + ); + containers.trim().split("\n").forEach(line => { + const [name, status] = line.split("|"); + if (name) response += `• ${name}: ${status}\n`; + }); + + // Session info (in-memory tracked) + response += `\n**Session:**\n`; + response += `• deepseek/deepseek-v4-flash\n`; + const totalTok = sessionTokens.input + sessionTokens.output; + response += `• ↑${sessionTokens.input.toLocaleString()} ↓${sessionTokens.output.toLocaleString()} CH${sessionTokens.cacheRead.toLocaleString()} R${totalTok.toLocaleString()}\n`; + if (sessionTokens.totalCost > 0) { + response += `• Cost: $${sessionTokens.totalCost.toFixed(4)}\n`; + } + try { + const zulipLine = execSync(`pm2 status --no-color 2>/dev/null | grep abiba-zulip | head -1`, + { timeout: 5000, encoding: "utf-8" }).trim(); + const cells = zulipLine.split("│").filter(c => c.trim()); + if (cells.length >= 6) { + const uptime = cells[3].replace(/\[\d+m|\[\d{2}m|\[m|\[7m|\[1m|\[36m|\[32m|\[39m|\[22m|\[90m/g,"").trim(); + response += `• Uptime: ${uptime}\n`; + } + } catch {} + } catch (err) { + response += `\n_Error fetching status: ${err.message}_`; + } + queue.sendMessage({ ...target, content: response }).catch(() => {}); + return true; + } + case "retry": { + if (lastUserMessage) { + // Re-inject the last user message + const label = replyInfo.type === "stream" + ? `@all-bots in #${replyInfo.streamName}` + : `DM from @${senderName}`; + response = `:repeat: Retrying last message...`; + queue.sendMessage({ ...target, content: response }).catch(() => { }); + queue.sendTypingNotification([senderId], "start").catch(() => { }); + if (isAgentBusy) { + pi.sendUserMessage(`[Zulip ${label}]: ${lastUserMessage}`, { deliverAs: "steer" }); + } + else { + pi.sendUserMessage(`[Zulip ${label}]: ${lastUserMessage}`); + } + console.log(`[zulip-ext] /retry: re-injecting "${lastUserMessage.slice(0, 60)}..."`); + } + else { + response = ":warning: No previous message to retry. Send me something first!"; + queue.sendMessage({ ...target, content: response }).catch(() => { }); + } + return true; + } + case "continue": { + const label = replyInfo.type === "stream" + ? `@all-bots in #${replyInfo.streamName}` + : `DM from @${senderName}`; + response = `:arrow_forward: Continuing...`; + queue.sendMessage({ ...target, content: response }).catch(() => { }); + queue.sendTypingNotification([senderId], "start").catch(() => { }); + if (isAgentBusy) { + pi.sendUserMessage(`[Zulip ${label}]: Continue your response from where you left off.`, { deliverAs: "steer" }); + } + else { + pi.sendUserMessage(`[Zulip ${label}]: Continue your response from where you left off.`); + } + console.log(`[zulip-ext] /continue: sent continue prompt`); + return true; + } + case "new": { + const label = replyInfo.type === "stream" + ? `@all-bots in #${replyInfo.streamName}` + : `DM from @${senderName}`; + const topicMsg = args + ? `Start fresh, focusing on: ${args}` + : `Start fresh.`; + response = `:broom: ${topicMsg}\n\n_Previous context cleared._`; + queue.sendMessage({ ...target, content: response }).catch(() => { }); + // Clear stored message so /retry won't reuse old context + lastUserMessage = null; + // Inject a reset message to the LLM + const resetMsg = `[System: Session reset requested by user. ${topicMsg} Treat this as a fresh conversation. Ignore all prior context.]`; + if (isAgentBusy) { + pi.sendUserMessage(resetMsg, { deliverAs: "steer" }); + } + else { + pi.sendUserMessage(resetMsg); + } + console.log(`[zulip-ext] /new: session reset`); + return true; + } + case "help": { + if (args) { + // Help for specific command + const specificCmd = args.toLowerCase(); + if (ZULIP_COMMANDS[specificCmd]) { + const c = ZULIP_COMMANDS[specificCmd]; + response = `**/${specificCmd}** — ${c.description}\n\nUsage: \`${c.usage}\``; + } + else { + response = `:question: Unknown command \`/${specificCmd}\`. Try \`/help\` to see all commands.`; + } + } + else { + response = `**Abiba Zulip Commands**\n\n` + + Object.entries(ZULIP_COMMANDS) + .map(([name, c]) => `• \`/${name}\` — ${c.description}`) + .join("\n") + + `\n\nSend any message to chat. Commands start with \`/\`.`; + } + queue.sendMessage({ ...target, content: response }).catch(() => { }); + return true; + } + default: { + response = `:question: Unknown command \`/${cmd}\`. Try \`/help\` to see available commands.`; + queue.sendMessage({ ...target, content: response }).catch(() => { }); + return true; + } + } + return true; + } + /** + * Queue a Zulip message into the current pi session via sendUserMessage + * and schedule the reply to be sent back on the next agent_end event. + * Sends a "Thinking..." placeholder to Zulip immediately for streaming UX. + */ + function injectIntoSession(msg, senderName, senderId, replyInfo) { + // Send typing indicator (fire-and-forget) + queue.sendTypingNotification([senderId], "start").catch(() => { }); + const replyId = ++zulipReplyIdCounter; + const label = replyInfo.type === "stream" + ? `@all-bots in #${replyInfo.streamName}` + : `DM from @${senderName}`; + // Send a placeholder "Thinking..." message to Zulip immediately + const placeholders = [ + ":robot: _Processing your message..._", + ":hourglass_flowing_sand: _Thinking..._", + ":brain: _Generating response..._", + ]; + const placeholder = placeholders[replyId % placeholders.length]; + const placeholderTarget = replyInfo.type === "private" + ? { type: "private", to: senderId } + : { type: "stream", to: replyInfo.streamId ?? replyInfo.streamName, subject: replyInfo.topic }; + queue.sendMessage({ ...placeholderTarget, content: placeholder }) + .then((msgId) => { + pendingZulipReplies.push({ id: replyId, ...replyInfo, zulipMessageId: msgId }); + console.log(`[zulip-ext] Placed [${replyId}] (msg ${msgId}) for ${senderName}`); + }) + .catch(() => { + pendingZulipReplies.push({ id: replyId, ...replyInfo }); + console.log(`[zulip-ext] Queued [${replyId}] for ${senderName} (no placeholder)`); + }); + // Inject into the pi session — use steer mode if agent is busy + if (isAgentBusy) { + pi.sendUserMessage(`[Zulip ${label}]: ${msg.content}`, { deliverAs: "steer" }); + console.log(`[zulip-ext] Steer [${replyId}]: queued while agent busy`); + } + else { + pi.sendUserMessage(`[Zulip ${label}]: ${msg.content}`); + console.log(`[zulip-ext] Injected [${replyId}]: from ${senderName}`); + } + } + /** + * Process a single Zulip event — injects DMs and @all-bots into the + * current pi session instead of spawning a subprocess. + */ + /** + * Gen 2: Check if sender is a known bot email to prevent echo loops. + */ + function isBotSender(senderEmail) { + return BOT_EMAILS.includes(senderEmail); + } + + async function processEvent(event) { + const msg = event.message; + // Ignore non-message events + if (event.type !== "message") + return; + // Update last activity time on any message event + lastActivityTime = Date.now(); + // Ignore own messages (sent by this bot) + if (msg.sender_email === config.zulip.email) + return; + // ---- Gen 2: Echo loop prevention ---- + if (isBotSender(msg.sender_email)) { + totalMessagesSkipped++; + if (totalMessagesSkipped % 10 === 1) { + console.log(`[zulip-ext] Skipped ${totalMessagesSkipped} bot msgs from ${msg.sender_email} (echo loop prevention)`); + } + return; + } + // ---- Heartbeat logging ---- + const now = Date.now(); + if (now - lastHeartbeatTime > HEARTBEAT_INTERVAL_MS) { + lastHeartbeatTime = now; + const queueInfo = queue ? `queue=${queue.queueId}` : "no-queue"; + console.log(`[zulip-ext] Heartbeat — processed=${processedCount} skipped=${totalMessagesSkipped} pending=${pendingZulipReplies.length} ${queueInfo}`); + } + // DM-first architecture (ADR-001, ADR-002): process private messages + if (msg.type === "private") { + processedCount++; + const senderName = msg.sender_full_name || msg.sender_email; + const senderId = msg.sender_id; + const isOwner = msg.sender_email === OWNER_EMAIL; + console.log(`[zulip-ext] DM from ${senderName} (id=${senderId}${isOwner ? ', OWNER' : ''}): ${msg.content.slice(0, 80)}...`); + // Check for Zulip commands (/retry, /continue, /help) + if (msg.content.startsWith("/")) { + if (handleZulipCommand(msg.content, senderName, senderId, { type: "private", senderId, senderName })) { + return; + } + } + // Store non-command message for /retry + lastUserMessage = msg.content; + lastNonBotMessageTime = Date.now(); + injectIntoSession(msg, senderName, senderId, { type: "private", senderId, senderName }); + return; + } + // Stream messages: respond to @mention of this bot or @all-bots (ADR-005, ADR-006) + const mentionedUsers = msg.mentioned_users || []; + const mentionedUserIds = mentionedUsers.map(u => u.user_id); + const isDirectMention = queue && queue.botUserId && mentionedUserIds.includes(queue.botUserId); + const isAllBots = isAllBotsMention(msg.content); + if (msg.type === "stream" && (isDirectMention || isAllBots)) { + processedCount++; + const streamName = typeof msg.display_recipient === "string" + ? msg.display_recipient + : "unknown"; + const topic = msg.subject || "general"; + const mentionType = isDirectMention ? "@mention" : "@all-bots"; + console.log(`[zulip-ext] ${mentionType} in #${streamName} > ${topic}`); + // Check for Zulip commands in stream too + if (msg.content.startsWith("/")) { + if (handleZulipCommand(msg.content, msg.sender_full_name || msg.sender_email, msg.sender_id, { + type: "stream", + streamName, + topic, + streamId: msg.stream_id, + })) { + return; + } + } + lastUserMessage = msg.content; + injectIntoSession(msg, msg.sender_full_name || msg.sender_email, msg.sender_id, { + type: "stream", + streamName, + topic, + streamId: msg.stream_id, + }); + } + } + /** + * Polling loop. + */ + async function startPolling() { + console.log(`[zulip-ext] Connecting to ${config.zulip.site} as ${config.zulip.email}...`); + try { + queue = await createZulipQueue(config); + connected = true; + retryCount = 0; + lastError = null; + console.log(`[zulip-ext] Connected, queue: ${queue.queueId} (last_event_id=${queue.lastEventId})`); + // Start polling + pollTimer = setInterval(async () => { + if (!queue) + return; + try { + const events = await queue.poll(); + lastError = null; // Clear stale errors on successful poll + // Update last activity time on any event + if (events.length > 0) { + lastActivityTime = Date.now(); + } + for (const event of events) { + await processEvent(event); + } + // Check if we're stuck: connected but no activity for >30 min + // Force queue re-registration to recover from silent expiry + if (connected && events.length === 0) { + const idleMs = Date.now() - lastActivityTime; + if (idleMs > STUCK_THRESHOLD_MS && queue) { + console.log(`[zulip-ext] STUCK: no events for ${Math.floor(idleMs/60000)}min — re-registering queue`); + connected = false; + if (pollTimer) clearInterval(pollTimer); + pollTimer = null; + setTimeout(() => startPolling(), 1000); + return; + } + } + } + catch (err) { + const errMsg = err instanceof Error ? err.message : String(err); + // Check if queue was deregistered (need to reconnect) — Gen 3: Broader detection + const isQueueExpired = errMsg.includes("BAD_EVENT_QUEUE_ID") || + errMsg.includes("queue_id") || + errMsg.includes("deregistered") || + errMsg.includes("Event queue") || + errMsg.includes("invalid queue"); + if (isQueueExpired) { + console.log(`[zulip-ext] Queue expired, reconnecting... (${errMsg.slice(0, 60)})`); + connected = false; + if (pollTimer) + clearInterval(pollTimer); + pollTimer = null; + setTimeout(() => startPolling(), config.retry_delay_ms ?? 5000); + return; + } + // Network errors — retry on next poll interval + lastError = errMsg; + console.error(`[zulip-ext] Poll error: ${errMsg}`); + } + }, config.poll_interval_ms ?? 3000); + } + catch (err) { + const errMsg = err instanceof Error ? err.message : String(err); + lastError = errMsg; + console.error(`[zulip-ext] Connection failed: ${errMsg}`); + if (retryCount < (config.max_retries ?? 3)) { + retryCount++; + const delay = (config.retry_delay_ms ?? 5000) * retryCount; + console.log(`[zulip-ext] Retrying in ${delay / 1000}s (attempt ${retryCount})...`); + setTimeout(() => startPolling(), delay); + } + else { + console.error(`[zulip-ext] Max retries (${config.max_retries}) reached. Giving up.`); + lastError = `Connection failed after ${config.max_retries} retries: ${errMsg}`; + } + } + } + // ------------------------------------------------------------------------- + // LiteLLM self-heal health check + // ------------------------------------------------------------------------- + /** + * Run health checks against the LiteLLM deployment. + * Returns a report object with issues_found, issues_fixed, etc. + * Only reports if something is wrong — silent on healthy. + */ + async function performHealthCheck(cfg) { + const backendHost = "192.168.68.116"; + const publicUrl = "https://litellm.sysloggh.net"; + const authHeader = "Basic " + Buffer.from(`${cfg.zulip.email}:${cfg.zulip.api_key}`).toString("base64"); + const report = { + checks_passed: 0, + checks_failed: 0, + issues_found: 0, + issues_fixed: 0, + issues_escalated: 0, + actions: [], + }; + + // Check 1: Public endpoints + const endpoints = [ + { path: "/ui/", name: "Admin UI" }, + { path: "/docs", name: "Swagger Docs" }, + { path: "/openapi.json", name: "OpenAPI" }, + { path: "/redoc", name: "ReDoc" }, + ]; + for (const ep of endpoints) { + try { + const res = await fetch(`${publicUrl}${ep.path}`, { method: "HEAD", signal: AbortSignal.timeout(10000) }); + if (res.ok) { + report.checks_passed++; + } else { + report.checks_failed++; + report.issues_found++; + report.issues_escalated++; + report.actions.push({ target: ep.name, action: "endpoint down", result: `HTTP ${res.status}` }); + } + } catch { + report.checks_failed++; + report.issues_found++; + report.issues_escalated++; + report.actions.push({ target: ep.name, action: "endpoint unreachable", result: "connection failed" }); + } + } + + // Check 2: Backend containers via SSH + try { + const sshBase = `/usr/bin/ssh -o StrictHostKeyChecking=no -o ConnectTimeout=5 -i /root/.ssh/id_ed25519`; + const cmd = `${sshBase} ${backendHost} "docker ps --format '{{.Names}}|{{.Status}}'" 2>&1`; + const containerList = execSync(cmd, { timeout: 20000, encoding: "utf-8", shell: "/bin/bash" }).trim(); + + const expected = ["harness-litellm", "harness-nginx", "harness-router", "harness-postgres", "harness-redis", "harness-dashboard"]; + const running = containerList.split("\n").filter(l => l.includes("(healthy)")).map(l => l.split("|")[0]); + const missing = expected.filter(e => !running.includes(e)); + + if (missing.length === 0) { + report.checks_passed++; + } else { + report.checks_failed++; + report.issues_found += missing.length; + for (const m of missing) { + // Try to fix: restart the container + try { + execSync( + `/usr/bin/ssh -o StrictHostKeyChecking=no -o ConnectTimeout=10 -i /root/.ssh/id_ed25519 ${backendHost} docker compose -f /opt/inference-harness/docker-compose.yml up -d ${m} 2>&1`, + { timeout: 30000, encoding: "utf-8" } + ); + // Wait and verify + await new Promise(r => setTimeout(r, 5000)); + const verifyOut = execSync( + `/usr/bin/ssh -o StrictHostKeyChecking=no -o ConnectTimeout=5 -i /root/.ssh/id_ed25519 ${backendHost} docker ps --filter name=${m} --format '{{.Status}}' 2>&1`, + { timeout: 10000, encoding: "utf-8" } + ).trim(); + + if (verifyOut.includes("healthy") || verifyOut.includes("Up")) { + report.issues_fixed++; + report.actions.push({ target: m, action: "restarted container", result: "healthy" }); + } else { + report.issues_escalated++; + report.actions.push({ target: m, action: "restart failed", result: verifyOut || "still down" }); + } + } catch (fixErr) { + report.issues_escalated++; + report.actions.push({ target: m, action: "restart error", result: fixErr.message }); + } + } + } + } catch (sshErr) { + report.checks_failed++; + report.issues_found++; + report.issues_escalated++; + report.actions.push({ target: "backend_host", action: "SSH failed", result: sshErr.message }); + } + + // Gitea auto-commit: if fixes were applied, commit and push + if (report.issues_fixed > 0) { + try { + const sshBase = `/usr/bin/ssh -o StrictHostKeyChecking=no -o ConnectTimeout=10 -i /root/.ssh/id_ed25519`; + const repoDir = "/opt/inference-harness"; + + // Build commit message from actions + let commitMsg = `auto-fix: \\`; + for (const a of report.actions) { + if (a.result === "healthy" || a.result?.includes("ok")) { + commitMsg += `${a.target}: ${a.action} — ${a.result}\\n`; + } + } + + // Run git commands via SSH + const gitCmds = [ + `cd ${repoDir} && git add -A`, + `cd ${repoDir} && git commit -m "${commitMsg.trim()}" 2>/dev/null || true`, + `cd ${repoDir} && git pull --no-rebase 2>/dev/null || true`, + `cd ${repoDir} && git push 2>&1`, + ]; + + for (const cmd of gitCmds) { + execSync(`${sshBase} ${backendHost} ${cmd}`, { timeout: 30000, encoding: "utf-8" }); + } + + const gitLog = execSync( + `${sshBase} ${backendHost} "cd ${repoDir} && git log --oneline -1"`, + { timeout: 10000, encoding: "utf-8" } + ).trim(); + + console.log(`[zulip-ext] Auto-commit pushed: ${gitLog}`); + report.git_commit = gitLog; + } catch (gitErr) { + console.error(`[zulip-ext] Auto-commit failed: ${gitErr.message}`); + } + } + + return report; + } + + // ------------------------------------------------------------------------- + // Extension lifecycle + // ------------------------------------------------------------------------- + pi.on("session_start", async (_event) => { + // Guard: only activate in the designated PM2 RPC process (ZULIP_EXTENSION_ACTIVE=true) + // Prevents double responses when TUI and RPC both load this extension + if (process.env.ZULIP_EXTENSION_ACTIVE !== "true") { + console.log(`[zulip-ext] Skipped — set ZULIP_EXTENSION_ACTIVE=true to enable (TUI vs RPC dedup)`); + return; + } + console.log(`[zulip-ext] Starting Zulip extension for ${config.agent.name} (${config.zulip.email})`); + if (!healthServer) { + healthServer = startHealthServer(config.health_port ?? 9200, getState); + } + startPolling(); + + // Schedule LiteLLM self-heal check every 4 hours + // Only DMs the user when issues are found/fixed — silent on healthy + const SELF_HEAL_INTERVAL_MS = 4 * 60 * 60 * 1000; // 4 hours + let healCycleCount = 0; + let issuesFixedSinceDigest = 0; + let digestsSentToday = 0; + let lastDigestDate = ""; + + async function runSelfHeal() { + healCycleCount++; + let report = null; + + try { + report = await performHealthCheck(config); + } catch (err) { + console.error(`[zulip-ext] Self-heal error: ${err.message}`); + return; + } + + // Jerome's Zulip user ID for DMs + const ownerZulipId = "9"; + + // Always log to console + console.log(`[zulip-ext] Self-heal #${healCycleCount}: ${report.checks_passed}/${report.checks_passed + report.checks_failed} passed, ${report.issues_fixed} fixed`); + + // Track for daily digest + issuesFixedSinceDigest += report.issues_fixed; + + // DM only if something happened + if (report.issues_found > 0 && queue) { + issuesFixedSinceDigest += report.issues_fixed; + + // Build DM message + let dmContent = report.issues_fixed > 0 + ? `:hammer_and_wrench: **LiteLLM Self-Heal — Fix Applied**\n\nFound ${report.issues_found} issue(s), fixed ${report.issues_fixed}, escalated ${report.issues_escalated}.` + : `:warning: **LiteLLM Self-Heal — Needs Attention**\n\nFound ${report.issues_found} issue(s) that could not be auto-fixed.`; + + if (report.actions && report.actions.length > 0) { + dmContent += `\n\n**Actions:**\n`; + for (const a of report.actions) { + dmContent += `• ${a.target}: ${a.action} — ${a.result}\n`; + } + } + + dmContent += `\\n\\n_Cycle #${healCycleCount} — next check in 4h_`; + + // Send DM to owner + try { + const formBody = new URLSearchParams(); + formBody.append("type", "private"); + formBody.append("to", JSON.stringify([ownerZulipId])); + formBody.append("content", dmContent); + + await fetch(`${config.zulip.site}/api/v1/messages`, { + method: "POST", + headers: { + Authorization: "Basic " + Buffer.from(`${config.zulip.email}:${config.zulip.api_key}`).toString("base64"), + "Content-Type": "application/x-www-form-urlencoded", + }, + body: formBody.toString(), + }); + console.log(`[zulip-ext] Self-heal alert sent (user_id=${ownerZulipId})`); + } catch (dmErr) { + console.error(`[zulip-ext] Failed to send self-heal DM: ${dmErr.message}`); + } + } + + // Send daily digest at end of day + const today = new Date().toISOString().slice(0, 10); + if (today !== lastDigestDate) { + // First run of a new day — send previous day's digest + if (lastDigestDate && issuesFixedSinceDigest > 0 && queue) { + const digestContent = `:bar_chart: **LiteLLM Self-Heal — Daily Digest (${lastDigestDate})**\n\n` + + `Total cycles: ${healCycleCount}\\n` + + `Issues fixed: ${issuesFixedSinceDigest}\\n` + + `Digests sent: ${digestsSentToday}`; + + try { + const formBody = new URLSearchParams(); + formBody.append("type", "private"); + formBody.append("to", JSON.stringify([ownerZulipId])); + formBody.append("content", digestContent); + await fetch(`${config.zulip.site}/api/v1/messages`, { + method: "POST", + headers: { + Authorization: "Basic " + Buffer.from(`${config.zulip.email}:${config.zulip.api_key}`).toString("base64"), + "Content-Type": "application/x-www-form-urlencoded", + }, + body: formBody.toString(), + }); + console.log(`[zulip-ext] Daily digest sent`); + } catch (digErr) { + console.error(`[zulip-ext] Digest send failed: ${digErr.message}`); + } + } + lastDigestDate = today; + issuesFixedSinceDigest = 0; + digestsSentToday++; + } + } + + // Run first check after 30 seconds (give queue time to stabilize) + setTimeout(async () => { + await runSelfHeal(); + // Then every 4 hours + setInterval(runSelfHeal, SELF_HEAL_INTERVAL_MS); + }, 30000); + + console.log(`[zulip-ext] Self-heal scheduled every 4 hours (first check in 30s)`); + }); + // ----------------------------------------------------------------------- + // Streaming + steering lifecycle + // ----------------------------------------------------------------------- + /** + * turn_start fires when a new LLM turn begins. Mark the agent as busy + * so subsequent Zulip DMs are steered instead of starting new turns. + */ + pi.on("turn_start", async () => { + isAgentBusy = true; + streamingAccumulator = ""; + streamingReplyId = pendingZulipReplies.length > 0 ? pendingZulipReplies[0].id : null; + }); + /** + * message_update fires on each streaming token. Accumulate partial text + * and periodically edit the Zulip placeholder message so the user sees + * the response come in live. + */ + pi.on("message_update", async (event) => { + const msg = event.message; + const content = msg.content; + let partialText = ""; + if (typeof content === "string") { + partialText = content; + } + else if (Array.isArray(content)) { + partialText = content + .filter((c) => c.type === "text") + .map((c) => c.text) + .join("\n"); + } + streamingAccumulator = partialText; + if (!streamingTimer && streamingAccumulator.length > 10) { + streamingTimer = setTimeout(async () => { + streamingTimer = null; + if (!streamingAccumulator || streamingReplyId === null) + return; + const pending = pendingZulipReplies.find((r) => r.id === streamingReplyId); + if (!pending || !pending.zulipMessageId) + return; + try { + const MAX_PREVIEW = 9500; + const preview = streamingAccumulator.length > MAX_PREVIEW + ? streamingAccumulator.slice(0, MAX_PREVIEW) + "\n\n_… still generating…_" + : streamingAccumulator; + await queue.editMessage(pending.zulipMessageId, preview); + } + catch { + // Non-critical during streaming + } + }, 2000); + } + }); + /** + * agent_end captures the final LLM response and relays it back to Zulip. + * If we sent a placeholder message earlier, edit it with the final content + * (instead of sending a new message) for a seamless streaming experience. + */ + pi.on("agent_end", async (event) => { + isAgentBusy = false; + // Track token usage from the LLM response + if (event.messages) { + for (const m of event.messages) { + if (m.usage) { + sessionTokens.input += m.usage.input || 0; + sessionTokens.output += m.usage.output || 0; + sessionTokens.cacheRead += m.usage.cacheRead || 0; + if (m.usage.cost && m.usage.cost.total) { + sessionTokens.totalCost += m.usage.cost.total; + } + } + } + } + if (streamingTimer) { + clearTimeout(streamingTimer); + streamingTimer = null; + } + if (pendingZulipReplies.length === 0) + return; + const reply = pendingZulipReplies.shift(); + // Extract final assistant response text + const allMsgs = event.messages; + const assistantMsgs = allMsgs.filter((m) => m.role === "assistant"); + const lastAssistant = assistantMsgs[assistantMsgs.length - 1]; + let responseText = ""; + if (lastAssistant) { + const content = lastAssistant.content; + if (typeof content === "string") { + responseText = content; + } + else if (Array.isArray(content)) { + responseText = content + .filter((c) => c.type === "text") + .map((c) => c.text) + .join("\n"); + } + } + // Stop typing indicator + if (reply.senderId) { + queue.sendTypingNotification([reply.senderId], "stop").catch(() => { }); + } + if (!responseText.trim()) { + console.log(`[zulip-ext] No assistant response for [${reply.id}] from ${reply.senderName}`); + return; + } + const MAX_ZULIP_MSG = 10000; + const truncated = responseText.length > MAX_ZULIP_MSG + ? responseText.slice(0, MAX_ZULIP_MSG) + "\n\n[...truncated at Zulip limit, see session for full output]" + : responseText; + try { + if (reply.zulipMessageId) { + // Edit the placeholder with the final response (streaming UX) + await queue.editMessage(reply.zulipMessageId, truncated); + console.log(`[zulip-ext] Finalized [${reply.id}] for ${reply.senderName} (${truncated.length} chars)`); + } + else { + // No placeholder — send as new message + const target = reply.type === "private" + ? { type: "private", to: reply.senderId } + : { type: "stream", to: reply.streamId ?? reply.streamName, subject: reply.topic }; + await queue.sendMessage({ ...target, content: truncated }); + console.log(`[zulip-ext] Replied to [${reply.id}] ${reply.senderName} (${truncated.length} chars)`); + } + } + catch (err) { + const errMsg = err instanceof Error ? err.message : String(err); + console.error(`[zulip-ext] Failed to finalize reply [${reply.id}]: ${errMsg}`); + lastError = errMsg; + } + }); + pi.on("session_shutdown", async () => { + console.log(`[zulip-ext] Shutting down...`); + if (streamingTimer) { + clearTimeout(streamingTimer); + streamingTimer = null; + } + if (pollTimer) { + clearInterval(pollTimer); + pollTimer = null; + } + if (healthServer) { + healthServer.close(); + healthServer = null; + } + connected = false; + queue = null; + }); + // Register /zulip-status command for diagnostics + pi.registerCommand("zulip-status", { + description: "Show Zulip extension status and stats", + handler: async (_args, ctx) => { + const state = getState(); + const lines = [ + "=== Zulip Extension Status ===", + "", + `Agent: ${state.agent}`, + `Email: ${state.email}`, + `Server: ${state.site}`, + `Owner: ${state.is_owner}`, + `Connected: ${state.connected ? "✅ Yes" : "❌ No"}`, + `Queue ID: ${state.queue_id ?? "N/A"}`, + `Last Event ID: ${state.last_event_id ?? "N/A"}`, + `Messages: ${state.messages_processed} processed`, + `Retries: ${state.retry_count}`, + `Health Port: :${config.health_port ?? 9200}`, + ]; + if (state.last_error) { + lines.push("", `Last Error: ${state.last_error}`); + } + lines.push("", "ADRs followed: DM-first (ADR-001), DM-routing (ADR-002), @all-bots content (ADR-006)"); + ctx.ui.notify(lines.join("\n"), "info"); + }, + }); + // Register zulip-send-test command + pi.registerCommand("zulip-send-test", { + description: "Send a test message to Zulip user: /zulip-send-test ", + handler: async (args, ctx) => { + const parts = args.trim().match(/^([^\s]+)\s+(.+)$/); + const targetEmail = parts?.[1] ?? config.agent.owner_email; + const message = parts?.[2] ?? "Test message from pi Zulip extension"; + if (!queue) { + ctx.ui.notify("Zulip queue not connected. Use /zulip-status to check.", "error"); + return; + } + try { + await queue.sendMessage({ + type: "private", + to: targetEmail, + content: message, + }); + ctx.ui.notify(`Test message sent to ${targetEmail}`, "info"); + } + catch (err) { + const errMsg = err instanceof Error ? err.message : String(err); + ctx.ui.notify(`Failed to send: ${errMsg}`, "error"); + } + }, + }); + // Register /zulip-test command — comprehensive self-test + pi.registerCommand("zulip-test", { + description: "Run comprehensive Zulip self-test (queue, identity, event reception)", + handler: async (_args, ctx) => { + const results = []; + const startTime = Date.now(); + + // Check 1: Queue connected + results.push({ + name: "queue_connected", + status: connected && queue !== null, + detail: connected ? `Queue: ${queue?.queueId ?? "?"}` : "Not connected", + }); + + // Check 2: Bot identity + results.push({ + name: "bot_identity", + status: queue?.botUserId !== null && queue?.botUserId !== undefined, + detail: `Bot user_id: ${queue?.botUserId ?? "unresolved"}`, + }); + + // Check 3: @all-bots resolution + results.push({ + name: "all_bots_resolved", + status: queue?.allBotsUserId !== null && queue?.allBotsUserId !== undefined && queue?.allBotsUserId > 1, + detail: `@all-bots user_id: ${queue?.allBotsUserId ?? "unresolved"} (config default was 1)`, + }); + + // Check 4: Health endpoint + results.push({ + name: "health_endpoint", + status: true, + detail: `Port :${config.health_port ?? 9200} responding`, + }); + + // Check 5: Poll loop active + const idleS = Math.floor((Date.now() - lastActivityTime) / 1000); + const isStuck = connected && idleS > (30 * 60); + results.push({ + name: "poll_loop", + status: connected && !isStuck, + detail: connected + ? `Active, last activity ${idleS}s ago${isStuck ? " (STUCK)" : ""}` + : "Not connected", + }); + + // Check 6: Echo prevention + results.push({ + name: "echo_prevention", + status: BOT_EMAILS.length > 0, + detail: `${BOT_EMAILS.length} bot emails tracked for echo prevention`, + }); + + // Check 7: Send loopback DM to verify API works + if (queue && connected) { + try { + const loopbackId = await queue.sendMessage({ + type: "private", + to: config.agent.owner_email, + content: `🔬 Zulip self-test completed at ${new Date().toISOString()}`, + }); + results.push({ + name: "api_send", + status: true, + detail: `Loopback DM sent (msg_id=${loopbackId})`, + }); + } catch (err) { + results.push({ + name: "api_send", + status: false, + detail: `Failed: ${err.message?.slice(0, 60)}`, + }); + } + } else { + results.push({ + name: "api_send", + status: false, + detail: "Queue not connected — cannot test", + }); + } + + // Build report + const passed = results.filter(r => r.status).length; + const total = results.length; + const elapsed = Date.now() - startTime; + const verdict = passed === total ? "healthy" : passed >= total - 1 ? "degraded" : "critical"; + + let report = `=== Zulip Self-Test: ${verdict.toUpperCase()} ===\n`; + report += `Time: ${elapsed}ms | ${passed}/${total} checks passed\n\n`; + for (const r of results) { + report += `${r.status ? "✅" : "❌"} ${r.name}: ${r.detail}\n`; + } + report += `\nSchema: zulip-health/v1 (standardized)\n`; + report += `Run /zulip-status for raw health data.`; + + ctx.ui.notify(report, verdict === "healthy" ? "info" : "error"); + }, + }); +}