Skip to content

Commit 3b7f3db

Browse files
committed
fix(telegram): drain outbound queue after polling reconnect
1 parent c6ddb1a commit 3b7f3db

3 files changed

Lines changed: 210 additions & 1 deletion

File tree

CHANGELOG.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ Docs: https://docs.openclaw.ai
1212
### Fixes
1313

1414
- Agents/media: preserve message-tool-only delivery for generated music and video completion handoffs, so group/channel completions do not finish without posting the generated attachment.
15+
- Telegram: drain queued outbound deliveries after polling reconnect confirms fresh `getUpdates` activity, so stale-socket and network recovery do not leave failed replies stranded. Fixes #50040. Refs #82175. Thanks @dmitriiforpost-commits and @shellyrocklobster.
1516
- LINE: acknowledge signed webhook events before agent processing so slow model replies do not cause LINE `request_timeout` delivery failures. Fixes #65375. Thanks @myericho.
1617
- TTS: preserve channel-derived voice-note delivery for `/tts audio` replies even when the provider output is not natively voice-compatible. (#82174) Thanks @xuruiray.
1718
- Codex/Lossless: keep Codex explicit compaction on native app-server threads while allowing Lossless through the context-engine slot; `openclaw doctor --fix` now migrates legacy `compaction.provider: "lossless-claw"` config to `plugins.slots.contextEngine`.

extensions/telegram/src/polling-session.test.ts

Lines changed: 159 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ const createTelegramBotMock = vi.hoisted(() => vi.fn());
99
const isRecoverableTelegramNetworkErrorMock = vi.hoisted(() => vi.fn(() => true));
1010
const computeBackoffMock = vi.hoisted(() => vi.fn(() => 0));
1111
const sleepWithAbortMock = vi.hoisted(() => vi.fn(async () => undefined));
12+
const drainPendingDeliveriesMock = vi.hoisted(() => vi.fn(async (_opts: unknown) => undefined));
1213

1314
vi.mock("@grammyjs/runner", () => ({
1415
run: runMock,
@@ -22,6 +23,10 @@ vi.mock("./network-errors.js", () => ({
2223
isRecoverableTelegramNetworkError: isRecoverableTelegramNetworkErrorMock,
2324
}));
2425

26+
vi.mock("openclaw/plugin-sdk/delivery-queue-runtime", () => ({
27+
drainPendingDeliveries: drainPendingDeliveriesMock,
28+
}));
29+
2530
vi.mock("./api-logging.js", () => ({
2631
withTelegramApiErrorLogging: async ({ fn }: { fn: () => Promise<unknown> }) => await fn(),
2732
}));
@@ -54,6 +59,18 @@ type TelegramApiMiddleware = (
5459
method: string,
5560
payload: unknown,
5661
) => Promise<unknown>;
62+
type DrainPendingDeliveriesCall = {
63+
drainKey: string;
64+
logLabel: string;
65+
selectEntry: (
66+
entry: {
67+
channel: string;
68+
accountId?: string;
69+
lastError?: string;
70+
},
71+
now: number,
72+
) => { match: boolean; bypassBackoff: boolean };
73+
};
5774
type AsyncVoidFn = () => Promise<void>;
5875
type MockCallSource = { mock: { calls: Array<Array<unknown>> } };
5976

@@ -164,6 +181,14 @@ function expectTelegramBotTransportSequence(firstTransport: unknown, secondTrans
164181
expect(createTelegramBotMock.mock.calls.at(1)?.[0]?.telegramTransport).toBe(secondTransport);
165182
}
166183

184+
function expectDrainPendingDeliveriesCall(index = 0): DrainPendingDeliveriesCall {
185+
const call = drainPendingDeliveriesMock.mock.calls[index]?.[0];
186+
if (!call || typeof call !== "object") {
187+
throw new Error(`Expected drainPendingDeliveries call ${index}`);
188+
}
189+
return call as DrainPendingDeliveriesCall;
190+
}
191+
167192
function makeTelegramTransport() {
168193
return {
169194
fetch: globalThis.fetch,
@@ -292,6 +317,7 @@ describe("TelegramPollingSession", () => {
292317
isRecoverableTelegramNetworkErrorMock.mockReset().mockReturnValue(true);
293318
computeBackoffMock.mockReset().mockReturnValue(0);
294319
sleepWithAbortMock.mockReset().mockResolvedValue(undefined);
320+
drainPendingDeliveriesMock.mockReset().mockResolvedValue(undefined);
295321
});
296322

297323
it("uses backoff helpers for recoverable polling retries", async () => {
@@ -454,6 +480,58 @@ describe("TelegramPollingSession", () => {
454480
}
455481
});
456482

483+
it("drains Telegram delivery queue after isolated ingress reports poll success", async () => {
484+
const abort = new AbortController();
485+
const init = vi.fn(async () => undefined);
486+
const bot = {
487+
api: {
488+
deleteWebhook: vi.fn(async () => true),
489+
config: { use: vi.fn() },
490+
},
491+
init,
492+
handleUpdate: vi.fn(async () => undefined),
493+
stop: vi.fn(async () => undefined),
494+
};
495+
createTelegramBotMock.mockReturnValueOnce(bot);
496+
let onMessage:
497+
| ((message: { type: "poll-success"; finishedAt: number; count: number }) => void)
498+
| undefined;
499+
let stopWorker: (() => void) | undefined;
500+
const workerDone = new Promise<void>((resolve) => {
501+
stopWorker = resolve;
502+
});
503+
const createWorker = vi.fn(() => ({
504+
onMessage: vi.fn((handler) => {
505+
onMessage = handler;
506+
return () => undefined;
507+
}),
508+
stop: vi.fn(async () => {
509+
stopWorker?.();
510+
}),
511+
task: vi.fn(async () => {
512+
await workerDone;
513+
}),
514+
}));
515+
516+
const session = createPollingSession({
517+
abortSignal: abort.signal,
518+
isolatedIngress: {
519+
enabled: true,
520+
createWorker,
521+
drainIntervalMs: 10,
522+
},
523+
});
524+
525+
const runPromise = session.runUntilAbort();
526+
await vi.waitFor(() => expect(init).toHaveBeenCalledTimes(1));
527+
onMessage?.({ type: "poll-success", finishedAt: Date.now(), count: 0 });
528+
529+
await vi.waitFor(() => expect(drainPendingDeliveriesMock).toHaveBeenCalledTimes(1));
530+
531+
abort.abort();
532+
await runPromise;
533+
});
534+
457535
it("lets isolated ingress drain interleave different Telegram topic lanes", async () => {
458536
const abort = new AbortController();
459537
const tempDir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-telegram-spool-"));
@@ -1175,6 +1253,87 @@ describe("TelegramPollingSession", () => {
11751253
});
11761254
});
11771255

1256+
it("drains Telegram delivery queue after getUpdates confirms polling reconnect", async () => {
1257+
const abort = new AbortController();
1258+
const botStop = vi.fn(async () => undefined);
1259+
const runnerStop = vi.fn(async () => undefined);
1260+
const getApiMiddleware = mockBotCapturingApiMiddleware(botStop);
1261+
const resolveFirstTask = mockLongRunningPollingCycle(runnerStop);
1262+
1263+
const session = createPollingSession({
1264+
abortSignal: abort.signal,
1265+
});
1266+
1267+
const runPromise = session.runUntilAbort();
1268+
const apiMiddleware = await waitForApiMiddleware(getApiMiddleware);
1269+
await apiMiddleware(
1270+
vi.fn(async () => []),
1271+
"getUpdates",
1272+
{ offset: 1 },
1273+
);
1274+
1275+
await vi.waitFor(() => expect(drainPendingDeliveriesMock).toHaveBeenCalledTimes(1));
1276+
const drain = expectDrainPendingDeliveriesCall();
1277+
expect(drain.drainKey).toBe("telegram:default");
1278+
expect(drain.logLabel).toBe("Telegram reconnect drain");
1279+
expect(drain.selectEntry({ channel: "telegram" }, Date.now())).toEqual({
1280+
match: true,
1281+
bypassBackoff: false,
1282+
});
1283+
expect(
1284+
drain.selectEntry(
1285+
{
1286+
channel: "telegram",
1287+
accountId: "default",
1288+
lastError: "Network request for 'sendMessage' failed!",
1289+
},
1290+
Date.now(),
1291+
),
1292+
).toEqual({
1293+
match: true,
1294+
bypassBackoff: true,
1295+
});
1296+
expect(drain.selectEntry({ channel: "telegram", accountId: "alerts" }, Date.now()).match).toBe(
1297+
false,
1298+
);
1299+
expect(drain.selectEntry({ channel: "whatsapp" }, Date.now()).match).toBe(false);
1300+
1301+
abort.abort();
1302+
resolveFirstTask();
1303+
await runPromise;
1304+
});
1305+
1306+
it("drains Telegram delivery queue after each getUpdates success", async () => {
1307+
const abort = new AbortController();
1308+
const botStop = vi.fn(async () => undefined);
1309+
const runnerStop = vi.fn(async () => undefined);
1310+
const getApiMiddleware = mockBotCapturingApiMiddleware(botStop);
1311+
const resolveFirstTask = mockLongRunningPollingCycle(runnerStop);
1312+
1313+
const session = createPollingSession({
1314+
abortSignal: abort.signal,
1315+
});
1316+
1317+
const runPromise = session.runUntilAbort();
1318+
const apiMiddleware = await waitForApiMiddleware(getApiMiddleware);
1319+
await apiMiddleware(
1320+
vi.fn(async () => []),
1321+
"getUpdates",
1322+
{ offset: 1 },
1323+
);
1324+
await apiMiddleware(
1325+
vi.fn(async () => []),
1326+
"getUpdates",
1327+
{ offset: 2 },
1328+
);
1329+
1330+
await vi.waitFor(() => expect(drainPendingDeliveriesMock).toHaveBeenCalledTimes(2));
1331+
1332+
abort.abort();
1333+
resolveFirstTask();
1334+
await runPromise;
1335+
});
1336+
11781337
it("keeps polling marked connected across recoverable restart cycles", async () => {
11791338
const abort = new AbortController();
11801339
const recoverableError = new Error("recoverable polling error");

extensions/telegram/src/polling-session.ts

Lines changed: 50 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import { type RunOptions, run } from "@grammyjs/runner";
22
import type { ChannelAccountSnapshot } from "openclaw/plugin-sdk/channel-contract";
33
import type { TelegramNetworkConfig } from "openclaw/plugin-sdk/config-contracts";
4+
import { drainPendingDeliveries } from "openclaw/plugin-sdk/delivery-queue-runtime";
45
import {
56
computeBackoff,
67
formatDurationPrecise,
@@ -44,6 +45,20 @@ const TELEGRAM_POLLING_CLIENT_TIMEOUT_FLOOR_SECONDS = Math.ceil(
4445
TELEGRAM_GET_UPDATES_REQUEST_TIMEOUT_MS / 1000,
4546
);
4647

48+
function normalizeTelegramAccountId(accountId?: string | null): string {
49+
return accountId?.trim() || "default";
50+
}
51+
52+
function shouldBypassReconnectBackoff(lastError?: string): boolean {
53+
if (!lastError) {
54+
return false;
55+
}
56+
return isRecoverableTelegramNetworkError(new Error(lastError), {
57+
context: "send",
58+
allowMessageMatch: true,
59+
});
60+
}
61+
4762
type TelegramBot = ReturnType<typeof createTelegramBot>;
4863

4964
const waitForGracefulStop = async (stop: () => Promise<void>) => {
@@ -114,6 +129,7 @@ export class TelegramPollingSession {
114129
#transportState: TelegramPollingTransportState;
115130
#status: ReturnType<typeof createTelegramPollingStatusPublisher>;
116131
#stallThresholdMs: number;
132+
#deliveryDrainInFlight = false;
117133

118134
constructor(private readonly opts: TelegramPollingSessionOpts) {
119135
this.#transportState = new TelegramPollingTransportState({
@@ -202,6 +218,35 @@ export class TelegramPollingSession {
202218
);
203219
}
204220

221+
#drainPendingDeliveriesAfterReconnect() {
222+
if (this.#deliveryDrainInFlight) {
223+
return;
224+
}
225+
this.#deliveryDrainInFlight = true;
226+
const accountId = normalizeTelegramAccountId(this.opts.accountId);
227+
void drainPendingDeliveries({
228+
drainKey: `telegram:${accountId}`,
229+
logLabel: "Telegram reconnect drain",
230+
cfg: this.opts.config,
231+
log: {
232+
info: (message) => this.opts.log(`[telegram][diag] ${message}`),
233+
warn: (message) => this.opts.log(`[telegram] ${message}`),
234+
error: (message) => this.opts.log(`[telegram] ${message}`),
235+
},
236+
selectEntry: (entry) => ({
237+
match:
238+
entry.channel === "telegram" && normalizeTelegramAccountId(entry.accountId) === accountId,
239+
bypassBackoff: shouldBypassReconnectBackoff(entry.lastError),
240+
}),
241+
})
242+
.catch((err) => {
243+
this.opts.log(`[telegram] reconnect delivery drain failed: ${formatErrorMessage(err)}`);
244+
})
245+
.finally(() => {
246+
this.#deliveryDrainInFlight = false;
247+
});
248+
}
249+
205250
async #createPollingBot(): Promise<TelegramBot | undefined> {
206251
const fetchAbortController = new AbortController();
207252
this.#activeFetchAbort = fetchAbortController;
@@ -359,6 +404,7 @@ export class TelegramPollingSession {
359404
}
360405
if (message.type === "poll-success") {
361406
this.#status.notePollSuccess(message.finishedAt);
407+
this.#drainPendingDeliveriesAfterReconnect();
362408
pollState.outcome = `ok:${message.count}`;
363409
return;
364410
}
@@ -428,7 +474,10 @@ export class TelegramPollingSession {
428474

429475
async #runPollingCycle(bot: TelegramBot): Promise<"continue" | "exit"> {
430476
const liveness = new TelegramPollingLivenessTracker({
431-
onPollSuccess: (finishedAt) => this.#status.notePollSuccess(finishedAt),
477+
onPollSuccess: (finishedAt) => {
478+
this.#status.notePollSuccess(finishedAt);
479+
this.#drainPendingDeliveriesAfterReconnect();
480+
},
432481
});
433482
bot.api.config.use(async (prev, method, payload, signal) => {
434483
if (method !== "getUpdates") {

0 commit comments

Comments
 (0)