返回 CodeWhale
index.mjs
1 import {
2 activeTurnBlock,
3 activeTurnKeyboard,
4 approvalKeyboard,
5 callbackAction,
6 commandAction,
7 compactRuntimeError,
8 controlKeyboard,
9 envFirst,
10 helpText,
11 isAllowed,
12 isGroupChat,
13 isTelegramMarkdownParseError,
14 latestRunningTurn,
15 looksLikePollingConflict,
16 pairingRefusalText,
17 parseBool,
18 parseCommand,
19 parseList,
20 preservedChatStateFields,
21 splitMessage,
22 stripGroupPrefix,
23 threadListKeyboard,
24 telegramIdentity,
25 telegramMessageBody,
26 telegramPollingConflictDelayMs,
27 telegramRetryDelayMs,
28 telegramSendRetryDelayMs
29 } from "./lib.mjs";
30 import {
31 createRuntimeClient,
32 readJsonSafe,
33 readSse,
34 ThreadStore as CoreThreadStore,
35 ApprovalOwnershipError,
36 decideApproval as decideRuntimeApproval,
37 settleStoredApproval
38 } from "../../bridge-core/src/lib.mjs";
39
40 const TYPING_INTERVAL_MS = 2000;
41 const TYPING_TIMEOUT_MS = 1500;
42 const LAST_SEQ_FLUSH_INTERVAL_MS = 2000;
43
44 class ThreadStore extends CoreThreadStore {
45 constructor(filePath) {
46 super(filePath, { messageLimit: 500, actions: true });
47 }
48 }
49
50 const config = {
51 botToken: requiredEnv("TELEGRAM_BOT_TOKEN"),
52 apiBase: (process.env.TELEGRAM_API_BASE || "https://api.telegram.org").replace(/\/+$/, ""),
53 runtimeUrl: (envFirst(process.env, "CODEWHALE_RUNTIME_URL", "DEEPSEEK_RUNTIME_URL") || "http://127.0.0.1:7878").replace(/\/+$/, ""),
54 runtimeToken: requiredEnvFirst("CODEWHALE_RUNTIME_TOKEN", "DEEPSEEK_RUNTIME_TOKEN"),
55 workspace: envFirst(process.env, "CODEWHALE_WORKSPACE", "DEEPSEEK_WORKSPACE") || process.cwd(),
56 model: envFirst(process.env, "CODEWHALE_MODEL", "DEEPSEEK_MODEL") || "auto",
57 mode: envFirst(process.env, "CODEWHALE_MODE", "DEEPSEEK_MODE") || "agent",
58 allowShell: parseBool(envFirst(process.env, "CODEWHALE_ALLOW_SHELL", "DEEPSEEK_ALLOW_SHELL"), true),
59 trustMode: parseBool(envFirst(process.env, "CODEWHALE_TRUST_MODE", "DEEPSEEK_TRUST_MODE"), false),
60 autoApprove: parseBool(envFirst(process.env, "CODEWHALE_AUTO_APPROVE", "DEEPSEEK_AUTO_APPROVE"), false),
61 allowlist: parseList(
62 envFirst(process.env, "TELEGRAM_CHAT_ALLOWLIST", "CODEWHALE_CHAT_ALLOWLIST", "DEEPSEEK_CHAT_ALLOWLIST")
63 ),
64 allowUnlisted: parseBool(
65 envFirst(process.env, "TELEGRAM_ALLOW_UNLISTED", "CODEWHALE_ALLOW_UNLISTED", "DEEPSEEK_ALLOW_UNLISTED"),
66 false
67 ),
68 threadMapPath:
69 process.env.TELEGRAM_THREAD_MAP_PATH ||
70 "/var/lib/codewhale-telegram-bridge/thread-map.json",
71 allowGroups: parseBool(process.env.TELEGRAM_ALLOW_GROUPS, false),
72 requirePrefixInGroup: parseBool(process.env.TELEGRAM_REQUIRE_PREFIX_IN_GROUP, true),
73 groupPrefix: process.env.TELEGRAM_GROUP_PREFIX || "/cw",
74 maxReplyChars: Math.min(Number(process.env.TELEGRAM_MAX_REPLY_CHARS || 3500), 4096),
75 pollTimeoutSeconds: Number(process.env.TELEGRAM_POLL_TIMEOUT_SECONDS || 50),
76 turnTimeoutMs: Number(envFirst(process.env, "CODEWHALE_TURN_TIMEOUT_MS", "DEEPSEEK_TURN_TIMEOUT_MS") || 900000)
77 };
78
79 const { runtimeJson, authHeaders } = createRuntimeClient(config);
80
81 const threadStore = await ThreadStore.open(config.threadMapPath);
82 const activeTurnTasks = new Map();
83 let stopping = false;
84 let updateOffset = threadStore.getCursor(
85 "telegram.update_offset",
86 Number(process.env.TELEGRAM_UPDATE_OFFSET || 0)
87 );
88
89 function requestStop() {
90 stopping = true;
91 abortActiveTurnStreams();
92 }
93
94 process.once("SIGINT", requestStop);
95 process.once("SIGTERM", requestStop);
96
97 console.log("Starting CodeWhale Telegram bridge");
98 console.log(`Runtime: ${config.runtimeUrl}`);
99 console.log(`Workspace: ${config.workspace}`);
100 if (!config.allowlist.length && !config.allowUnlisted) {
101 console.log("No allowlist configured. Incoming chats will receive their IDs and be refused.");
102 }
103
104 await configureBotCommands().catch((error) => {
105 console.error("failed to configure Telegram bot command menu", error);
106 });
107 void reattachActiveTurns().catch((error) => {
108 console.error("failed to reattach active Telegram bridge turns", error);
109 });
110 await pollTelegram();
111
112 async function configureBotCommands() {
113 await telegramApi("setMyCommands", {
114 commands: [
115 { command: "menu", description: "Open CodeWhale controls" },
116 { command: "status", description: "Show runtime and workspace status" },
117 { command: "threads", description: "List recent runtime threads" },
118 { command: "new", description: "Create a new thread" },
119 { command: "interrupt", description: "Interrupt the active turn" },
120 { command: "compact", description: "Compact the current thread" },
121 { command: "help", description: "Show command help" }
122 ]
123 });
124 }
125
126 async function pollTelegram() {
127 let pollingConflictAttempts = 0;
128 while (!stopping) {
129 try {
130 const updates = await telegramApi("getUpdates", {
131 offset: updateOffset || undefined,
132 timeout: config.pollTimeoutSeconds,
133 allowed_updates: ["message", "callback_query"]
134 });
135 pollingConflictAttempts = 0;
136 for (const update of updates || []) {
137 try {
138 await handleIncomingUpdate(update);
139 await markUpdateHandled(update);
140 } catch (error) {
141 console.error("failed to handle incoming Telegram update", error);
142 break;
143 }
144 }
145 } catch (error) {
146 if (looksLikePollingConflict(error)) {
147 const waitMs = telegramPollingConflictDelayMs(pollingConflictAttempts);
148 pollingConflictAttempts += 1;
149 if (waitMs == null) {
150 throw new Error(
151 "Telegram getUpdates conflict; another bridge is polling this bot token. Stop the other bridge process or use a different token."
152 );
153 }
154 console.warn(
155 `Telegram getUpdates conflict; another bridge is polling this bot. Retrying in ${Math.round(waitMs / 1000)}s.`
156 );
157 await delay(waitMs);
158 continue;
159 }
160 pollingConflictAttempts = 0;
161 const waitMs = telegramRetryDelayMs(error);
162 console.error(`Telegram poll failed: ${error.message}. Retrying in ${Math.round(waitMs / 1000)}s.`);
163 await delay(waitMs);
164 }
165 }
166 }
167
168 async function markUpdateHandled(update) {
169 if (update.update_id == null) return;
170 if (!Number.isFinite(Number(update.update_id))) return;
171 const nextOffset = Math.max(updateOffset, Number(update.update_id) + 1);
172 if (nextOffset === updateOffset) return;
173 updateOffset = nextOffset;
174 await threadStore.setCursor("telegram.update_offset", updateOffset);
175 }
176
177 async function handleIncomingUpdate(update) {
178 if (update.callback_query) {
179 if (await isReplayCallbackUpdate(update)) return;
180 await handleCallbackQuery(update.callback_query);
181 return;
182 }
183
184 const identity = telegramIdentity(update);
185 if (!identity.chatId || !identity.messageId) return;
186 if (identity.isBot) return;
187
188 const messageKey = `${identity.chatId}:${identity.messageId}`;
189 if (await threadStore.recordMessage(messageKey)) return;
190
191 if (!identity.text) {
192 await sendText(identity.chatId, "Only text messages are supported in this first bridge.");
193 return;
194 }
195
196 const scoped = stripGroupPrefix(identity.text, {
197 chatType: identity.chatType,
198 requirePrefix: config.requirePrefixInGroup,
199 prefix: config.groupPrefix
200 });
201 if (!scoped.accepted) return;
202
203 if (isGroupChat(identity.chatType) && !config.allowGroups) {
204 await sendText(
205 identity.chatId,
206 "Group chat control is disabled for this bridge. DM the bot, or set TELEGRAM_ALLOW_GROUPS=true and allowlist this chat."
207 );
208 return;
209 }
210
211 if (!isAllowed(identity, config.allowlist, config.allowUnlisted)) {
212 await sendText(identity.chatId, pairingRefusalText(identity));
213 return;
214 }
215
216 const command = parseCommand(scoped.text);
217 await rememberAuthorizedIdentity(identity);
218 await handleCommand(identity.chatId, command, identity);
219 }
220
221 async function rememberAuthorizedIdentity({ chatId, chatType, userId, username, isBot }) {
222 await threadStore.patchChat(chatId, {
223 authorizedIdentity: { chatId, chatType, userId, username, isBot }
224 });
225 }
226
227 async function isReplayCallbackUpdate(update) {
228 if (update.update_id == null) return false;
229 return threadStore.recordMessage(`callback:${update.update_id}`);
230 }
231
232 async function handleCommand(chatId, command, identity) {
233 const action = commandAction(command);
234 switch (action.kind) {
235 case "help":
236 await sendText(chatId, helpText(), { replyMarkup: controlKeyboard() });
237 return;
238 case "menu":
239 await sendMenu(chatId);
240 return;
241 case "status":
242 await sendStatus(chatId);
243 return;
244 case "threads":
245 await sendThreads(chatId);
246 return;
247 case "new_thread": {
248 const state = await ensureThread(chatId, { forceNew: true });
249 await sendText(chatId, `Created thread ${state.threadId}`, { replyMarkup: controlKeyboard() });
250 return;
251 }
252 case "resume":
253 await resumeThread(chatId, action.threadId);
254 return;
255 case "interrupt":
256 await interruptActiveTurn(chatId);
257 return;
258 case "compact":
259 await compactThread(chatId);
260 return;
261 case "approval":
262 await decideApproval(chatId, action, identity);
263 return;
264 case "set_model":
265 await setChatModel(chatId, action.modelName);
266 return;
267 case "prompt":
268 startPromptTurn(chatId, action.prompt, identity);
269 return;
270 default:
271 await sendText(chatId, helpText(), { replyMarkup: controlKeyboard() });
272 }
273 }
274
275 async function handleCallbackQuery(query) {
276 const chat = query.message?.chat || {};
277 const from = query.from || {};
278 const identity = {
279 chatId: chat.id != null ? String(chat.id) : "",
280 messageId: query.message?.message_id != null ? String(query.message.message_id) : "",
281 chatType: chat.type || "",
282 userId: from.id != null ? String(from.id) : "",
283 username: from.username ? `@${from.username}` : "",
284 firstName: from.first_name || "",
285 isBot: Boolean(from.is_bot)
286 };
287
288 if (!identity.chatId || !query.id) return;
289 if (identity.isBot) return;
290
291 if (isGroupChat(identity.chatType) && !config.allowGroups) {
292 await answerCallback(query.id, "Group control is disabled.");
293 return;
294 }
295 if (!isAllowed(identity, config.allowlist, config.allowUnlisted)) {
296 await answerCallback(query.id, "This chat is not allowlisted.");
297 return;
298 }
299
300 const action = callbackAction(query.data);
301 if (!action) {
302 await answerCallback(query.id, "Unknown action.");
303 return;
304 }
305
306 await rememberAuthorizedIdentity(identity);
307 answerCallback(query.id, "Working...").catch((error) => {
308 console.warn("failed to acknowledge Telegram callback", error);
309 });
310 await handleModalAction(identity.chatId, action, query, identity);
311 }
312
313 async function handleModalAction(chatId, action, query = null, identity = null) {
314 switch (action.kind) {
315 case "help":
316 await sendText(chatId, helpText(), { replyMarkup: controlKeyboard() });
317 return;
318 case "status":
319 await sendStatus(chatId);
320 return;
321 case "threads":
322 await sendThreads(chatId);
323 return;
324 case "new_thread": {
325 const state = await ensureThread(chatId, { forceNew: true });
326 await sendText(chatId, `Created thread ${state.threadId}`, { replyMarkup: controlKeyboard() });
327 return;
328 }
329 case "interrupt":
330 await interruptActiveTurn(chatId);
331 return;
332 case "compact":
333 await compactThread(chatId);
334 return;
335 case "set_model":
336 await setChatModel(chatId, action.modelName);
337 return;
338 case "stored_action":
339 await handleStoredAction(chatId, action, query, identity);
340 return;
341 default:
342 await sendText(chatId, helpText(), { replyMarkup: controlKeyboard() });
343 }
344 }
345
346 async function handleStoredAction(chatId, action, query = null, identity = null) {
347 // Approval buttons also belong to the thread still bound to this chat.
348 const state = await threadStore.getChat(chatId);
349 const owner = { chatId, threadId: state?.threadId, actorId: turnActor(identity) };
350 const stored = await threadStore.getAction(action.token, owner);
351 if (!stored) {
352 await sendText(chatId, "That action expired. Open /menu and try again.");
353 return;
354 }
355
356 if (stored.kind === "resume") {
357 await threadStore.takeAction(action.token, owner);
358 await resumeThread(chatId, stored.threadId);
359 return;
360 }
361
362 if (stored.kind === "approval") {
363 const suffix = action.suffix || "";
364 const decision = suffix === "deny" ? "deny" : "allow";
365 const remember = suffix === "remember";
366 // The token is consumed only after the Runtime accepted the decision.
367 try {
368 await settleStoredApproval(threadStore, action.token, owner, async () => {
369 if (!(await decideApproval(chatId, { decision, approvalId: stored.approvalId, remember, turnId: stored.owner.turnId }, identity))) {
370 throw new ApprovalOwnershipError("not delivered");
371 }
372 });
373 } catch (error) {
374 if (error instanceof ApprovalOwnershipError) return;
375 throw error;
376 }
377 if (query?.message?.message_id) {
378 await editMessageReplyMarkup(chatId, query.message.message_id, null).catch(() => {});
379 }
380 return;
381 }
382
383 await sendText(chatId, "That action is no longer supported.");
384 }
385
386 async function sendMenu(chatId) {
387 const state = await threadStore.getChat(chatId);
388 await sendText(
389 chatId,
390 [
391 "CodeWhale controls",
392 state?.threadId ? `thread=${state.threadId}` : "thread=(new on first prompt)",
393 `model=${state?.model || config.model}`
394 ].join("\n"),
395 { replyMarkup: controlKeyboard() }
396 );
397 }
398
399 async function ensureThread(chatId, { forceNew = false } = {}) {
400 const existing = await threadStore.getChat(chatId);
401 if (existing?.threadId && !forceNew) return existing;
402
403 const effectiveModel = existing?.model || config.model;
404 const thread = await runtimeJson("/v1/threads", {
405 method: "POST",
406 body: {
407 model: effectiveModel,
408 workspace: config.workspace,
409 mode: config.mode,
410 allow_shell: config.allowShell,
411 trust_mode: config.trustMode,
412 auto_approve: config.autoApprove,
413 archived: false,
414 system_prompt:
415 "You are being controlled from a Telegram phone chat. Keep status updates concise. Ask for tool approvals when needed; do not assume mobile messages imply blanket approval."
416 }
417 });
418
419 const state = {
420 ...preservedChatStateFields(existing),
421 threadId: thread.id,
422 lastSeq: 0,
423 activeTurnId: null,
424 updatedAt: new Date().toISOString()
425 };
426 await threadStore.setChat(chatId, state);
427 return state;
428 }
429
430 function startPromptTurn(chatId, prompt, identity) {
431 if (activeTurnTasks.has(chatId)) {
432 void sendText(chatId, "Thread already has an active turn. Wait for it to finish or send /interrupt.", {
433 replyMarkup: activeTurnKeyboard()
434 }).catch((error) => {
435 console.error("failed to report active Telegram bridge turn", error);
436 });
437 return;
438 }
439
440 const controller = new AbortController();
441 const task = { controller };
442 activeTurnTasks.set(chatId, task);
443 void runPrompt(chatId, prompt, { signal: controller.signal, actorId: turnActor(identity) })
444 .catch((error) => {
445 console.error("failed to run Telegram bridge prompt", error);
446 })
447 .finally(() => {
448 if (activeTurnTasks.get(chatId) === task) {
449 activeTurnTasks.delete(chatId);
450 }
451 });
452 }
453
454 function abortActiveTurnStreams() {
455 for (const task of activeTurnTasks.values()) {
456 task.controller?.abort();
457 }
458 }
459
460 async function clearActiveTurn(chatId) {
461 await threadStore.patchChat(chatId, {
462 activeTurnId: null,
463 updatedAt: new Date().toISOString()
464 }).catch((error) => {
465 console.error("failed to clear Telegram bridge active turn", error);
466 });
467 }
468
469 function startTrackedTurnStream(chatId, threadId, turnId, sinceSeq) {
470 if (activeTurnTasks.has(chatId)) return false;
471
472 const controller = new AbortController();
473 const task = { controller };
474 activeTurnTasks.set(chatId, task);
475 void streamTurnEvents(chatId, threadId, turnId, sinceSeq, { signal: controller.signal })
476 .catch((error) => {
477 console.error("failed to stream Telegram bridge turn", error);
478 })
479 .finally(async () => {
480 if (activeTurnTasks.get(chatId) === task) {
481 activeTurnTasks.delete(chatId);
482 }
483 if (!stopping) {
484 await clearActiveTurn(chatId);
485 }
486 });
487 return true;
488 }
489
490 function turnActor(identity) {
491 return identity?.userId && !identity.isBot ? `telegram-user:${identity.userId}` : "";
492 }
493
494 async function runPrompt(chatId, prompt, options = {}) {
495 if (!prompt.trim()) {
496 await sendText(chatId, helpText(), { replyMarkup: controlKeyboard() });
497 return;
498 }
499 if (!options.actorId) throw new ApprovalOwnershipError("turn needs its initiating human");
500 const state = await ensureThread(chatId);
501 const effectiveModel = state?.model || config.model;
502 const detail = await runtimeJson(`/v1/threads/${encodeURIComponent(state.threadId)}`);
503 const activeBlock = activeTurnBlock(detail, state);
504 if (activeBlock) {
505 await threadStore.patchChat(chatId, {
506 activeTurnId: activeBlock.turnId,
507 updatedAt: new Date().toISOString()
508 });
509 await sendText(chatId, activeBlock.message, { replyMarkup: activeTurnKeyboard() });
510 return;
511 }
512 if (state.activeTurnId) {
513 await threadStore.patchChat(chatId, { activeTurnId: null });
514 }
515 const sinceSeq = Number(detail.latest_seq || state.lastSeq || 0);
516
517 const turnResponse = await runtimeJson(
518 `/v1/threads/${encodeURIComponent(state.threadId)}/turns`,
519 {
520 method: "POST",
521 body: {
522 prompt,
523 input_summary: prompt.slice(0, 200),
524 model: effectiveModel,
525 mode: config.mode,
526 allow_shell: config.allowShell,
527 trust_mode: config.trustMode,
528 auto_approve: config.autoApprove
529 }
530 }
531 );
532
533 const turnId = turnResponse.turn?.id;
534 // A turn the Runtime accepted keeps streaming even when its origin cannot be
535 // recorded; its approvals then fail closed and are decided from the TUI.
536 if (turnId && options.actorId) {
537 await threadStore.recordTurnOrigin(chatId, state.threadId, turnId, options.actorId);
538 }
539 await threadStore.patchChat(chatId, {
540 activeTurnId: turnId || null,
541 lastSeq: sinceSeq,
542 updatedAt: new Date().toISOString()
543 });
544 await sendTurnText(chatId, `Started turn ${turnId || "(unknown)"}`, {
545 replyMarkup: activeTurnKeyboard()
546 });
547
548 try {
549 await streamTurnEvents(chatId, state.threadId, turnId, sinceSeq, options);
550 } finally {
551 if (!stopping) {
552 await clearActiveTurn(chatId);
553 }
554 }
555 }
556
557 async function reattachActiveTurns() {
558 for (const [chatId, state] of threadStore.listChats()) {
559 if (!state?.threadId || !state.activeTurnId) continue;
560 const identity = state.authorizedIdentity;
561 // Legacy state has no verified sender provenance. Never recover delivery
562 // under a previous allowlist or group policy, even to report completion.
563 if (!identity || identity.chatId !== chatId || identity.isBot ||
564 !["private", "group", "supergroup"].includes(identity.chatType) ||
565 (isGroupChat(identity.chatType) && !config.allowGroups) ||
566 !isAllowed(identity, config.allowlist, config.allowUnlisted)) continue;
567
568 const detail = await runtimeJson(`/v1/threads/${encodeURIComponent(state.threadId)}`);
569 const runningTurn = latestRunningTurn(detail);
570 if (!runningTurn) {
571 await threadStore.patchChat(chatId, {
572 activeTurnId: null,
573 lastSeq: Number(detail.latest_seq || state.lastSeq || 0),
574 updatedAt: new Date().toISOString()
575 });
576 await sendText(chatId, `Bridge restarted. No active turn remains for ${state.threadId}.`);
577 continue;
578 }
579
580 const turnId = runningTurn.id || state.activeTurnId;
581 const sinceSeq = Number(state.lastSeq || 0);
582 await threadStore.patchChat(chatId, {
583 activeTurnId: turnId,
584 updatedAt: new Date().toISOString()
585 });
586 await sendTurnText(
587 chatId,
588 `Bridge restarted. Reattaching to active turn ${turnId} from seq ${sinceSeq}.`
589 );
590 startTrackedTurnStream(chatId, state.threadId, turnId, sinceSeq);
591 }
592 }
593
594 async function streamTurnEvents(chatId, threadId, turnId, sinceSeq, options = {}) {
595 const controller = new AbortController();
596 let timedOut = false;
597 const timeout = setTimeout(() => {
598 timedOut = true;
599 controller.abort();
600 }, config.turnTimeoutMs);
601 const abortFromCaller = () => controller.abort();
602 if (options.signal?.aborted) {
603 controller.abort();
604 } else {
605 options.signal?.addEventListener("abort", abortFromCaller, { once: true });
606 }
607 let responseText = "";
608 let latestSeq = sinceSeq;
609 let flushedSeq = sinceSeq;
610 let lastSeqFlushAt = 0;
611 let sentProgressAt = Date.now();
612 let typingPaused = false;
613 let typingInFlight = false;
614
615 async function flushLastSeq(force = false) {
616 if (latestSeq <= flushedSeq) return;
617 if (!force && Date.now() - lastSeqFlushAt < LAST_SEQ_FLUSH_INTERVAL_MS) return;
618 await threadStore.patchChat(chatId, { lastSeq: latestSeq });
619 flushedSeq = latestSeq;
620 lastSeqFlushAt = Date.now();
621 }
622
623 const tickTyping = async () => {
624 if (stopping || typingPaused || typingInFlight) return;
625 typingInFlight = true;
626 try {
627 await sendTypingAction(chatId);
628 } catch (error) {
629 console.warn("failed to send Telegram typing action", error);
630 } finally {
631 typingInFlight = false;
632 }
633 };
634 const typingTimer = setInterval(() => {
635 void tickTyping();
636 }, TYPING_INTERVAL_MS);
637 typingTimer.unref?.();
638
639 try {
640 void tickTyping();
641 const response = await fetch(
642 `${config.runtimeUrl}/v1/threads/${encodeURIComponent(threadId)}/events?since_seq=${sinceSeq}`,
643 {
644 headers: authHeaders(),
645 signal: controller.signal
646 }
647 );
648 if (!response.ok) {
649 const body = await readJsonSafe(response);
650 throw new Error(compactRuntimeError(response.status, body));
651 }
652
653 for await (const event of readSse(response)) {
654 if (!event.data) continue;
655 const record = JSON.parse(event.data);
656 latestSeq = Math.max(latestSeq, Number(record.seq || 0));
657 await flushLastSeq(false);
658
659 if (turnId && record.turn_id && record.turn_id !== turnId) continue;
660 const lifecycleStatus =
661 record.event === "turn.lifecycle"
662 ? record.payload?.turn?.status || record.payload?.status
663 : null;
664 const stopTypingEvent =
665 record.event === "turn.completed" ||
666 ["failed", "canceled", "interrupted"].includes(lifecycleStatus);
667 if (typingPaused && record.event !== "approval.required" && !stopTypingEvent) {
668 typingPaused = false;
669 void tickTyping();
670 }
671
672 if (record.event === "item.delta" && record.payload?.kind === "agent_message") {
673 responseText += record.payload.delta || "";
674 const now = Date.now();
675 if (responseText.length > config.maxReplyChars && now - sentProgressAt > 15000) {
676 await sendTurnText(chatId, responseText.slice(0, config.maxReplyChars));
677 responseText = responseText.slice(config.maxReplyChars);
678 sentProgressAt = now;
679 }
680 }
681
682 if (record.event === "approval.required") {
683 typingPaused = true;
684 const approval = record.payload || {};
685 const approvalId = approval.approval_id || approval.id;
686 if (!approvalId) {
687 await sendTurnText(
688 chatId,
689 [
690 "Approval required",
691 `tool=${approval.tool_name || "unknown"}`,
692 approval.description || "",
693 "",
694 "No approval_id was provided by the runtime; use /status and retry from the TUI."
695 ]
696 .filter(Boolean)
697 .join("\n"),
698 { replyMarkup: controlKeyboard() }
699 );
700 continue;
701 }
702 const origin = threadStore.turnOrigin(chatId, threadId, turnId);
703 if (!origin) {
704 await sendTurnText(chatId, "Approval has no recorded initiating human; decide from the TUI.");
705 continue;
706 }
707 const actionToken = await threadStore.putAction(
708 { kind: "approval", approvalId },
709 { chatId, threadId, turnId, actorId: origin.actorId }
710 );
711 await sendTurnText(
712 chatId,
713 [
714 "Approval required",
715 `tool=${approval.tool_name || "unknown"}`,
716 `approval_id=${approvalId}`,
717 approval.description || "",
718 "",
719 `Tap a button, or reply /allow ${approvalId}`,
720 `Reply /deny ${approvalId}`
721 ]
722 .filter(Boolean)
723 .join("\n"),
724 { replyMarkup: approvalKeyboard(actionToken) }
725 );
726 }
727
728 if (record.event === "turn.completed") {
729 typingPaused = true;
730 const turn = record.payload?.turn || {};
731 const status = turn.status || "completed";
732 const error = turn.error ? `\n${turn.error}` : "";
733 if (status !== "completed") {
734 await sendTurnText(chatId, `Turn ${status}.${error}`.trim(), {
735 replyMarkup: controlKeyboard()
736 });
737 } else {
738 await sendTurnText(chatId, responseText.trim() || "Turn completed.", {
739 replyMarkup: controlKeyboard()
740 });
741 }
742 return;
743 }
744
745 if (record.event === "turn.lifecycle") {
746 if (["failed", "canceled", "interrupted"].includes(lifecycleStatus)) {
747 typingPaused = true;
748 await sendTurnText(chatId, `Turn ${lifecycleStatus}.`, { replyMarkup: controlKeyboard() });
749 return;
750 }
751 }
752 }
753 } catch (error) {
754 if (error.name === "AbortError") {
755 if (timedOut) {
756 await sendTurnText(chatId, `Turn timed out after ${Math.round(config.turnTimeoutMs / 1000)}s.`);
757 } else if (!stopping) {
758 await sendTurnText(chatId, "Turn stream aborted.");
759 }
760 return;
761 }
762 throw error;
763 } finally {
764 clearInterval(typingTimer);
765 clearTimeout(timeout);
766 options.signal?.removeEventListener("abort", abortFromCaller);
767 await flushLastSeq(true);
768 }
769 }
770
771 async function sendStatus(chatId) {
772 const [health, runtimeInfo, workspace] = await Promise.all([
773 runtimeJson("/health", { auth: false }),
774 runtimeJson("/v1/runtime/info"),
775 runtimeJson("/v1/workspace/status")
776 ]);
777 await sendText(
778 chatId,
779 [
780 `runtime=${health.status || "unknown"}`,
781 `version=${runtimeInfo.version || "unknown"}`,
782 `bind=${runtimeInfo.bind_host}:${runtimeInfo.port}`,
783 `auth_required=${runtimeInfo.auth_required}`,
784 `workspace=${workspace.workspace}`,
785 `git_repo=${workspace.git_repo}`,
786 workspace.branch ? `branch=${workspace.branch}` : "",
787 `staged=${workspace.staged} unstaged=${workspace.unstaged} untracked=${workspace.untracked}`
788 ]
789 .filter(Boolean)
790 .join("\n"),
791 { replyMarkup: controlKeyboard() }
792 );
793 }
794
795 async function sendThreads(chatId) {
796 const threads = await runtimeJson("/v1/threads/summary?limit=8&include_archived=true");
797 if (!threads.length) {
798 await sendText(chatId, "No runtime threads yet.", { replyMarkup: controlKeyboard() });
799 return;
800 }
801 const actions = [];
802 for (const [index, thread] of threads.slice(0, 8).entries()) {
803 const token = await threadStore.putAction(
804 { kind: "resume", threadId: thread.id },
805 { chatId }
806 );
807 actions.push({ token, label: `Resume ${index + 1}` });
808 }
809 await sendText(
810 chatId,
811 threads
812 .map((thread, index) => {
813 const status = thread.latest_turn_status || "none";
814 return `${index + 1}. ${thread.id} [${status}] ${thread.title || thread.preview || ""}`;
815 })
816 .join("\n"),
817 { replyMarkup: threadListKeyboard(actions) }
818 );
819 }
820
821 async function resumeThread(chatId, args) {
822 const threadId = args.trim();
823 if (!threadId) {
824 await sendText(chatId, "Usage: /resume <thread_id>");
825 return;
826 }
827 const detail = await runtimeJson(`/v1/threads/${encodeURIComponent(threadId)}`);
828 const existing = await threadStore.getChat(chatId);
829 await threadStore.setChat(chatId, {
830 ...preservedChatStateFields(existing),
831 threadId,
832 lastSeq: Number(detail.latest_seq || 0),
833 activeTurnId: null,
834 updatedAt: new Date().toISOString()
835 });
836 await sendText(chatId, `Resumed thread ${threadId}`, { replyMarkup: controlKeyboard() });
837 }
838
839 async function interruptActiveTurn(chatId) {
840 const state = await threadStore.getChat(chatId);
841 if (!state?.threadId) {
842 await sendText(chatId, "No runtime thread recorded for this chat.");
843 return;
844 }
845 const detail = await runtimeJson(`/v1/threads/${encodeURIComponent(state.threadId)}`);
846 const runningTurn = latestRunningTurn(detail);
847 const turnId = state.activeTurnId || runningTurn?.id;
848 if (!turnId) {
849 await sendText(chatId, "No active turn recorded for this chat.");
850 return;
851 }
852 await runtimeJson(
853 `/v1/threads/${encodeURIComponent(state.threadId)}/turns/${encodeURIComponent(
854 turnId
855 )}/interrupt`,
856 { method: "POST" }
857 );
858 await threadStore.patchChat(chatId, {
859 activeTurnId: turnId,
860 updatedAt: new Date().toISOString()
861 });
862 await sendText(chatId, `Interrupt requested for ${turnId}`, { replyMarkup: controlKeyboard() });
863 }
864
865 async function compactThread(chatId) {
866 const state = await ensureThread(chatId);
867 const result = await runtimeJson(`/v1/threads/${encodeURIComponent(state.threadId)}/compact`, {
868 method: "POST",
869 body: { reason: "telegram bridge request" }
870 });
871 await sendText(chatId, `Compaction started: ${result.turn?.id || "unknown turn"}`, {
872 replyMarkup: activeTurnKeyboard()
873 });
874 }
875
876 /** Deliver a decision for an approval pending on this chat's thread. */
877 async function decideApproval(chatId, action, identity) {
878 const decision = action.decision;
879 const { approvalId, remember } = action;
880 if (!approvalId) {
881 await sendText(
882 chatId,
883 `Usage: /${decision} <approval_id>${decision === "allow" ? " [remember]" : ""}`
884 );
885 return false;
886 }
887 try {
888 await decideRuntimeApproval(runtimeJson, { store: threadStore, chatId, actorId: turnActor(identity), approvalId, decision, remember, turnId: action.turnId });
889 } catch (error) {
890 if (!(error instanceof ApprovalOwnershipError)) throw error;
891 await sendText(chatId, `Approval ${approvalId} is not waiting in this chat.`);
892 return false;
893 }
894 await sendText(chatId, `Approval ${approvalId}: ${decision}${remember ? " and remember" : ""}`);
895 return true;
896 }
897
898 async function setChatModel(chatId, modelName) {
899 if (!modelName || modelName === "default") {
900 await threadStore.patchChat(chatId, {
901 model: null,
902 updatedAt: new Date().toISOString()
903 });
904 await sendText(chatId, `Reset per-chat model. Using bridge default: ${config.model}`, {
905 replyMarkup: controlKeyboard()
906 });
907 return;
908 }
909 await threadStore.patchChat(chatId, {
910 model: modelName,
911 updatedAt: new Date().toISOString()
912 });
913 await sendText(chatId, `Per-chat model set to: ${modelName}`, { replyMarkup: controlKeyboard() });
914 }
915
916 async function sendText(chatId, text, options = {}) {
917 const chunks = splitMessage(text, config.maxReplyChars);
918 for (const [index, chunk] of chunks.entries()) {
919 const body = {
920 chat_id: chatId,
921 ...telegramMessageBody(chunk, { markdown: true, maxChars: config.maxReplyChars }),
922 disable_web_page_preview: true
923 };
924 if (options.replyMarkup && index === chunks.length - 1) {
925 body.reply_markup = options.replyMarkup;
926 }
927 try {
928 await telegramApi("sendMessage", body);
929 } catch (error) {
930 if (!isTelegramMarkdownParseError(error)) throw error;
931 const fallbackBody = {
932 chat_id: chatId,
933 ...telegramMessageBody(chunk, { markdown: false, maxChars: config.maxReplyChars }),
934 disable_web_page_preview: true
935 };
936 if (options.replyMarkup && index === chunks.length - 1) {
937 fallbackBody.reply_markup = options.replyMarkup;
938 }
939 await telegramApi("sendMessage", fallbackBody);
940 }
941 }
942 }
943
944 async function sendTypingAction(chatId) {
945 const controller = new AbortController();
946 const timeout = setTimeout(() => controller.abort(), TYPING_TIMEOUT_MS);
947 try {
948 await telegramApi(
949 "sendChatAction",
950 {
951 chat_id: chatId,
952 action: "typing"
953 },
954 { signal: controller.signal }
955 );
956 } finally {
957 clearTimeout(timeout);
958 }
959 }
960
961 async function sendTurnText(chatId, text, options = {}) {
962 try {
963 await sendText(chatId, text, options);
964 } catch (error) {
965 console.error("failed to send Telegram turn update", error);
966 }
967 }
968
969 async function answerCallback(callbackQueryId, text = "") {
970 await telegramApi("answerCallbackQuery", {
971 callback_query_id: callbackQueryId,
972 text: text.slice(0, 200),
973 show_alert: false
974 });
975 }
976
977 async function editMessageReplyMarkup(chatId, messageId, replyMarkup) {
978 await telegramApi("editMessageReplyMarkup", {
979 chat_id: chatId,
980 message_id: messageId,
981 reply_markup: replyMarkup
982 });
983 }
984
985 async function telegramApi(method, body = {}, options = {}) {
986 for (let attempt = 0; ; attempt += 1) {
987 try {
988 return await telegramApiOnce(method, body, options);
989 } catch (error) {
990 const retryMs = method === "sendMessage" ? telegramSendRetryDelayMs(error, attempt) : null;
991 if (retryMs == null) throw error;
992 console.warn(
993 `Telegram ${method} failed: ${error.message}. Retrying in ${Math.round(retryMs / 1000)}s.`
994 );
995 await delay(retryMs);
996 }
997 }
998 }
999
1000 async function telegramApiOnce(method, body = {}, options = {}) {
1001 const response = await fetch(`${config.apiBase}/bot${config.botToken}/${method}`, {
1002 method: "POST",
1003 headers: { "content-type": "application/json" },
1004 body: JSON.stringify(body),
1005 signal: options.signal
1006 });
1007 const payload = await readJsonSafe(response);
1008 if (!response.ok || payload?.ok === false) {
1009 const error = new Error(
1010 payload?.description || `Telegram API request failed (${response.status})`
1011 );
1012 error.errorCode = payload?.error_code || response.status;
1013 error.description = payload?.description || "";
1014 error.parameters = payload?.parameters || {};
1015 throw error;
1016 }
1017 return payload.result;
1018 }
1019
1020 function requiredEnv(name) {
1021 const value = process.env[name];
1022 if (!value || !value.trim()) {
1023 throw new Error(`${name} is required`);
1024 }
1025 return value.trim();
1026 }
1027
1028 function requiredEnvFirst(...names) {
1029 const value = envFirst(process.env, ...names);
1030 if (!value) {
1031 throw new Error(`${names.join(" or ")} is required`);
1032 }
1033 return value;
1034 }
1035
1036 function delay(ms) {
1037 return new Promise((resolve) => setTimeout(resolve, ms));
1038 }
1039
1039 lines Plain Text