paperclip/patches/@chat-adapter__telegram@4.3...

239 lines
9.2 KiB
Diff

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<string>);
/** 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<string>;
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)) {