From c0f3318d9fda8584d1f8f604235ebdc813d1bdbb Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=A8=8B=E5=BA=8F=E5=91=98=E9=98=BF=E6=B1=9F=28Relakkes?= =?UTF-8?q?=29?= Date: Tue, 1 Sep 2026 21:22:43 +0800 Subject: [PATCH] fix(telegram): preserve complete streamed replies (#1264) --- .../__tests__/stream-delivery.test.ts | 198 ++++++++++++ adapters/telegram/format.ts | 39 ++- adapters/telegram/index.ts | 159 +-------- adapters/telegram/stream-delivery.ts | 303 ++++++++++++++++++ 4 files changed, 546 insertions(+), 153 deletions(-) create mode 100644 adapters/telegram/__tests__/stream-delivery.test.ts create mode 100644 adapters/telegram/stream-delivery.ts diff --git a/adapters/telegram/__tests__/stream-delivery.test.ts b/adapters/telegram/__tests__/stream-delivery.test.ts new file mode 100644 index 00000000..6432779a --- /dev/null +++ b/adapters/telegram/__tests__/stream-delivery.test.ts @@ -0,0 +1,198 @@ +import { describe, expect, it } from 'bun:test' +import { TelegramStreamDelivery, type TelegramDeliveryApi } from '../stream-delivery.js' + +type Failure = unknown | null + +function createFakeApi(options: { + editFailures?: Failure[] + sendFailure?: (text: string, attempt: number) => Failure +} = {}) { + let nextMessageId = 1 + let nonPlaceholderSendAttempts = 0 + const messages = new Map() + const editAttempts: string[] = [] + const sendAttempts: string[] = [] + const editFailures = [...(options.editFailures ?? [])] + + const api: TelegramDeliveryApi = { + async sendMessage(_chatId, text) { + sendAttempts.push(text) + if (text !== '▍') { + nonPlaceholderSendAttempts += 1 + const failure = options.sendFailure?.(text, nonPlaceholderSendAttempts) + if (failure) throw failure + } + const messageId = nextMessageId++ + messages.set(messageId, text) + return { message_id: messageId } + }, + async editMessageText(_chatId, messageId, text) { + editAttempts.push(text) + const failure = editFailures.shift() + if (failure) throw failure + messages.set(messageId, text) + return {} + }, + } + + return { api, messages, editAttempts, sendAttempts } +} + +function createDelivery( + api: TelegramDeliveryApi, + options: { + bufferCharThreshold?: number + delays?: number[] + warnings?: string[] + } = {}, +) { + return new TelegramStreamDelivery(api, { + bufferIntervalMs: 60_000, + bufferCharThreshold: options.bufferCharThreshold ?? 20_000, + sleep: async (ms) => { options.delays?.push(ms) }, + warn: (message) => { options.warnings?.push(message) }, + }) +} + +async function streamReply( + delivery: TelegramStreamDelivery, + chatId: string, + deltas: string[], +): Promise { + await delivery.handleEvent(chatId, { type: 'content_start', blockType: 'text' }) + for (const text of deltas) { + await delivery.handleEvent(chatId, { type: 'content_delta', text }) + } + await delivery.handleEvent(chatId, { type: 'message_complete' }) +} + +describe('TelegramStreamDelivery', () => { + it('delivers the complete reply through the real stream transition', async () => { + const fake = createFakeApi() + const delivery = createDelivery(fake.api) + const answer = '第一句完整。第二句也完整。第三句收尾。' + + await streamReply(delivery, '42', ['第一句完整。', '第二句也完整。', '第三句收尾。']) + + expect([...fake.messages.values()]).toEqual([answer]) + expect(delivery.hasState('42')).toBe(false) + }) + + it('preserves failed streaming deltas and falls back to a complete send when final edits fail', async () => { + const rateLimit = { error_code: 429, parameters: { retry_after: 1 } } + const fake = createFakeApi({ + editFailures: [null, rateLimit, rateLimit, rateLimit, rateLimit], + }) + const delays: number[] = [] + const warnings: string[] = [] + const delivery = createDelivery(fake.api, { + bufferCharThreshold: 1, + delays, + warnings, + }) + const answer = '第一句完整。第二句也完整。第三句收尾。' + + await delivery.handleEvent('42', { type: 'content_start', blockType: 'text' }) + await delivery.handleEvent('42', { type: 'content_delta', text: '第一句完整。' }) + await Bun.sleep(0) + await delivery.handleEvent('42', { type: 'content_delta', text: '第二句也完整。第三句收尾。' }) + await Bun.sleep(0) + + expect(delivery.getDeliveryState('42')).toEqual({ + desiredText: answer, + immutableOffset: 0, + placeholderDeliveredText: '第一句完整。', + finalChunksDelivered: 0, + }) + + await delivery.handleEvent('42', { type: 'message_complete' }) + + expect(fake.editAttempts).toHaveLength(5) + expect(fake.sendAttempts.at(-1)).toBe(answer) + expect(delays).toEqual([1000, 1000]) + expect(warnings.some((message) => message.includes('streaming update failed'))).toBe(true) + expect(warnings.some((message) => message.includes('falling back to sendMessage'))).toBe(true) + expect(delivery.hasState('42')).toBe(false) + }) + + it('respects retry_after before retrying a final edit', async () => { + const fake = createFakeApi({ + editFailures: [{ error_code: 429, parameters: { retry_after: 2 } }], + }) + const delays: number[] = [] + const delivery = createDelivery(fake.api, { delays }) + + await streamReply(delivery, '42', ['complete after retry']) + + expect(fake.editAttempts).toEqual(['complete after retry', 'complete after retry']) + expect(delays).toEqual([2000]) + expect([...fake.messages.values()]).toEqual(['complete after retry']) + }) + + it('retries a transient middle send without losing or duplicating long-reply chunks', async () => { + const answer = 'x'.repeat(9000) + let failed = false + const fake = createFakeApi({ + sendFailure(text) { + if (!failed && text.length === 4000) { + failed = true + return new TypeError('fetch failed') + } + return null + }, + }) + const delays: number[] = [] + const delivery = createDelivery(fake.api, { delays }) + + await streamReply(delivery, '42', [answer]) + + expect([...fake.messages.values()].join('')).toBe(answer) + expect(fake.sendAttempts.filter((text) => text.length === 4000)).toHaveLength(2) + expect(delays).toEqual([250]) + expect(delivery.hasState('42')).toBe(false) + }) + + it('resumes a long reply from the confirmed prefix after an intermediate send fails', async () => { + const answer = 'y'.repeat(9000) + let failed = false + const fake = createFakeApi({ + sendFailure() { + if (!failed) { + failed = true + return new TypeError('stream socket reset') + } + return null + }, + }) + const warnings: string[] = [] + const delivery = createDelivery(fake.api, { + bufferCharThreshold: 1, + warnings, + }) + + await delivery.handleEvent('42', { type: 'content_start', blockType: 'text' }) + await delivery.handleEvent('42', { type: 'content_delta', text: answer }) + await Bun.sleep(0) + await delivery.handleEvent('42', { type: 'message_complete' }) + + expect([...fake.messages.values()].join('')).toBe(answer) + expect(warnings.some((message) => message.includes('streaming update failed'))).toBe(true) + expect(delivery.hasState('42')).toBe(false) + }) + + it('retains delivery state when the final fallback cannot be sent', async () => { + const fake = createFakeApi({ + editFailures: [new Error('edit rejected')], + sendFailure: () => new Error('send rejected'), + }) + const delivery = createDelivery(fake.api) + + await delivery.handleEvent('42', { type: 'content_start', blockType: 'text' }) + await delivery.handleEvent('42', { type: 'content_delta', text: 'must survive' }) + + await expect( + delivery.handleEvent('42', { type: 'message_complete' }), + ).rejects.toThrow('send rejected') + expect(delivery.hasState('42')).toBe(true) + }) +}) diff --git a/adapters/telegram/format.ts b/adapters/telegram/format.ts index 4c9f5e50..fd3bfa95 100644 --- a/adapters/telegram/format.ts +++ b/adapters/telegram/format.ts @@ -36,6 +36,12 @@ export type TelegramStreamingUpdate = { activeChunk: string } +export type TelegramStreamingChunk = { + text: string + remainder: string + consumedLength: number +} + export function planTelegramStreamingUpdate( currentText: string, deltaText: string, @@ -50,9 +56,9 @@ export function planTelegramStreamingUpdate( let remaining = fullText while (formatTelegramOutboundText(remaining).length > limit) { - const [sealed, rest] = splitOneStreamingChunk(remaining, limit) - sealedChunks.push(sealed) - remaining = rest + const chunk = splitTelegramStreamingChunk(remaining, limit) + sealedChunks.push(chunk.text) + remaining = chunk.remainder if (!remaining) break } @@ -60,7 +66,10 @@ export function planTelegramStreamingUpdate( return { sealedChunks, activeChunk: remaining } } -function splitOneStreamingChunk(text: string, limit: number): [string, string] { +export function splitTelegramStreamingChunk( + text: string, + limit: number, +): TelegramStreamingChunk { const roughLimit = Math.min(limit, text.length) const candidates = [ text.lastIndexOf('\n\n', roughLimit), @@ -73,20 +82,36 @@ function splitOneStreamingChunk(text: string, limit: number): [string, string] { const splitAt = includeDelimiter(text, candidate) const sealed = text.slice(0, splitAt).trimEnd() if (sealed && formatTelegramOutboundText(sealed).length <= limit) { - return [sealed, text.slice(splitAt).trimStart()] + return buildStreamingChunk(text, sealed, splitAt) } } const chunks = splitMessage(formatTelegramOutboundText(text), limit) const firstFormattedChunk = chunks[0] ?? text.slice(0, limit) if (firstFormattedChunk.length < text.length && text.startsWith(firstFormattedChunk)) { - return [firstFormattedChunk.trimEnd(), text.slice(firstFormattedChunk.length).trimStart()] + return buildStreamingChunk(text, firstFormattedChunk.trimEnd(), firstFormattedChunk.length) } const splitAt = Math.max(1, Math.min(limit, text.length)) - return [text.slice(0, splitAt).trimEnd(), text.slice(splitAt).trimStart()] + return buildStreamingChunk(text, text.slice(0, splitAt).trimEnd(), splitAt) } function includeDelimiter(text: string, splitAt: number): number { return text[splitAt] === '\n' || text[splitAt] === '.' ? splitAt + 1 : splitAt } + +function buildStreamingChunk( + source: string, + text: string, + splitAt: number, +): TelegramStreamingChunk { + let consumedLength = splitAt + while (consumedLength < source.length && /\s/.test(source[consumedLength]!)) { + consumedLength += 1 + } + return { + text, + remainder: source.slice(consumedLength), + consumedLength, + } +} diff --git a/adapters/telegram/index.ts b/adapters/telegram/index.ts index 4dd81779..f1f9b4c6 100644 --- a/adapters/telegram/index.ts +++ b/adapters/telegram/index.ts @@ -8,21 +8,17 @@ import { Bot, InlineKeyboard, type Context } from 'grammy' import * as path from 'node:path' import { WsBridge, type ServerMessage } from '../common/ws-bridge.js' -import { MessageBuffer } from '../common/message-buffer.js' import { MessageDedup } from '../common/message-dedup.js' import { enqueue } from '../common/chat-queue.js' import { loadConfig } from '../common/config.js' import { formatImStatus, formatPermissionRequest, - splitMessage, } from '../common/format.js' import { buildTelegramThinkingUpdate, - formatTelegramOutboundText, - formatTelegramStreamingText, - planTelegramStreamingUpdate, } from './format.js' +import { TelegramStreamDelivery } from './stream-delivery.js' import { formatPermissionDecisionStatus, formatPermissionInstructions, @@ -44,9 +40,6 @@ import { sendSafeOutboundImage } from '../common/attachment/outbound-image.js' import { syncTelegramBotCommands } from './menu.js' import { createTelegramRuntimeCommandController, registerAuthorizedTelegramCommand, registerTelegramExtendedCommands, shouldProcessTelegramMessage, tryHandleTelegramSelectionCallback } from './commands.js' -const TELEGRAM_TEXT_LIMIT = 4000 // leave margin below 4096 -const TELEGRAM_STREAMING_TEXT_LIMIT = TELEGRAM_TEXT_LIMIT - 2 // reserve room for cursor - // ---------- init ---------- const config = loadConfig() @@ -57,6 +50,7 @@ if (!config.telegram.botToken) { const bot = new Bot(config.telegram.botToken) const bridge = new WsBridge(config.serverUrl, 'tg') +const streamDelivery = new TelegramStreamDelivery(bot.api) const dedup = new MessageDedup() const sessionStore = new SessionStore() const { httpClient, defaultWorkDir } = createAdapterClient(config, config.telegram) @@ -66,13 +60,7 @@ attachmentStore.gc().catch((err) => { console.warn('[Telegram] AttachmentStore.gc failed:', err instanceof Error ? err.message : err) }) -// Track placeholder messages for streaming updates -const placeholders = new Map() -// Track accumulated text per chat for streaming -const accumulatedText = new Map() const accumulatedThinkingText = new Map() -// Message buffers per chat -const buffers = new Map() // Track chats waiting for project selection const pendingProjectSelection = new Map() const runtimeStates = new Map() @@ -100,17 +88,6 @@ const commandController = createTelegramRuntimeCommandController({ botApi: bot.a // ---------- helpers ---------- -function getBuffer(chatId: string): MessageBuffer { - let buf = buffers.get(chatId) - if (!buf) { - buf = new MessageBuffer(async (text, isComplete) => { - await flushToTelegram(chatId, text, isComplete) - }) - buffers.set(chatId, buf) - } - return buf -} - function getRuntimeState(chatId: string): ChatRuntimeState { let state = runtimeStates.get(chatId) if (!state) { @@ -121,10 +98,8 @@ function getRuntimeState(chatId: string): ChatRuntimeState { } function clearTransientChatState(chatId: string): void { - placeholders.delete(chatId) - accumulatedText.delete(chatId) + streamDelivery.clear(chatId) accumulatedThinkingText.delete(chatId) - buffers.get(chatId)?.reset() const runtime = getRuntimeState(chatId) runtime.state = 'idle' runtime.verb = undefined @@ -215,72 +190,6 @@ async function buildStatusText(chatId: string): Promise { }) } -async function flushToTelegram(chatId: string, newText: string, isComplete: boolean): Promise { - const numericChatId = Number(chatId) - const prev = accumulatedText.get(chatId) ?? '' - - const placeholder = placeholders.get(chatId) - - if (placeholder) { - if (isComplete) { - const fullText = prev + newText - accumulatedText.set(chatId, fullText) - const chunks = splitMessage(formatTelegramOutboundText(fullText), TELEGRAM_TEXT_LIMIT) - try { - await bot.api.editMessageText(numericChatId, placeholder.messageId, chunks[0]!) - } catch { /* ignore */ } - for (let i = 1; i < chunks.length; i++) { - await bot.api.sendMessage(numericChatId, chunks[i]!) - } - } else { - const { sealedChunks, activeChunk } = planTelegramStreamingUpdate( - prev, - newText, - TELEGRAM_STREAMING_TEXT_LIMIT, - ) - accumulatedText.set(chatId, activeChunk) - try { - const firstSealedChunk = sealedChunks.shift() - if (firstSealedChunk) { - const firstSealedFormattedChunks = splitMessage( - formatTelegramOutboundText(firstSealedChunk), - TELEGRAM_TEXT_LIMIT, - ) - await bot.api.editMessageText(numericChatId, placeholder.messageId, firstSealedFormattedChunks[0]!) - for (let i = 1; i < firstSealedFormattedChunks.length; i++) { - await bot.api.sendMessage(numericChatId, firstSealedFormattedChunks[i]!) - } - for (const chunk of sealedChunks) { - const formattedChunks = splitMessage(formatTelegramOutboundText(chunk), TELEGRAM_TEXT_LIMIT) - for (const formattedChunk of formattedChunks) { - await bot.api.sendMessage(numericChatId, formattedChunk) - } - } - const sent = await bot.api.sendMessage(numericChatId, formatTelegramStreamingText(activeChunk)) - placeholders.set(chatId, { chatId, messageId: sent.message_id }) - } else { - await bot.api.editMessageText(numericChatId, placeholder.messageId, formatTelegramStreamingText(activeChunk)) - } - } catch { /* ignore */ } - } - } else if (isComplete && (prev + newText).trim()) { - const fullText = prev + newText - accumulatedText.set(chatId, fullText) - const chunks = splitMessage(formatTelegramOutboundText(fullText), TELEGRAM_TEXT_LIMIT) - for (const chunk of chunks) { - await bot.api.sendMessage(numericChatId, chunk) - } - } else { - accumulatedText.set(chatId, prev + newText) - } - - if (isComplete) { - placeholders.delete(chatId) - accumulatedText.delete(chatId) - buffers.get(chatId)?.reset() - } -} - // ---------- session management ---------- async function ensureSession(chatId: string): Promise { @@ -365,7 +274,6 @@ async function dispatchOutboundMedia(chatId: string, pending: PendingUpload): Pr async function handleServerMessage(chatId: string, msg: ServerMessage): Promise { const numericChatId = Number(chatId) - const buf = getBuffer(chatId) const runtime = getRuntimeState(chatId) switch (msg.type) { @@ -375,10 +283,8 @@ async function handleServerMessage(chatId: string, msg: ServerMessage): Promise< case 'status': runtime.state = msg.state runtime.verb = typeof msg.verb === 'string' ? msg.verb : undefined - if (msg.state === 'thinking' && !placeholders.has(chatId)) { - const sent = await bot.api.sendMessage(numericChatId, '💭 思考中...') - placeholders.set(chatId, { chatId, messageId: sent.message_id }) - accumulatedText.set(chatId, '') + if (msg.state === 'thinking' && !streamDelivery.hasState(chatId)) { + await streamDelivery.ensurePlaceholder(chatId, '💭 思考中...') accumulatedThinkingText.set(chatId, '') } break @@ -386,38 +292,18 @@ async function handleServerMessage(chatId: string, msg: ServerMessage): Promise< case 'content_start': if (msg.blockType === 'text') { accumulatedThinkingText.delete(chatId) - if (!placeholders.has(chatId)) { - const sent = await bot.api.sendMessage(numericChatId, '▍') - placeholders.set(chatId, { chatId, messageId: sent.message_id }) - accumulatedText.set(chatId, '') - } + await streamDelivery.handleEvent(chatId, { type: 'content_start', blockType: msg.blockType }) } else if (msg.blockType === 'tool_use') { // Finalize current text placeholder before tool calls, // so text after tools gets a fresh message - await buf.complete() - // If placeholder still exists (buffer was already empty), clean up directly - if (placeholders.has(chatId)) { - const text = accumulatedText.get(chatId) - if (text?.trim()) { - try { - await bot.api.editMessageText( - numericChatId, - placeholders.get(chatId)!.messageId, - formatTelegramOutboundText(text), - ) - } catch { /* ignore */ } - } - placeholders.delete(chatId) - accumulatedText.delete(chatId) - buffers.get(chatId)?.reset() - } + await streamDelivery.complete(chatId) } break case 'content_delta': if (msg.text) { accumulatedThinkingText.delete(chatId) - buf.append(msg.text) + await streamDelivery.handleEvent(chatId, { type: 'content_delta', text: msg.text }) const newUploads = getTgWatcher(chatId).feed(msg.text) for (const pending of newUploads) { void dispatchOutboundMedia(chatId, pending) @@ -426,7 +312,7 @@ async function handleServerMessage(chatId: string, msg: ServerMessage): Promise< break case 'thinking': - if (placeholders.has(chatId)) { + if (streamDelivery.getPlaceholderMessageId(chatId) !== undefined) { const update = buildTelegramThinkingUpdate( accumulatedThinkingText.get(chatId) ?? '', msg.text, @@ -435,7 +321,7 @@ async function handleServerMessage(chatId: string, msg: ServerMessage): Promise< try { await bot.api.editMessageText( numericChatId, - placeholders.get(chatId)!.messageId, + streamDelivery.getPlaceholderMessageId(chatId)!, update.messageText, ) } catch { /* ignore */ } @@ -470,24 +356,8 @@ async function handleServerMessage(chatId: string, msg: ServerMessage): Promise< case 'message_complete': runtime.state = 'idle' runtime.verb = undefined - await buf.complete() - // Ensure placeholder is always cleaned up even if buffer was already empty - if (placeholders.has(chatId)) { - const text = accumulatedText.get(chatId) - if (text?.trim()) { - try { - const chunks = splitMessage(formatTelegramOutboundText(text), TELEGRAM_TEXT_LIMIT) - await bot.api.editMessageText(numericChatId, placeholders.get(chatId)!.messageId, chunks[0]!) - for (let i = 1; i < chunks.length; i++) { - await bot.api.sendMessage(numericChatId, chunks[i]!) - } - } catch { /* ignore */ } - } - placeholders.delete(chatId) - accumulatedText.delete(chatId) - accumulatedThinkingText.delete(chatId) - buffers.get(chatId)?.reset() - } + await streamDelivery.handleEvent(chatId, { type: 'message_complete' }) + accumulatedThinkingText.delete(chatId) break case 'error': @@ -541,10 +411,7 @@ async function startNewSession(chatId: string, query?: string): Promise { bridge.resetSession(chatId) sessionStore.delete(chatId) - placeholders.delete(chatId) - accumulatedText.delete(chatId) - buffers.get(chatId)?.reset() - buffers.delete(chatId) + streamDelivery.clear(chatId) pendingProjectSelection.delete(chatId) commandController.clearPendingSelections(chatId) pendingPermissions.delete(chatId) diff --git a/adapters/telegram/stream-delivery.ts b/adapters/telegram/stream-delivery.ts new file mode 100644 index 00000000..23730a69 --- /dev/null +++ b/adapters/telegram/stream-delivery.ts @@ -0,0 +1,303 @@ +import { MessageBuffer } from '../common/message-buffer.js' +import { splitMessage } from '../common/format.js' +import { + formatTelegramOutboundText, + formatTelegramStreamingText, + splitTelegramStreamingChunk, +} from './format.js' + +const TELEGRAM_TEXT_LIMIT = 4000 +const TELEGRAM_STREAMING_TEXT_LIMIT = TELEGRAM_TEXT_LIMIT - 2 +const DEFAULT_MAX_ATTEMPTS = 3 +const DEFAULT_NETWORK_RETRY_MS = 250 + +export type TelegramDeliveryApi = { + sendMessage: (chatId: number, text: string) => Promise<{ message_id: number }> + editMessageText: (chatId: number, messageId: number, text: string) => Promise +} + +export type TelegramStreamEvent = + | { type: 'content_start'; blockType: string } + | { type: 'content_delta'; text?: string } + | { type: 'message_complete' } + +type FinalDeliveryPlan = { + chunks: string[] + mode: 'edit' | 'send' + nextChunkIndex: number +} + +type DeliveryState = { + fullText: string + immutableOffset: number + placeholderMessageId?: number + placeholderDeliveredText: string + finalPlan?: FinalDeliveryPlan +} + +export type TelegramStreamDeliveryOptions = { + bufferIntervalMs?: number + bufferCharThreshold?: number + maxAttempts?: number + sleep?: (ms: number) => Promise + warn?: (message: string) => void +} + +export type TelegramDeliveryStateSnapshot = { + desiredText: string + immutableOffset: number + placeholderDeliveredText: string + finalChunksDelivered: number +} + +export class TelegramStreamDelivery { + private readonly states = new Map() + private readonly buffers = new Map() + private readonly maxAttempts: number + private readonly sleep: (ms: number) => Promise + private readonly warn: (message: string) => void + + constructor( + private readonly api: TelegramDeliveryApi, + private readonly options: TelegramStreamDeliveryOptions = {}, + ) { + this.maxAttempts = Math.max(1, options.maxAttempts ?? DEFAULT_MAX_ATTEMPTS) + this.sleep = options.sleep ?? (async (ms) => { + await new Promise((resolve) => setTimeout(resolve, ms)) + }) + this.warn = options.warn ?? ((message) => console.warn(message)) + } + + async handleEvent(chatId: string, event: TelegramStreamEvent): Promise { + switch (event.type) { + case 'content_start': + if (event.blockType === 'text') await this.ensurePlaceholder(chatId, '▍') + break + case 'content_delta': + if (event.text) await this.append(chatId, event.text) + break + case 'message_complete': + await this.complete(chatId) + break + } + } + + hasState(chatId: string): boolean { + return this.states.has(chatId) + } + + getPlaceholderMessageId(chatId: string): number | undefined { + return this.states.get(chatId)?.placeholderMessageId + } + + getDeliveryState(chatId: string): TelegramDeliveryStateSnapshot | undefined { + const state = this.states.get(chatId) + if (!state) return undefined + return { + desiredText: state.fullText, + immutableOffset: state.immutableOffset, + placeholderDeliveredText: state.placeholderDeliveredText, + finalChunksDelivered: state.finalPlan?.nextChunkIndex ?? 0, + } + } + + async ensurePlaceholder(chatId: string, text: string): Promise { + const existing = this.states.get(chatId)?.placeholderMessageId + if (existing !== undefined) return existing + + const sent = await this.api.sendMessage(Number(chatId), text) + this.states.set(chatId, { + fullText: '', + immutableOffset: 0, + placeholderMessageId: sent.message_id, + placeholderDeliveredText: '', + }) + this.getBuffer(chatId) + return sent.message_id + } + + async append(chatId: string, text: string): Promise { + if (!this.states.has(chatId)) await this.ensurePlaceholder(chatId, '▍') + this.getBuffer(chatId).append(text) + } + + async complete(chatId: string): Promise { + const buffer = this.buffers.get(chatId) + if (buffer) await buffer.complete() + + const state = this.states.get(chatId) + if (!state) return + if (!state.fullText.trim()) { + this.clear(chatId) + return + } + + await this.deliverFinal(chatId, state) + this.clear(chatId) + } + + clear(chatId: string): void { + this.buffers.get(chatId)?.reset() + this.buffers.delete(chatId) + this.states.delete(chatId) + } + + private getBuffer(chatId: string): MessageBuffer { + let buffer = this.buffers.get(chatId) + if (!buffer) { + buffer = new MessageBuffer( + async (text, isComplete) => { + await this.acceptBufferedText(chatId, text, isComplete) + }, + this.options.bufferIntervalMs, + this.options.bufferCharThreshold, + ) + this.buffers.set(chatId, buffer) + } + return buffer + } + + private async acceptBufferedText( + chatId: string, + text: string, + isComplete: boolean, + ): Promise { + const state = this.states.get(chatId) + if (!state) return + + state.fullText += text + if (isComplete) return + + try { + await this.deliverStreaming(chatId, state) + } catch (err) { + this.warn(`[Telegram] streaming update failed for ${chatId}; retaining buffered text: ${formatError(err)}`) + } + } + + private async deliverStreaming(chatId: string, state: DeliveryState): Promise { + const numericChatId = Number(chatId) + + while (true) { + const remaining = state.fullText.slice(state.immutableOffset) + if (formatTelegramOutboundText(remaining).length <= TELEGRAM_STREAMING_TEXT_LIMIT) { + const messageText = formatTelegramStreamingText(remaining) + if (state.placeholderMessageId === undefined) { + const sent = await this.api.sendMessage(numericChatId, messageText) + state.placeholderMessageId = sent.message_id + } else { + await this.api.editMessageText(numericChatId, state.placeholderMessageId, messageText) + } + state.placeholderDeliveredText = remaining + return + } + + const sealed = splitTelegramStreamingChunk(remaining, TELEGRAM_STREAMING_TEXT_LIMIT) + const messageText = formatTelegramOutboundText(sealed.text) + if (state.placeholderMessageId === undefined) { + await this.api.sendMessage(numericChatId, messageText) + } else { + await this.api.editMessageText(numericChatId, state.placeholderMessageId, messageText) + } + + state.immutableOffset += sealed.consumedLength + state.placeholderMessageId = undefined + state.placeholderDeliveredText = '' + } + } + + private async deliverFinal(chatId: string, state: DeliveryState): Promise { + const numericChatId = Number(chatId) + if (!state.finalPlan) { + const remaining = state.fullText.slice(state.immutableOffset) + const chunks = splitMessage(formatTelegramOutboundText(remaining), TELEGRAM_TEXT_LIMIT) + state.finalPlan = { + chunks, + mode: state.placeholderMessageId === undefined ? 'send' : 'edit', + nextChunkIndex: 0, + } + } + const plan = state.finalPlan + + if (plan.mode === 'edit' && plan.nextChunkIndex === 0) { + try { + await this.withFinalRetry( + () => this.api.editMessageText( + numericChatId, + state.placeholderMessageId!, + plan.chunks[0]!, + ), + ) + state.placeholderDeliveredText = plan.chunks[0]! + plan.nextChunkIndex = 1 + } catch (err) { + this.warn(`[Telegram] final edit failed for ${chatId}; falling back to sendMessage: ${formatError(err)}`) + plan.mode = 'send' + } + } + + while (plan.nextChunkIndex < plan.chunks.length) { + const chunk = plan.chunks[plan.nextChunkIndex]! + await this.withFinalRetry(() => this.api.sendMessage(numericChatId, chunk)) + plan.nextChunkIndex += 1 + } + } + + private async withFinalRetry(operation: () => Promise): Promise { + let attempt = 1 + while (true) { + try { + return await operation() + } catch (err) { + const retryDelayMs = getRetryDelayMs(err, attempt) + if (attempt >= this.maxAttempts || retryDelayMs === null) throw err + await this.sleep(retryDelayMs) + attempt += 1 + } + } + } +} + +function getRetryDelayMs(err: unknown, attempt: number): number | null { + if (isRecord(err)) { + const errorCode = numberField(err, 'error_code') ?? numberField(err, 'status') ?? numberField(err, 'statusCode') + if (errorCode === 429) { + const parameters = isRecord(err.parameters) ? err.parameters : undefined + const retryAfterSeconds = parameters ? numberField(parameters, 'retry_after') : undefined + return Math.max(0, retryAfterSeconds ?? 1) * 1000 + } + } + + if (err instanceof TypeError) return DEFAULT_NETWORK_RETRY_MS * 2 ** (attempt - 1) + + const code = isRecord(err) && typeof err.code === 'string' ? err.code : '' + if (/^(ECONNRESET|ECONNREFUSED|EPIPE|ETIMEDOUT|EAI_AGAIN|UND_ERR_)/.test(code)) { + return DEFAULT_NETWORK_RETRY_MS * 2 ** (attempt - 1) + } + + const name = err instanceof Error ? err.name : '' + const message = formatError(err) + if (name === 'HttpError' || /fetch failed|network error|socket hang up/i.test(message)) { + return DEFAULT_NETWORK_RETRY_MS * 2 ** (attempt - 1) + } + + return null +} + +function isRecord(value: unknown): value is Record { + return typeof value === 'object' && value !== null +} + +function numberField(value: Record, key: string): number | undefined { + const field = value[key] + return typeof field === 'number' ? field : undefined +} + +function formatError(err: unknown): string { + if (err instanceof Error) return err.message + try { + return JSON.stringify(err) + } catch { + return String(err) + } +}