fix(telegram): preserve complete streamed replies (#1264)

This commit is contained in:
程序员阿江(Relakkes)
2026-09-01 21:22:43 +08:00
parent b1ecce7812
commit c0f3318d9f
4 changed files with 546 additions and 153 deletions
@@ -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<number, string>()
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<void> {
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)
})
})
+32 -7
View File
@@ -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,
}
}
+13 -146
View File
@@ -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<string, { chatId: string; messageId: number }>()
// Track accumulated text per chat for streaming
const accumulatedText = new Map<string, string>()
const accumulatedThinkingText = new Map<string, string>()
// Message buffers per chat
const buffers = new Map<string, MessageBuffer>()
// Track chats waiting for project selection
const pendingProjectSelection = new Map<string, boolean>()
const runtimeStates = new Map<string, ChatRuntimeState>()
@@ -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<string> {
})
}
async function flushToTelegram(chatId: string, newText: string, isComplete: boolean): Promise<void> {
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<boolean> {
@@ -365,7 +274,6 @@ async function dispatchOutboundMedia(chatId: string, pending: PendingUpload): Pr
async function handleServerMessage(chatId: string, msg: ServerMessage): Promise<void> {
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<void> {
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)
+303
View File
@@ -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<unknown>
}
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<void>
warn?: (message: string) => void
}
export type TelegramDeliveryStateSnapshot = {
desiredText: string
immutableOffset: number
placeholderDeliveredText: string
finalChunksDelivered: number
}
export class TelegramStreamDelivery {
private readonly states = new Map<string, DeliveryState>()
private readonly buffers = new Map<string, MessageBuffer>()
private readonly maxAttempts: number
private readonly sleep: (ms: number) => Promise<void>
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<void>((resolve) => setTimeout(resolve, ms))
})
this.warn = options.warn ?? ((message) => console.warn(message))
}
async handleEvent(chatId: string, event: TelegramStreamEvent): Promise<void> {
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<number> {
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<void> {
if (!this.states.has(chatId)) await this.ensurePlaceholder(chatId, '▍')
this.getBuffer(chatId).append(text)
}
async complete(chatId: string): Promise<void> {
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<void> {
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<void> {
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<void> {
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<T>(operation: () => Promise<T>): Promise<T> {
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<string, unknown> {
return typeof value === 'object' && value !== null
}
function numberField(value: Record<string, unknown>, 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)
}
}