Skip to content

Commit 7270bb9

Browse files
authored
fix(telegram): drain outbound queue after polling reconnect (#82227)
* fix(telegram): drain outbound queue after polling reconnect * fix(telegram): keep reconnect drain backoff-safe
1 parent 6ca9de1 commit 7270bb9

3 files changed

Lines changed: 200 additions & 1 deletion

File tree

CHANGELOG.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@ Docs: https://docs.openclaw.ai
1414
- Telegram/active-memory: run blocking memory recall through the Telegram provider for direct-message turns even when the hook context carries the raw chat id, preventing embedded recall from launching against an invalid numeric channel. Fixes #82177. Thanks @cslash-zz.
1515
- Control UI/WebChat: keep optimistic image messages from embedding large inline `data:` previews and preserve image-only user turns in chat history, avoiding browser stack overflows when sending image attachments. Fixes #82182. Thanks @ExploreSheep.
1616
- 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.
17+
- 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.
1718
- Agents: strip Gemini/Gemma `<final>` tags with attributes or self-closing syntax from delivered replies, including strict final-tag streaming enforcement. Fixes #65867.
1819
- macOS/update: disarm legacy `ai.openclaw.update.*` LaunchAgents when `openclaw update` starts from one, preventing KeepAlive relaunch loops that repeatedly restart the Gateway and replay update continuations. Fixes #82167.
1920
- Agents/replay: strip internal runtime-context metadata and `NO_REPLY` sentinels from provider replay and pending final-delivery recovery so restart and heartbeat resumes do not feed control text back to the model. Fixes #76629. Thanks @fuyizheng3120, @bryan-chx, and @cael-dandelion-cult.

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: false,
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: 40 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,10 @@ 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+
4752
type TelegramBot = ReturnType<typeof createTelegramBot>;
4853

4954
const waitForGracefulStop = async (stop: () => Promise<void>) => {
@@ -114,6 +119,7 @@ export class TelegramPollingSession {
114119
#transportState: TelegramPollingTransportState;
115120
#status: ReturnType<typeof createTelegramPollingStatusPublisher>;
116121
#stallThresholdMs: number;
122+
#deliveryDrainInFlight = false;
117123

118124
constructor(private readonly opts: TelegramPollingSessionOpts) {
119125
this.#transportState = new TelegramPollingTransportState({
@@ -202,6 +208,35 @@ export class TelegramPollingSession {
202208
);
203209
}
204210

211+
#drainPendingDeliveriesAfterReconnect() {
212+
if (this.#deliveryDrainInFlight) {
213+
return;
214+
}
215+
this.#deliveryDrainInFlight = true;
216+
const accountId = normalizeTelegramAccountId(this.opts.accountId);
217+
void drainPendingDeliveries({
218+
drainKey: `telegram:${accountId}`,
219+
logLabel: "Telegram reconnect drain",
220+
cfg: this.opts.config,
221+
log: {
222+
info: (message) => this.opts.log(`[telegram][diag] ${message}`),
223+
warn: (message) => this.opts.log(`[telegram] ${message}`),
224+
error: (message) => this.opts.log(`[telegram] ${message}`),
225+
},
226+
selectEntry: (entry) => ({
227+
match:
228+
entry.channel === "telegram" && normalizeTelegramAccountId(entry.accountId) === accountId,
229+
bypassBackoff: false,
230+
}),
231+
})
232+
.catch((err) => {
233+
this.opts.log(`[telegram] reconnect delivery drain failed: ${formatErrorMessage(err)}`);
234+
})
235+
.finally(() => {
236+
this.#deliveryDrainInFlight = false;
237+
});
238+
}
239+
205240
async #createPollingBot(): Promise<TelegramBot | undefined> {
206241
const fetchAbortController = new AbortController();
207242
this.#activeFetchAbort = fetchAbortController;
@@ -359,6 +394,7 @@ export class TelegramPollingSession {
359394
}
360395
if (message.type === "poll-success") {
361396
this.#status.notePollSuccess(message.finishedAt);
397+
this.#drainPendingDeliveriesAfterReconnect();
362398
pollState.outcome = `ok:${message.count}`;
363399
return;
364400
}
@@ -428,7 +464,10 @@ export class TelegramPollingSession {
428464

429465
async #runPollingCycle(bot: TelegramBot): Promise<"continue" | "exit"> {
430466
const liveness = new TelegramPollingLivenessTracker({
431-
onPollSuccess: (finishedAt) => this.#status.notePollSuccess(finishedAt),
467+
onPollSuccess: (finishedAt) => {
468+
this.#status.notePollSuccess(finishedAt);
469+
this.#drainPendingDeliveriesAfterReconnect();
470+
},
432471
});
433472
bot.api.config.use(async (prev, method, payload, signal) => {
434473
if (method !== "getUpdates") {

0 commit comments

Comments
 (0)