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}