Parent directory

telegram.ts

14071 bytes
  1import { Bot, InlineKeyboard, type Context } from "grammy";
  2import type { Chat } from "grammy/types";
  3import type { Part } from "@opencode-ai/sdk";
  4import * as log from "./log.js";
  5import type { PermissionEvent } from "./events.js";
  6import {
  7  escapeMarkdownV2,
  8  formatParts,
  9  formatTextParts,
 10  splitMessage,
 11} from "./format.js";
 12import {
 13  getOrCreateSession,
 14  createNewSession,
 15  listSessions,
 16  getSessionId,
 17  switchSession,
 18  abortSession,
 19  sendPrompt,
 20  replyPermission,
 21  listAgents,
 22  getAgent,
 23  setAgent,
 24} from "./opencode.js";
 25import { registerSession, unregisterSession } from "./events.js";
 26
 27export const BOT_COMMANDS = [
 28  { command: "new", description: "New session" },
 29  { command: "sessions", description: "List and switch sessions" },
 30  { command: "abort", description: "Abort current session" },
 31  { command: "agent", description: "Switch agent (/agent <name>)" },
 32];
 33
 34function isChannel(chat: Chat): boolean {
 35  return chat.type === "channel" || chat.type === "supergroup";
 36}
 37
 38function formatAsQuote(text: string): string {
 39  const escaped = escapeMarkdownV2(text);
 40  const lines = escaped.split("\n");
 41  const quoted = lines.map((line) => `> ${line}`).join("\n");
 42  return quoted;
 43}
 44
 45const THROTTLE_MS = 2000;
 46
 47// Track channel membership for quoting
 48let channelChatId: number | null = null;
 49
 50// Telegram callback data is limited to 64 bytes. Permission IDs are too long,
 51// so we store them in a map keyed by a short incrementing counter.
 52let permCounter = 0;
 53const pendingPerms = new Map<
 54  string,
 55  { sessionId: string; permissionId: string }
 56>();
 57
 58function formatPartsPreview(parts: Part[]): string {
 59  const text = formatParts(parts);
 60  if (!text) return escapeMarkdownV2("thinking...");
 61  // Truncate to fit Telegram's 4096 limit for edits
 62  if (text.length > 4000) return text.slice(0, 4000) + escapeMarkdownV2("...");
 63  return text;
 64}
 65
 66async function editMessage(
 67  ctx: Context,
 68  messageId: number,
 69  text: string,
 70): Promise<void> {
 71  try {
 72    await ctx.api.editMessageText(ctx.chat!.id, messageId, text, {
 73      parse_mode: "MarkdownV2",
 74    });
 75  } catch {
 76    // Edit may fail if text unchanged or message deleted — ignore
 77  }
 78}
 79
 80function formatPermissionMessage(perm: PermissionEvent): string {
 81  const lines = [escapeMarkdownV2(`Permission: ${perm.permission}`)];
 82  if (perm.patterns.length > 0) {
 83    lines.push(escapeMarkdownV2(perm.patterns.join("\n")));
 84  }
 85  return lines.join("\n");
 86}
 87
 88function trackChannel(ctx: Context): void {
 89  const chat = ctx.chat;
 90  if (!chat) return;
 91  if (isChannel(chat)) {
 92    channelChatId = chat.id;
 93  }
 94}
 95
 96export function createBot(token: string, allowedUsers: number[]): Bot {
 97  const bot = new Bot(token);
 98
 99  const allowed = new Set(allowedUsers);
100
101  // Auth middleware — reject users not in allowlist
102  bot.use(async (ctx, next) => {
103    const userId = ctx.from?.id;
104    if (!userId || !allowed.has(userId)) {
105      await ctx.reply("Not authorized.");
106      return;
107    }
108    trackChannel(ctx);
109    await next();
110  });
111
112  bot.command("new", async (ctx) => {
113    const chatId = ctx.chat.id;
114    log.info(`[cmd] /new chat=${chatId}`);
115    try {
116      const sessionId = await createNewSession(chatId);
117      log.info(`[session] created session=${sessionId} chat=${chatId}`);
118      await ctx.reply(
119        `New session created: \`${escapeMarkdownV2(sessionId)}\``,
120        {
121          parse_mode: "MarkdownV2",
122        },
123      );
124    } catch (err) {
125      log.error(`[cmd] /new error:`, err);
126      await ctx.reply(`Failed to create session: ${String(err)}`);
127    }
128  });
129
130  bot.command("sessions", async (ctx) => {
131    log.info(`[cmd] /sessions chat=${ctx.chat.id}`);
132    try {
133      const sessions = await listSessions();
134      if (sessions.length === 0) {
135        await ctx.reply("No sessions found\\.", { parse_mode: "MarkdownV2" });
136        return;
137      }
138
139      const currentId = getSessionId(ctx.chat.id);
140      const lines = sessions.map((s) => {
141        const marker = s.id === currentId ? " \\(active\\)" : "";
142        const title = escapeMarkdownV2(s.title || "untitled");
143        const id = escapeMarkdownV2(s.id.slice(0, 8));
144        return `• \`${id}\` ${title}${marker}`;
145      });
146
147      const keyboard = new InlineKeyboard();
148      for (const s of sessions) {
149        if (s.id === currentId) continue;
150        const label = (s.title || "untitled").slice(0, 30);
151        keyboard.text(label, `switch:${s.id}`).row();
152      }
153
154      await ctx.reply(lines.join("\n"), {
155        parse_mode: "MarkdownV2",
156        reply_markup:
157          keyboard.inline_keyboard.length > 0 ? keyboard : undefined,
158      });
159    } catch (err) {
160      log.error(`[cmd] /sessions error:`, err);
161      await ctx.reply(`Failed to list sessions: ${String(err)}`);
162    }
163  });
164
165  bot.command("abort", async (ctx) => {
166    const sessionId = getSessionId(ctx.chat.id);
167    log.info(`[cmd] /abort chat=${ctx.chat.id} session=${sessionId ?? "none"}`);
168    if (!sessionId) {
169      await ctx.reply("No active session to abort.");
170      return;
171    }
172    try {
173      await abortSession(sessionId);
174      log.info(`[session] aborted session=${sessionId}`);
175      await ctx.reply("Session aborted\\.", { parse_mode: "MarkdownV2" });
176    } catch (err) {
177      log.error(`[cmd] /abort error:`, err);
178      await ctx.reply(`Failed to abort: ${String(err)}`);
179    }
180  });
181
182  // Handle callback queries (session switch, permissions)
183  bot.on("callback_query:data", async (ctx) => {
184    const data = ctx.callbackQuery.data;
185
186    // Session switch: switch:<sessionId>
187    if (data.startsWith("switch:")) {
188      const sessionId = data.slice("switch:".length);
189      const chatId = ctx.chat?.id;
190      if (!chatId || !sessionId) {
191        await ctx.answerCallbackQuery({ text: "Invalid switch data" });
192        return;
193      }
194      switchSession(chatId, sessionId);
195      log.info(`[session] switched to session=${sessionId} chat=${chatId}`);
196      await ctx.answerCallbackQuery({ text: "Session switched" });
197      await ctx.editMessageReplyMarkup({ reply_markup: undefined });
198      return;
199    }
200
201    // Format: p:<allow|deny>:<key>
202    if (!data.startsWith("p:")) {
203      await ctx.answerCallbackQuery();
204      return;
205    }
206    const [, action, key] = data.split(":");
207    const perm = key ? pendingPerms.get(key) : undefined;
208    if (!action || !perm) {
209      await ctx.answerCallbackQuery({ text: "Permission expired" });
210      return;
211    }
212    pendingPerms.delete(key!);
213    const responseMap: Record<string, "once" | "always" | "reject"> = {
214      a: "once",
215      s: "always",
216      d: "reject",
217    };
218    const permResponse = responseMap[action] ?? "reject";
219    try {
220      await replyPermission(perm.sessionId, perm.permissionId, permResponse);
221      log.info(
222        `[permission] ${permResponse} session=${perm.sessionId} perm=${perm.permissionId}`,
223      );
224      await ctx.answerCallbackQuery({ text: `Permission: ${permResponse}` });
225      await ctx.deleteMessage();
226    } catch (err) {
227      log.error(`[permission] reply error:`, err);
228      await ctx.answerCallbackQuery({ text: `Error: ${String(err)}` });
229    }
230  });
231
232  bot.command("agent", async (ctx) => {
233    const chatId = ctx.chat.id;
234    const name = ctx.match?.trim().toLowerCase();
235    log.info(`[cmd] /agent chat=${chatId} arg=${name ?? "(none)"}`);
236
237    if (!name) {
238      // Show current agent and list available
239      try {
240        const agents = await listAgents();
241        const current = getAgent(chatId) ?? escapeMarkdownV2("(default)");
242        const list = agents.map((a) => escapeMarkdownV2(a)).join(", ");
243        await ctx.reply(
244          `Current agent: \`${current}\`\nAvailable: ${list}\n\nUsage: /agent \\<name\\>`,
245          { parse_mode: "MarkdownV2" },
246        );
247      } catch (err) {
248        log.error(`[cmd] /agent list error:`, err);
249        await ctx.reply(`Failed to list agents: ${String(err)}`);
250      }
251      return;
252    }
253
254    try {
255      const agents = await listAgents();
256      const match = agents.find((a) => a.toLowerCase() === name);
257      if (!match) {
258        const list = agents.map((a) => escapeMarkdownV2(a)).join(", ");
259        await ctx.reply(
260          `Unknown agent \`${escapeMarkdownV2(name)}\`\\. Available: ${list}`,
261          { parse_mode: "MarkdownV2" },
262        );
263        return;
264      }
265      setAgent(chatId, match);
266      log.info(`[agent] set agent=${match} chat=${chatId}`);
267      await ctx.reply(`Agent set to \`${escapeMarkdownV2(match)}\`\\.`, {
268        parse_mode: "MarkdownV2",
269      });
270    } catch (err) {
271      log.error(`[cmd] /agent error:`, err);
272      await ctx.reply(`Error: ${String(err)}`);
273    }
274  });
275
276  // Text message handler — forward to OpenCode with streaming
277  bot.on("message:text", async (ctx) => {
278    const chatId = ctx.chat.id;
279    log.info(`[prompt] message received chat=${chatId}`);
280    try {
281      const { sessionId, fallback } = await getOrCreateSession(chatId);
282      if (fallback) {
283        log.info(
284          `[session] previous session gone, created new session=${sessionId} chat=${chatId}`,
285        );
286        await ctx.reply(
287          escapeMarkdownV2(
288            "Previous session no longer available. Started a new session.",
289          ),
290          { parse_mode: "MarkdownV2" },
291        );
292      }
293
294      // Typing indicator, refreshed every 4s
295      let typingInterval: ReturnType<typeof setInterval> | null = setInterval(
296        () => {
297          void ctx.api.sendChatAction(chatId, "typing");
298        },
299        4000,
300      );
301      void ctx.api.sendChatAction(chatId, "typing");
302
303      // Send "thinking..." immediately, then edit-in-place as streaming arrives
304      let lastEditTime = 0;
305      let editTimer: ReturnType<typeof setTimeout> | null = null;
306      let latestPreview = "";
307      const thinkingMsg = await ctx.api.sendMessage(
308        chatId,
309        escapeMarkdownV2("thinking..."),
310        { parse_mode: "MarkdownV2" },
311      );
312      const responseMsgId = thinkingMsg.message_id;
313
314      const cleanup = () => {
315        if (editTimer) {
316          clearTimeout(editTimer);
317          editTimer = null;
318        }
319        if (typingInterval) {
320          clearInterval(typingInterval);
321          typingInterval = null;
322        }
323        unregisterSession(sessionId);
324      };
325
326      const flushEdit = () => {
327        if (latestPreview) {
328          const textToSend =
329            channelChatId !== null
330              ? formatAsQuote(latestPreview)
331              : latestPreview;
332          void editMessage(ctx, responseMsgId, textToSend);
333          lastEditTime = Date.now();
334        }
335      };
336
337      // Register SSE handlers for streaming and permissions before firing prompt
338      registerSession(
339        sessionId,
340        // onPart — send first message on first data, then throttled edits
341        (parts: Part[]) => {
342          latestPreview = formatPartsPreview(parts);
343          const elapsed = Date.now() - lastEditTime;
344          if (elapsed >= THROTTLE_MS) {
345            if (editTimer) clearTimeout(editTimer);
346            editTimer = null;
347            flushEdit();
348          } else if (!editTimer) {
349            editTimer = setTimeout(flushEdit, THROTTLE_MS - elapsed);
350          }
351        },
352        // onPermission
353        async (perm: PermissionEvent) => {
354          log.info(
355            `[permission] request permission=${perm.permission} session=${perm.sessionID} perm=${perm.id}`,
356          );
357          const key = String(++permCounter);
358          pendingPerms.set(key, {
359            sessionId: perm.sessionID,
360            permissionId: perm.id,
361          });
362          const keyboard = new InlineKeyboard()
363            .text("Allow", `p:a:${key}`)
364            .text("Session", `p:s:${key}`)
365            .text("Deny", `p:d:${key}`);
366          await ctx.api.sendMessage(chatId, formatPermissionMessage(perm), {
367            parse_mode: "MarkdownV2",
368            reply_markup: keyboard,
369          });
370        },
371      );
372
373      // Fire prompt without blocking grammY's update loop (permissions need callback handling)
374      const userText = ctx.message.text;
375      const activeAgent = getAgent(chatId);
376      log.info(
377        `[prompt] sending to session=${sessionId} agent=${activeAgent ?? "(default)"}`,
378      );
379      sendPrompt(sessionId, userText, activeAgent)
380        .then(async (parts) => {
381          cleanup();
382          // Final edit of the thinking message with full content (tools + text)
383          const fullContent = formatParts(parts);
384          if (fullContent) {
385            const editText =
386              fullContent.length > 4000
387                ? fullContent.slice(0, 4000) + escapeMarkdownV2("...")
388                : fullContent;
389            const editToSend =
390              channelChatId !== null ? formatAsQuote(editText) : editText;
391            await editMessage(ctx, responseMsgId, editToSend);
392          }
393          // Send final summary (text-only) as new message, but only if it
394          // differs from what's already in the edited message
395          const textContent = formatTextParts(parts);
396          if (textContent && textContent !== fullContent) {
397            const chunks = splitMessage(textContent);
398            log.info(
399              `[prompt] done session=${sessionId} chunks=${chunks.length}`,
400            );
401            for (const chunk of chunks) {
402              const textToSend =
403                channelChatId !== null ? formatAsQuote(chunk) : chunk;
404              await ctx.api.sendMessage(chatId, textToSend, {
405                parse_mode: "MarkdownV2",
406              });
407            }
408          }
409        })
410        .catch(async (err) => {
411          cleanup();
412          log.error(`[prompt] error session=${sessionId}:`, err);
413          const errText = escapeMarkdownV2(`Error: ${String(err)}`);
414          const errTextToSend =
415            channelChatId !== null ? formatAsQuote(errText) : errText;
416          await editMessage(ctx, responseMsgId, errTextToSend);
417        });
418    } catch (err) {
419      log.error(`[prompt] unhandled error chat=${chatId}:`, err);
420      await ctx.reply(`Error: ${String(err)}`);
421    }
422  });
423
424  return bot;
425}