import { readFileSync, existsSync, mkdirSync, statSync, unlinkSync, } from "fs"; import { Bot } from "grammy"; import { join } from "path"; // ─── Config ─── const BOT_TOKEN = process.env.TELEGRAM_BOT_TOKEN!; const ALLOWED_USERS = (process.env.TELEGRAM_ALLOWED_USERS || "") .split(",") .map((s) => s.trim()) .filter(Boolean); const SSH_USER = process.env.SSH_USER || "root"; const SSH_HOST = process.env.SSH_HOST || "159.195.60.33"; const SSH_PORT = process.env.SSH_PORT || "22"; const PI_MODEL = process.env.PI_MODEL || "claude-sonnet-4-6"; const PI_TIMEOUT_MS = 10 * 60 * 1000; // 10 min const SESSION_DIR = join(import.meta.dir, "sessions"); const SESSION_WARN_KB = 80_000; const SESSION_MAX_KB = 100_000; const MAX_RETRY = 1; const RATE_LIMIT_WINDOW = 60_000; const RATE_LIMIT_MAX = 15; // ─── State ─── const activeProcesses = new Map(); const rateLimits = new Map(); // Ensure session dir if (!existsSync(SESSION_DIR)) mkdirSync(SESSION_DIR, { recursive: true }); // ─── Infrastructure context ─── function loadInfraContext(): string { try { const infraPath = join(import.meta.dir, "../.pi/infra.md"); if (existsSync(infraPath)) return readFileSync(infraPath, "utf-8"); } catch {} return ""; } const INFRA_CONTEXT = loadInfraContext(); // ─── System Prompt ─── const SYSTEM_PROMPT = `You are a DevOps management bot for CharityRight, operated via Telegram. INFRASTRUCTURE: - Primary server: ${SSH_USER}@${SSH_HOST}:${SSH_PORT} - Containers: cr-server-new (CharityRight), qc-server-new (QuikCue) via Incus - Services: Docker, Dokploy, Postgres, Redis, sync services ${INFRA_CONTEXT ? `\nDETAILED INFRA:\n${INFRA_CONTEXT}` : ""} RULES: 1. ALL server commands MUST use: ssh ${SSH_USER}@${SSH_HOST} -p ${SSH_PORT} "command" 2. For containers: ssh ${SSH_USER}@${SSH_HOST} "incus exec cr-server-new -- command" 3. NEVER run commands directly on the local machine 4. NEVER read local files, .env files, or secrets 5. NEVER run destructive commands: rm -rf /, mkfs, dd, DROP DATABASE 6. Format responses clearly with status indicators 7. Be concise — this is a mobile chat interface`; // ─── Dangerous command patterns ─── const DANGEROUS_PATTERNS = [ /rm\s+-rf\s+\/(?!\w)/, /mkfs\./, /dd\s+if=/, /:(){ :\|:& };:/, /shutdown/, /reboot/, /init\s+0/, /DROP\s+(DATABASE|TABLE)/i, /TRUNCATE\s+TABLE/i, /DELETE\s+FROM\s+\w+\s*;/i, ]; // ─── Bot ─── const bot = new Bot(BOT_TOKEN); function isAllowed(userId: number): boolean { if (ALLOWED_USERS.length === 0) return false; return ALLOWED_USERS.includes(String(userId)); } function checkRateLimit(userId: number): boolean { const now = Date.now(); const timestamps = rateLimits.get(userId) || []; const recent = timestamps.filter((t) => now - t < RATE_LIMIT_WINDOW); recent.push(now); rateLimits.set(userId, recent); return recent.length <= RATE_LIMIT_MAX; } function getSessionPath(userId: number): string { return join(SESSION_DIR, `user-${userId}.jsonl`); } function checkSession(userId: number): { exists: boolean; sizeKB: number } { const path = getSessionPath(userId); if (!existsSync(path)) return { exists: false, sizeKB: 0 }; const stat = statSync(path); return { exists: true, sizeKB: Math.round(stat.size / 1024) }; } // ─── SSH helper ─── async function sshExec( command: string, timeoutMs = 30_000 ): Promise<{ stdout: string; stderr: string; exitCode: number }> { const proc = Bun.spawn( [ "ssh", "-o", "StrictHostKeyChecking=no", "-o", "ConnectTimeout=10", "-p", SSH_PORT, `${SSH_USER}@${SSH_HOST}`, command, ], { stdout: "pipe", stderr: "pipe" } ); const timer = setTimeout(() => proc.kill(), timeoutMs); const [stdout, stderr] = await Promise.all([ new Response(proc.stdout).text(), new Response(proc.stderr).text(), ]); clearTimeout(timer); const exitCode = await proc.exited; return { stdout: stdout.trim(), stderr: stderr.trim(), exitCode }; } // ─── Pi Streaming ─── async function callPiStreaming( userId: number, message: string, signal: AbortSignal, ctx: any, retries = MAX_RETRY ): Promise { const session = getSessionPath(userId); const hasSession = existsSync(session); const args = [ "pi", "--mode", "json", "-p", "--no-extensions", "--model", PI_MODEL, "--tools", "bash", "--thinking", "off", "--append-system-prompt", SYSTEM_PROMPT, ]; if (hasSession) { args.push("--session", session, "--continue"); } else { args.push("--session", session); } args.push(message); // Strip secrets from child env const { ANTHROPIC_API_KEY: _removed, ...cleanEnv } = process.env; const proc = Bun.spawn({ cmd: args, cwd: "/tmp", stdout: "pipe", stderr: "pipe", env: cleanEnv, }); const timeout = setTimeout(() => { console.log(`⏰ [${userId}] Pi timed out`); proc.kill(); }, PI_TIMEOUT_MS); const abortHandler = () => proc.kill(); signal.addEventListener("abort", abortHandler, { once: true }); let finalText = ""; let toolCalls = 0; let progressMsgId: number | null = null; let lastProgressUpdate = 0; try { const reader = proc.stdout.getReader(); const decoder = new TextDecoder(); let buffer = ""; while (true) { const { done, value } = await reader.read(); if (done) break; buffer += decoder.decode(value, { stream: true }); const lines = buffer.split("\n"); buffer = lines.pop() || ""; for (const line of lines) { if (!line.trim()) continue; try { const event = JSON.parse(line); // Tool execution — update single progress message if (event.type === "tool_execution_start") { toolCalls++; const input = event.tool_input; let preview = ""; if (input?.command) { preview = input.command.substring(0, 100); } else if (input?.path) { preview = `📄 ${input.path}`; } else { preview = event.tool_name || "working"; } const progressText = `⚙️ [${toolCalls}] ${preview}...`; const now = Date.now(); // Throttle to every 2s if (now - lastProgressUpdate > 2000) { try { if (progressMsgId) { await ctx.api.editMessageText(ctx.chat.id, progressMsgId, progressText); } else { const sent = await ctx.reply(progressText); progressMsgId = sent.message_id; } lastProgressUpdate = now; } catch {} } } // Agent end — extract final assistant response if (event.type === "agent_end" && event.messages) { const allText: string[] = []; for (const msg of event.messages) { if (msg.role === "assistant" && Array.isArray(msg.content)) { for (const block of msg.content) { if (block.type === "text" && block.text) allText.push(block.text); } } } if (allText.length > 0) { finalText = allText[allText.length - 1]; } } } catch {} } } // Drain stderr const stderrText = await new Response(proc.stderr).text(); clearTimeout(timeout); signal.removeEventListener("abort", abortHandler); const exitCode = await proc.exited; // Delete progress message if (progressMsgId) { try { await ctx.api.deleteMessage(ctx.chat.id, progressMsgId); } catch {} } // Handle corrupted session — retry with counter if ( exitCode !== 0 && hasSession && retries > 0 && (stderrText.includes("session") || stderrText.includes("corrupt") || stderrText.includes("parse")) ) { console.log(`🔄 [${userId}] Corrupted session — clearing and retrying (${retries} left)...`); try { unlinkSync(session); } catch {} return callPiStreaming(userId, message, signal, ctx, retries - 1); } if (!finalText.trim()) { if (stderrText) return `⚠️ Pi error: ${stderrText.substring(0, 500)}`; return "⚠️ No response from Pi agent."; } return finalText; } catch (error: any) { clearTimeout(timeout); signal.removeEventListener("abort", abortHandler); if (progressMsgId) { try { await ctx.api.deleteMessage(ctx.chat.id, progressMsgId); } catch {} } if (signal.aborted) return "🛑 Task cancelled."; return `❌ Error: ${error.message}`; } } // ─── Command Handlers ─── bot.command("start", async (ctx) => { if (!isAllowed(ctx.from?.id || 0)) return; await ctx.reply( `🤖 CharityRight DevOps Bot\n\n` + `Commands:\n` + `/status — Server health check\n` + `/ssh <cmd> — Run SSH command\n` + `/health — Pi Worker status\n` + `/asana — List Asana tasks\n` + `/clear — Reset conversation\n` + `/stop — Cancel running task\n\n` + `Or just type naturally — I'll figure it out.`, { parse_mode: "HTML" } ); }); bot.command("status", async (ctx) => { if (!isAllowed(ctx.from?.id || 0)) return; const msg = await ctx.reply("🔍 Checking server status..."); try { const { stdout, exitCode } = await sshExec( "uptime && echo '---' && df -h / | tail -1 && echo '---' && free -h | grep Mem && echo '---' && incus list -f compact 2>/dev/null || true" ); const status = exitCode === 0 ? "✅" : "⚠️"; await ctx.api.editMessageText( ctx.chat.id, msg.message_id, `${status} Server Status\n
${stdout.substring(0, 3500)}
`, { parse_mode: "HTML" } ); } catch (error: any) { await ctx.api.editMessageText(ctx.chat.id, msg.message_id, `❌ SSH failed: ${error.message}`); } }); bot.command("ssh", async (ctx) => { if (!isAllowed(ctx.from?.id || 0)) return; const cmd = ctx.match; if (!cmd) { await ctx.reply("Usage: /ssh "); return; } if (DANGEROUS_PATTERNS.some((p) => p.test(cmd))) { await ctx.reply("🚫 Blocked — dangerous command detected."); return; } const msg = await ctx.reply(`⚙️ Running: ${cmd.substring(0, 200)}`, { parse_mode: "HTML" }); try { const { stdout, stderr, exitCode } = await sshExec(cmd); const output = stdout || stderr || "(no output)"; const status = exitCode === 0 ? "✅" : `⚠️ Exit ${exitCode}`; await ctx.api.editMessageText( ctx.chat.id, msg.message_id, `${status}\n
${output.substring(0, 3500)}
`, { parse_mode: "HTML" } ); } catch (error: any) { await ctx.api.editMessageText(ctx.chat.id, msg.message_id, `❌ ${error.message}`); } }); bot.command("health", async (ctx) => { if (!isAllowed(ctx.from?.id || 0)) return; const msg = await ctx.reply("🔍 Checking Pi Worker..."); try { const { stdout } = await sshExec("incus exec cr-server-new -- curl -s http://localhost:8787/health"); const health = JSON.parse(stdout); const statusEmoji = health.status === "ok" ? "✅" : "⚠️"; await ctx.api.editMessageText( ctx.chat.id, msg.message_id, `${statusEmoji} Pi Worker\n` + `Status: ${health.status}\n` + `Uptime: ${Math.floor(health.uptime / 60)}m\n` + `Task Loop: ${health.taskLoop ? "✅" : "❌"}\n` + `Improvement Loop: ${health.improvementLoop ? "✅" : "❌"}`, { parse_mode: "HTML" } ); } catch (error: any) { await ctx.api.editMessageText(ctx.chat.id, msg.message_id, `❌ Pi Worker unreachable: ${error.message}`); } }); bot.command("asana", async (ctx) => { if (!isAllowed(ctx.from?.id || 0)) return; const msg = await ctx.reply("🔍 Fetching Asana tasks..."); try { const token = process.env.ASANA_ACCESS_TOKEN; if (!token) { await ctx.api.editMessageText(ctx.chat.id, msg.message_id, "❌ No ASANA_ACCESS_TOKEN"); return; } const res = await fetch( "https://app.asana.com/api/1.0/user_task_lists/1212299150915258/tasks?opt_fields=name,completed,assignee_status&limit=20", { headers: { Authorization: `Bearer ${token}` } } ); const data = (await res.json()) as any; const tasks = data.data || []; const open = tasks.filter((t: any) => !t.completed); let text = `📋 My Tasks (${open.length} open)\n\n`; for (const t of open.slice(0, 15)) { text += `• ${t.name.substring(0, 60)}\n`; } if (open.length > 15) text += `\n... and ${open.length - 15} more`; await ctx.api.editMessageText(ctx.chat.id, msg.message_id, text, { parse_mode: "HTML" }); } catch (error: any) { await ctx.api.editMessageText(ctx.chat.id, msg.message_id, `❌ Asana error: ${error.message}`); } }); bot.command("clear", async (ctx) => { if (!isAllowed(ctx.from?.id || 0)) return; const userId = ctx.from!.id; const session = getSessionPath(userId); try { unlinkSync(session); } catch {} await ctx.reply("🧹 Session cleared."); }); bot.command("stop", async (ctx) => { const userId = ctx.from?.id || 0; const controller = activeProcesses.get(userId); if (controller) { controller.abort(); activeProcesses.delete(userId); await ctx.reply("🛑 Stopped."); } else { await ctx.reply("ℹ️ Nothing running."); } }); bot.command("cancel", async (ctx) => { const userId = ctx.from?.id || 0; const controller = activeProcesses.get(userId); if (controller) { controller.abort(); activeProcesses.delete(userId); await ctx.reply("🛑 Cancelled."); } else { await ctx.reply("ℹ️ Nothing running."); } }); // ─── Natural language messages ─── bot.on("message:text", async (ctx) => { const userId = ctx.from?.id || 0; if (!isAllowed(userId)) return; if (!checkRateLimit(userId)) { await ctx.reply("⏳ Slow down — rate limit reached."); return; } const msg = ctx.message.text; // Block concurrent tasks per user if (activeProcesses.has(userId)) { await ctx.reply("⏳ Already processing a task. Use /stop to cancel it first."); return; } // Session size management const sessionInfo = checkSession(userId); if (sessionInfo.sizeKB > SESSION_WARN_KB) { if (sessionInfo.sizeKB > SESSION_MAX_KB) { try { unlinkSync(getSessionPath(userId)); } catch {} await ctx.reply( `🧹 Session auto-cleared (${sessionInfo.sizeKB}KB exceeded ${SESSION_MAX_KB}KB limit). Starting fresh.` ); } else { await ctx.reply(`⚠️ Session is ${sessionInfo.sizeKB}KB/${SESSION_MAX_KB}KB. Consider /clear soon.`); } } const controller = new AbortController(); activeProcesses.set(userId, controller); try { const response = await callPiStreaming(userId, msg, controller.signal, ctx); // Send response in chunks (Telegram 4096 char limit) const maxLen = 4000; for (let i = 0; i < response.length; i += maxLen) { const chunk = response.substring(i, i + maxLen); await ctx.reply(chunk); } } catch (error: any) { if (!controller.signal.aborted) { await ctx.reply(`❌ Error: ${error.message}`); } } finally { activeProcesses.delete(userId); } }); // ─── Error handler ─── bot.catch((err) => { const msg = err.message || String(err); // 409 is expected during restart — ignore it if (msg.includes("409") || msg.includes("Conflict")) { console.log("⚠️ Telegram 409 conflict (normal during restart, retrying...)"); return; } console.error("Bot error:", msg); }); // ─── Startup ─── console.log("🤖 CharityRight DevOps Bot starting..."); const BOT_API = `https://api.telegram.org/bot${BOT_TOKEN}`; let pollOffset = 0; let running = true; async function pollLoop() { // Grammy requires bot.init() before handleUpdate() can work await bot.init(); console.log("✅ Bot online — polling for messages"); while (running) { try { const res = await fetch(`${BOT_API}/getUpdates?offset=${pollOffset}&limit=100&timeout=30`); const data = (await res.json()) as any; if (!data.ok) { if (data.error_code === 409) { console.log("⚠️ 409 conflict — waiting 10s..."); await new Promise((r) => setTimeout(r, 10_000)); continue; } console.error("Poll error:", data.description); await new Promise((r) => setTimeout(r, 3000)); continue; } for (const update of data.result || []) { pollOffset = update.update_id + 1; try { await bot.handleUpdate(update); } catch (e: any) { console.error("Handler error:", e.message?.substring(0, 200)); } } } catch (e: any) { console.error("Network error:", e.message?.substring(0, 100)); await new Promise((r) => setTimeout(r, 3000)); } } } await pollLoop(); const cleanup = () => { console.log("🛑 Shutting down..."); running = false; for (const [, controller] of activeProcesses) controller.abort(); process.exit(0); }; process.on("SIGINT", cleanup); process.on("SIGTERM", cleanup); process.on("unhandledRejection", (err: any) => { console.error("Unhandled:", err?.message?.substring(0, 200) || err); });