diff --git a/dist/index.d.ts b/dist/index.d.ts --- a/dist/index.d.ts +++ b/dist/index.d.ts @@ -53,6 +53,8 @@ botToken?: string | (() => string | Promise); /** Logger instance for error reporting. Defaults to ConsoleLogger. */ logger?: Logger; + /** Maximum bytes buffered for one downloaded Telegram attachment. Defaults to Telegram's 25 MB ceiling. */ + maxDownloadBytes?: number; /** Optional long-polling configuration for getUpdates flow. */ longPolling?: TelegramLongPollingConfig; /** @@ -646,6 +648,7 @@ protected readonly botTokenProvider: () => Promise; protected readonly staticBotToken?: string; protected readonly apiBaseUrl: string; + protected readonly maxDownloadBytes: number; protected readonly secretToken?: string; protected readonly mentionOnReply: boolean; private botIdentityPromise; @@ -669,6 +672,8 @@ private pollingTask; private pollingActive; private nextDraftId; + /** Paperclip durable private-draft stop handshake; not task cancellation. */ + readonly paperclipDraftStopVersion: 1; private richMessagesAvailable; get botUserId(): string | undefined; get userName(): string; diff --git a/dist/index.js b/dist/index.js --- a/dist/index.js +++ b/dist/index.js @@ -771,9 +771,10 @@ var TELEGRAM_API_BASE = "https://api.telegram.org"; var TELEGRAM_FILE_LIMIT = 25 * 1024 * 1024; var TELEGRAM_FILE_TIMEOUT_MS = 3e4; -async function readTelegramFile(response, fileId) { +var TELEGRAM_API_TIMEOUT_MS = 45e3; +async function readTelegramFile(response, fileId, maxDownloadBytes = TELEGRAM_FILE_LIMIT) { const declared = Number(response.headers.get("content-length")); - if (Number.isFinite(declared) && declared > TELEGRAM_FILE_LIMIT) { + if (Number.isFinite(declared) && declared > maxDownloadBytes) { await response.body?.cancel(); throw new NetworkError( "telegram", @@ -792,7 +793,7 @@ break; } size += value.length; - if (size > TELEGRAM_FILE_LIMIT) { + if (size > maxDownloadBytes) { await reader.cancel(); throw new NetworkError( "telegram", @@ -989,6 +990,7 @@ botTokenProvider; staticBotToken; apiBaseUrl; + maxDownloadBytes; secretToken; mentionOnReply; botIdentityPromise = null; @@ -1012,6 +1014,7 @@ pollingTask = null; pollingActive = false; nextDraftId = Math.max(1, Date.now() % 2147483647); + paperclipDraftStopVersion = 1; richMessagesAvailable = true; get botUserId() { return this._botUserId; @@ -1038,6 +1041,12 @@ this.apiBaseUrl = trimTrailingSlashes( config.apiUrl ?? config.apiBaseUrl ?? process.env.TELEGRAM_API_BASE_URL ?? TELEGRAM_API_BASE ); + this.maxDownloadBytes = this.clampInteger( + config.maxDownloadBytes, + TELEGRAM_FILE_LIMIT, + 1, + TELEGRAM_FILE_LIMIT + ); this.secretToken = config.secretToken ?? process.env.TELEGRAM_WEBHOOK_SECRET_TOKEN; this.allowUnverifiedWebhooks = config.allowUnverifiedWebhooks ?? process.env.TELEGRAM_ALLOW_UNVERIFIED_WEBHOOKS === "true"; this.mentionOnReply = config.mentionOnReply ?? process.env.TELEGRAM_MENTION_ON_REPLY === "true"; @@ -1949,6 +1958,9 @@ }); } async stream(threadId, textStream, options) { + if (options?.paperclipDraftControl && (!this.nativeStreaming || !this.isDM(threadId))) { + throw new ValidationError2("telegram", "Durable draft control requires native private streaming"); + } if (this.nativeStreaming && this.isDM(threadId)) { return await this.nativeDraftStream(threadId, textStream, options); } @@ -2089,6 +2101,29 @@ } async nativeDraftStream(threadId, textStream, options) { const parsedThread = this.resolveThreadId(threadId); + const draftControl = options?.paperclipDraftControl; + if (draftControl && (draftControl.version !== 1 || !Number.isSafeInteger(draftControl.draftId) || draftControl.draftId <= 0 || draftControl.draftId > 2147483647 || typeof draftControl.beforeDraft !== "function" || typeof draftControl.beforeFinal !== "function")) { + throw new ValidationError2("telegram", "Invalid durable draft control"); + } + let draftStopped = false; + let draftControlFailed = false; + const sendControlledDraft = async (method, payload) => { + if (draftControl) { + try { + if (!await draftControl.beforeDraft()) { + draftStopped = true; + return; + } + } catch (error) { + draftControlFailed = true; + throw error; + } + } + return await this.telegramFetch(method, { + ...payload, + ...draftControl ? { can_stop: true, keep_on_stop: false } : {} + }); + }; const updateIntervalMs = this.clampInteger( options?.updateIntervalMs, TELEGRAM_DEFAULT_STREAM_UPDATE_INTERVAL_MS, @@ -2096,7 +2131,7 @@ Number.MAX_SAFE_INTEGER ); const renderer = new StreamingMarkdownRenderer2(); - const draftId = this.createDraftId(); + const draftId = draftControl?.draftId ?? this.createDraftId(); let accumulated = ""; let lastDraftText = null; let lastFlushAt = 0; @@ -2124,7 +2159,7 @@ let draftText = text2; if (streamUsesRich) { try { - await this.telegramFetch("sendRichMessageDraft", { + await sendControlledDraft("sendRichMessageDraft", { chat_id: parsedThread.chatId, message_thread_id: parsedThread.messageThreadId, draft_id: draftId, @@ -2136,6 +2171,7 @@ lastFlushAt = Date.now(); return; } catch (error) { + if (draftControlFailed) throw error; if (!this.canFallbackFromRichMessage(error, "sendRichMessageDraft")) { draftStreamingEnabled = false; this.logger.warn("Telegram rich draft streaming update failed", { @@ -2158,7 +2194,7 @@ } try { if (useMarkdown) { - await this.telegramFetch("sendMessageDraft", { + await sendControlledDraft("sendMessageDraft", { chat_id: parsedThread.chatId, message_thread_id: parsedThread.messageThreadId, draft_id: draftId, @@ -2166,7 +2202,7 @@ parse_mode: toBotApiParseMode("MarkdownV2") }); } else { - await this.telegramFetch("sendMessageDraft", { + await sendControlledDraft("sendMessageDraft", { chat_id: parsedThread.chatId, message_thread_id: parsedThread.messageThreadId, draft_id: draftId, @@ -2176,11 +2212,12 @@ lastDraftText = draftText; lastFlushAt = Date.now(); } catch (error) { + if (draftControlFailed) throw error; if (useMarkdown && this.isTelegramMarkdownParseError(error)) { streamUsesMarkdown = false; const plainDraftText = renderPlainText(accumulated); try { - await this.telegramFetch("sendMessageDraft", { + await sendControlledDraft("sendMessageDraft", { chat_id: parsedThread.chatId, message_thread_id: parsedThread.messageThreadId, draft_id: draftId, @@ -2189,6 +2226,7 @@ lastDraftText = plainDraftText; lastFlushAt = Date.now(); } catch (retryError) { + if (draftControlFailed) throw retryError; draftStreamingEnabled = false; this.logger.warn("Telegram draft streaming update failed", { error: String(retryError), @@ -2230,6 +2268,7 @@ renderer.push(text2); if (Date.now() - lastFlushAt >= updateIntervalMs) { await flushDraft(); + if (draftStopped) return { paperclipDraftStopped: true }; } } if (!accumulated.trim()) { @@ -2240,6 +2279,9 @@ } const finalMarkdown = renderer.finish(); await flushDraft(); + if (draftStopped || draftControl && !await draftControl.beforeFinal()) { + return { paperclipDraftStopped: true }; + } if (streamUsesRich) { const markdown = truncateRichMarkdown(finalMarkdown); try { @@ -2625,7 +2667,7 @@ ); } try { - return await readTelegramFile(response, fileId); + return await readTelegramFile(response, fileId, this.maxDownloadBytes); } catch (error) { if (error instanceof NetworkError) { throw error; @@ -3439,6 +3481,12 @@ async telegramFetch(method, payload, request) { const botToken = this.staticBotToken ?? await this.resolveBotToken(); const url = `${this.apiBaseUrl}/bot${botToken}/${method}`; + const requestTimeoutMs = method === "getUpdates" && payload && typeof payload === "object" && typeof payload.timeout === "number" ? Math.max(TELEGRAM_API_TIMEOUT_MS, (payload.timeout + 5) * 1e3) : TELEGRAM_API_TIMEOUT_MS; + const timeoutSignal = AbortSignal.timeout(requestTimeoutMs); + const signal = request?.signal ? AbortSignal.any([ + request.signal, + timeoutSignal + ]) : timeoutSignal; let response; try { response = await fetch(url, { @@ -3447,7 +3495,7 @@ "Content-Type": "application/json" }, body: payload instanceof FormData ? payload : JSON.stringify(payload ?? {}), - signal: request?.signal + signal }); } catch (error) { if (this.isAbortError(error)) {