Skip to content

Commit f876955

Browse files
committed
fix(telegram): honor long flood-wait retry_after on outbound sends
1 parent fd18e39 commit f876955

7 files changed

Lines changed: 113 additions & 3 deletions

File tree

extensions/telegram/src/bot/delivery.send.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@ import {
1313
isTelegramQuoteParamError,
1414
removeTelegramNativeQuoteParam,
1515
} from "../reply-parameters.js";
16+
import { TELEGRAM_OUTBOUND_RETRY_AFTER_CAP_MS } from "../retry-after.js";
1617
import {
1718
buildTelegramRichMessage,
1819
getTelegramRichRawApi,
@@ -33,6 +34,7 @@ function createTelegramDeliverySendRetry() {
3334
return createTelegramRetryRunner({
3435
shouldRetry: (err) => isSafeToRetrySendError(err) || isTelegramRateLimitError(err),
3536
strictShouldRetry: true,
37+
retryAfterMaxDelayMs: TELEGRAM_OUTBOUND_RETRY_AFTER_CAP_MS,
3638
});
3739
}
3840

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,2 @@
1+
// Matches draft preview suspension: final Telegram replies should wait through routine flood windows.
2+
export const TELEGRAM_OUTBOUND_RETRY_AFTER_CAP_MS = 60_000;

extensions/telegram/src/send.test.ts

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2219,6 +2219,42 @@ describe("sendMessageTelegram", () => {
22192219
vi.useRealTimers();
22202220
});
22212221

2222+
it("honors long Telegram retry_after hints above the default send retry cap", async () => {
2223+
vi.useFakeTimers();
2224+
const chatId = "123";
2225+
const sendMessage = vi
2226+
.fn()
2227+
.mockRejectedValueOnce({
2228+
message: "429 Too Many Requests",
2229+
response: { parameters: { retry_after: 45 } },
2230+
})
2231+
.mockResolvedValueOnce({
2232+
message_id: 2,
2233+
chat: { id: chatId },
2234+
});
2235+
const api = { sendMessage } as unknown as {
2236+
sendMessage: typeof sendMessage;
2237+
};
2238+
const setTimeoutSpy = vi.spyOn(global, "setTimeout");
2239+
2240+
const promise = sendMessageTelegram(chatId, "hi", {
2241+
cfg: TELEGRAM_TEST_CFG,
2242+
token: "tok",
2243+
api,
2244+
retry: { attempts: 2, minDelayMs: 0, maxDelayMs: 30_000, jitter: 0 },
2245+
});
2246+
2247+
await vi.advanceTimersByTimeAsync(44_999);
2248+
expect(sendMessage).toHaveBeenCalledTimes(1);
2249+
2250+
await vi.advanceTimersByTimeAsync(1);
2251+
await expect(promise).resolves.toEqual({ messageId: "2", chatId });
2252+
expect(firstMockCall(setTimeoutSpy, "setTimeout call")[1]).toBe(45_000);
2253+
expect(sendMessage).toHaveBeenCalledTimes(2);
2254+
setTimeoutSpy.mockRestore();
2255+
vi.useRealTimers();
2256+
});
2257+
22222258
it("retries wrapped pre-connect HttpError sends", async () => {
22232259
vi.useFakeTimers();
22242260
const chatId = "123";

extensions/telegram/src/send.ts

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,7 @@ import {
4848
removeTelegramNativeQuoteParam,
4949
resolveTelegramSendThreadSpec,
5050
} from "./reply-parameters.js";
51+
import { TELEGRAM_OUTBOUND_RETRY_AFTER_CAP_MS } from "./retry-after.js";
5152
import {
5253
buildTelegramRichMessage,
5354
getTelegramRichRawApi,
@@ -659,6 +660,7 @@ function createTelegramRequestWithDiag(params: {
659660
account: ResolvedTelegramAccount;
660661
retry?: RetryConfig;
661662
verbose?: boolean;
663+
retryAfterMaxDelayMs?: number;
662664
shouldRetry?: (err: unknown) => boolean;
663665
/** When true, the shouldRetry predicate is used exclusively without the TELEGRAM_RETRY_RE fallback. */
664666
strictShouldRetry?: boolean;
@@ -668,6 +670,9 @@ function createTelegramRequestWithDiag(params: {
668670
retry: params.retry,
669671
configRetry: params.account.config.retry,
670672
verbose: params.verbose,
673+
...(params.retryAfterMaxDelayMs !== undefined
674+
? { retryAfterMaxDelayMs: params.retryAfterMaxDelayMs }
675+
: {}),
671676
...(params.shouldRetry ? { shouldRetry: params.shouldRetry } : {}),
672677
...(params.strictShouldRetry ? { strictShouldRetry: true } : {}),
673678
});
@@ -747,6 +752,7 @@ function createTelegramNonIdempotentRequestWithDiag(params: {
747752
retry: params.retry,
748753
verbose: params.verbose,
749754
useApiErrorLogging: params.useApiErrorLogging,
755+
retryAfterMaxDelayMs: TELEGRAM_OUTBOUND_RETRY_AFTER_CAP_MS,
750756
shouldRetry: (err) => isSafeToRetrySendError(err) || isTelegramRateLimitError(err),
751757
strictShouldRetry: true,
752758
});

src/infra/retry-policy.test.ts

Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -203,4 +203,55 @@ describe("createChannelApiRetryRunner", () => {
203203
await expect(promise).resolves.toBe("ok");
204204
expect(fn).toHaveBeenCalledTimes(2);
205205
});
206+
207+
it("keeps retry_after hints capped by maxDelayMs by default", async () => {
208+
vi.useFakeTimers();
209+
210+
const runner = createChannelApiRetryRunner({
211+
retry: { attempts: 2, minDelayMs: 0, maxDelayMs: 30_000, jitter: 0 },
212+
});
213+
const fn = vi
214+
.fn()
215+
.mockRejectedValueOnce({
216+
message: "429 Too Many Requests",
217+
response: { parameters: { retry_after: 45 } },
218+
})
219+
.mockResolvedValue("ok");
220+
221+
const promise = runner(fn, "test");
222+
223+
expect(fn).toHaveBeenCalledTimes(1);
224+
await vi.advanceTimersByTimeAsync(29_999);
225+
expect(fn).toHaveBeenCalledTimes(1);
226+
227+
await vi.advanceTimersByTimeAsync(1);
228+
await expect(promise).resolves.toBe("ok");
229+
expect(fn).toHaveBeenCalledTimes(2);
230+
});
231+
232+
it("honors retry_after above maxDelayMs when a separate retry-after cap is configured", async () => {
233+
vi.useFakeTimers();
234+
235+
const runner = createChannelApiRetryRunner({
236+
retry: { attempts: 2, minDelayMs: 0, maxDelayMs: 30_000, jitter: 0 },
237+
retryAfterMaxDelayMs: 60_000,
238+
});
239+
const fn = vi
240+
.fn()
241+
.mockRejectedValueOnce({
242+
message: "429 Too Many Requests",
243+
response: { parameters: { retry_after: 45 } },
244+
})
245+
.mockResolvedValue("ok");
246+
247+
const promise = runner(fn, "test");
248+
249+
expect(fn).toHaveBeenCalledTimes(1);
250+
await vi.advanceTimersByTimeAsync(44_999);
251+
expect(fn).toHaveBeenCalledTimes(1);
252+
253+
await vi.advanceTimersByTimeAsync(1);
254+
await expect(promise).resolves.toBe("ok");
255+
expect(fn).toHaveBeenCalledTimes(2);
256+
});
206257
});

src/infra/retry-policy.ts

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -95,6 +95,7 @@ export function createChannelApiRetryRunner(params: {
9595
retry?: RetryConfig;
9696
configRetry?: RetryConfig;
9797
verbose?: boolean;
98+
retryAfterMaxDelayMs?: number;
9899
shouldRetry?: (err: unknown) => boolean;
99100
/**
100101
* When true, the custom shouldRetry predicate is used exclusively —
@@ -116,6 +117,9 @@ export function createChannelApiRetryRunner(params: {
116117
label,
117118
shouldRetry,
118119
retryAfterMs: getChannelApiRetryAfterMs,
120+
...(params.retryAfterMaxDelayMs !== undefined
121+
? { retryAfterMaxDelayMs: params.retryAfterMaxDelayMs }
122+
: {}),
119123
onRetry: params.verbose
120124
? (info) => {
121125
const maxRetries = Math.max(1, info.maxAttempts - 1);

src/infra/retry.ts

Lines changed: 12 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@ export type RetryOptions = RetryConfig & {
2727
label?: string;
2828
shouldRetry?: (err: unknown, attempt: number) => boolean;
2929
retryAfterMs?: (err: unknown) => number | undefined;
30+
retryAfterMaxDelayMs?: number;
3031
onRetry?: (info: RetryInfo) => void;
3132
};
3233

@@ -136,6 +137,13 @@ export async function retryAsync<T>(
136137
Number.isFinite(resolved.maxDelayMs) && resolved.maxDelayMs > 0
137138
? resolved.maxDelayMs
138139
: Number.POSITIVE_INFINITY;
140+
const retryAfterMaxDelayMs =
141+
options.retryAfterMaxDelayMs === undefined
142+
? maxDelayMs
143+
: Math.max(
144+
minDelayMs,
145+
resolveRetryDelayMs(Math.round(clampNumber(options.retryAfterMaxDelayMs, maxDelayMs, 0))),
146+
);
139147
const jitter = resolved.jitter;
140148
const shouldRetry = options.shouldRetry ?? (() => true);
141149
let lastErr: unknown;
@@ -154,7 +162,8 @@ export async function retryAsync<T>(
154162
const baseDelay = hasRetryAfter
155163
? Math.max(retryAfterMs, minDelayMs)
156164
: minDelayMs * 2 ** (attempt - 1);
157-
let delay = Math.min(baseDelay, maxDelayMs);
165+
const delayCap = hasRetryAfter ? retryAfterMaxDelayMs : maxDelayMs;
166+
let delay = Math.min(baseDelay, delayCap);
158167
// Server-supplied Retry-After is a lower-bound contract with the
159168
// upstream rate limiter; symmetric jitter would let roughly half the
160169
// retries land before the requested time and invite escalation. Use
@@ -181,9 +190,9 @@ export async function retryAsync<T>(
181190
// (`retryAfterMs > maxDelayMs`), where the contract is already
182191
// unsatisfiable and we gain spread without adding a violation.
183192
const canHonorRetryAfter =
184-
hasRetryAfter && typeof retryAfterMs === "number" && retryAfterMs <= maxDelayMs;
193+
hasRetryAfter && typeof retryAfterMs === "number" && retryAfterMs <= delayCap;
185194
delay = applyJitter(delay, jitter, canHonorRetryAfter ? "positive" : "symmetric");
186-
delay = Math.min(Math.max(delay, minDelayMs), maxDelayMs);
195+
delay = Math.min(Math.max(delay, minDelayMs), delayCap);
187196

188197
options.onRetry?.({
189198
attempt,

0 commit comments

Comments
 (0)