Skip to content

Commit 7c9f901

Browse files
committed
fix: surface stalled telegram ingress backlog
1 parent 4add9ba commit 7c9f901

6 files changed

Lines changed: 450 additions & 20 deletions

File tree

CHANGELOG.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ Docs: https://docs.openclaw.ai
2323
- 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.
2424
- 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.
2525
- 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.
26+
- Telegram: mark isolated polling ingress unhealthy when a spooled inbound backlog stalls while Bot API polling still succeeds, so gateway/channel health no longer stays green after Telegram DM processing wedges. Fixes #82175. Thanks @shellyrocklobster.
2627
- Agents: strip Gemini/Gemma `<final>` tags with attributes or self-closing syntax from delivered replies, including strict final-tag streaming enforcement. Fixes #65867.
2728
- 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.
2829
- 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: 309 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -71,6 +71,12 @@ type DrainPendingDeliveriesCall = {
7171
now: number,
7272
) => { match: boolean; bypassBackoff: boolean };
7373
};
74+
type WorkerPollSuccessListener = (message: {
75+
type: "poll-success";
76+
offset: null;
77+
count: number;
78+
finishedAt: number;
79+
}) => void;
7480
type AsyncVoidFn = () => Promise<void>;
7581
type MockCallSource = { mock: { calls: Array<Array<unknown>> } };
7682

@@ -882,6 +888,309 @@ describe("TelegramPollingSession", () => {
882888
}
883889
});
884890

891+
it("keeps active spooled lanes blocked across isolated ingress restarts", async () => {
892+
vi.useFakeTimers({ shouldAdvanceTime: true });
893+
const abort = new AbortController();
894+
const tempDir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-telegram-spool-"));
895+
let releaseRegularTurn: (() => void) | undefined;
896+
const regularTurnDone = new Promise<void>((resolve) => {
897+
releaseRegularTurn = resolve;
898+
});
899+
const handleUpdate = vi.fn(async () => {
900+
await regularTurnDone;
901+
});
902+
createTelegramBotMock.mockImplementation(() => ({
903+
api: {
904+
deleteWebhook: vi.fn(async () => true),
905+
config: { use: vi.fn() },
906+
},
907+
init: vi.fn(async () => undefined),
908+
handleUpdate,
909+
stop: vi.fn(async () => undefined),
910+
}));
911+
await writeTelegramSpooledUpdate({
912+
spoolDir: tempDir,
913+
update: {
914+
update_id: 42,
915+
message: { text: "summarize this", chat: { id: -100, type: "supergroup" } },
916+
},
917+
});
918+
919+
let workerTaskCalls = 0;
920+
let stopWorker: (() => void) | undefined;
921+
const workerDone = new Promise<void>((resolve) => {
922+
stopWorker = resolve;
923+
});
924+
const createWorker = vi.fn(() => ({
925+
onMessage: vi.fn(() => () => undefined),
926+
stop: vi.fn(async () => {
927+
stopWorker?.();
928+
}),
929+
task: vi.fn(async () => {
930+
workerTaskCalls += 1;
931+
if (workerTaskCalls === 1) {
932+
return;
933+
}
934+
await workerDone;
935+
}),
936+
}));
937+
938+
try {
939+
const session = createPollingSession({
940+
abortSignal: abort.signal,
941+
isolatedIngress: {
942+
enabled: true,
943+
spoolDir: tempDir,
944+
createWorker,
945+
drainIntervalMs: 100,
946+
},
947+
});
948+
949+
const runPromise = session.runUntilAbort();
950+
await vi.waitFor(() => expect(handleUpdate).toHaveBeenCalledTimes(1));
951+
await vi.advanceTimersByTimeAsync(16_000);
952+
await vi.waitFor(() => expect(createWorker).toHaveBeenCalledTimes(2));
953+
expect(handleUpdate).toHaveBeenCalledTimes(1);
954+
955+
releaseRegularTurn?.();
956+
await vi.advanceTimersByTimeAsync(1_000);
957+
await vi.waitFor(async () =>
958+
expect(
959+
(await listTelegramSpooledUpdates({ spoolDir: tempDir })).map(
960+
(update) => update.updateId,
961+
),
962+
).toEqual([]),
963+
);
964+
abort.abort();
965+
await vi.advanceTimersByTimeAsync(20_000);
966+
await runPromise;
967+
} finally {
968+
releaseRegularTurn?.();
969+
vi.useRealTimers();
970+
await fs.rm(tempDir, { recursive: true, force: true });
971+
}
972+
});
973+
974+
it("keeps active spooled lanes blocked across account restarts", async () => {
975+
vi.useFakeTimers({ shouldAdvanceTime: true });
976+
const firstAbort = new AbortController();
977+
const secondAbort = new AbortController();
978+
const tempDir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-telegram-spool-"));
979+
let releaseRegularTurn: (() => void) | undefined;
980+
const regularTurnDone = new Promise<void>((resolve) => {
981+
releaseRegularTurn = resolve;
982+
});
983+
const handleUpdate = vi.fn(async () => {
984+
await regularTurnDone;
985+
});
986+
createTelegramBotMock.mockImplementation(() => ({
987+
api: {
988+
deleteWebhook: vi.fn(async () => true),
989+
config: { use: vi.fn() },
990+
},
991+
init: vi.fn(async () => undefined),
992+
handleUpdate,
993+
stop: vi.fn(async () => undefined),
994+
}));
995+
await writeTelegramSpooledUpdate({
996+
spoolDir: tempDir,
997+
update: {
998+
update_id: 42,
999+
message: { text: "summarize this", chat: { id: -100, type: "supergroup" } },
1000+
},
1001+
});
1002+
1003+
const createWorker = vi.fn(() => {
1004+
let stopWorker: (() => void) | undefined;
1005+
const workerDone = new Promise<void>((resolve) => {
1006+
stopWorker = resolve;
1007+
});
1008+
return {
1009+
onMessage: vi.fn(() => () => undefined),
1010+
stop: vi.fn(async () => {
1011+
stopWorker?.();
1012+
}),
1013+
task: vi.fn(async () => {
1014+
await workerDone;
1015+
}),
1016+
};
1017+
});
1018+
1019+
try {
1020+
const firstSession = createPollingSession({
1021+
abortSignal: firstAbort.signal,
1022+
isolatedIngress: {
1023+
enabled: true,
1024+
spoolDir: tempDir,
1025+
createWorker,
1026+
drainIntervalMs: 100,
1027+
},
1028+
});
1029+
1030+
const firstRunPromise = firstSession.runUntilAbort();
1031+
await vi.waitFor(() => expect(handleUpdate).toHaveBeenCalledTimes(1));
1032+
firstAbort.abort();
1033+
await vi.advanceTimersByTimeAsync(16_000);
1034+
await firstRunPromise;
1035+
1036+
const secondSession = createPollingSession({
1037+
abortSignal: secondAbort.signal,
1038+
isolatedIngress: {
1039+
enabled: true,
1040+
spoolDir: tempDir,
1041+
createWorker,
1042+
drainIntervalMs: 100,
1043+
},
1044+
});
1045+
const secondRunPromise = secondSession.runUntilAbort();
1046+
await vi.waitFor(() => expect(createWorker).toHaveBeenCalledTimes(2));
1047+
await vi.advanceTimersByTimeAsync(1_000);
1048+
expect(handleUpdate).toHaveBeenCalledTimes(1);
1049+
1050+
releaseRegularTurn?.();
1051+
await vi.advanceTimersByTimeAsync(1_000);
1052+
await vi.waitFor(async () =>
1053+
expect(
1054+
(await listTelegramSpooledUpdates({ spoolDir: tempDir })).map(
1055+
(update) => update.updateId,
1056+
),
1057+
).toEqual([]),
1058+
);
1059+
secondAbort.abort();
1060+
await vi.advanceTimersByTimeAsync(20_000);
1061+
await secondRunPromise;
1062+
} finally {
1063+
releaseRegularTurn?.();
1064+
vi.useRealTimers();
1065+
await fs.rm(tempDir, { recursive: true, force: true });
1066+
}
1067+
});
1068+
1069+
it("marks isolated ingress unhealthy when a spooled backlog wedges while polling stays live", async () => {
1070+
vi.useFakeTimers({ shouldAdvanceTime: true });
1071+
const abort = new AbortController();
1072+
const tempDir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-telegram-spool-"));
1073+
const log = vi.fn();
1074+
const setStatus = vi.fn();
1075+
let releaseRegularTurn: (() => void) | undefined;
1076+
const regularTurnDone = new Promise<void>((resolve) => {
1077+
releaseRegularTurn = resolve;
1078+
});
1079+
const handleUpdate = vi.fn(async () => {
1080+
await regularTurnDone;
1081+
});
1082+
createTelegramBotMock.mockImplementation(() => ({
1083+
api: {
1084+
deleteWebhook: vi.fn(async () => true),
1085+
config: { use: vi.fn() },
1086+
},
1087+
init: vi.fn(async () => undefined),
1088+
handleUpdate,
1089+
stop: vi.fn(async () => undefined),
1090+
}));
1091+
for (const updateId of [42, 43]) {
1092+
await writeTelegramSpooledUpdate({
1093+
spoolDir: tempDir,
1094+
update: {
1095+
update_id: updateId,
1096+
message: { text: `dm ${updateId}`, chat: { id: 123, type: "private" } },
1097+
},
1098+
});
1099+
}
1100+
1101+
const workerListeners: WorkerPollSuccessListener[] = [];
1102+
const createWorker = vi.fn(() => {
1103+
let stopWorker: (() => void) | undefined;
1104+
const workerDone = new Promise<void>((resolve) => {
1105+
stopWorker = resolve;
1106+
});
1107+
return {
1108+
onMessage: vi.fn((listener: WorkerPollSuccessListener) => {
1109+
workerListeners.push(listener);
1110+
return () => undefined;
1111+
}),
1112+
stop: vi.fn(async () => {
1113+
stopWorker?.();
1114+
}),
1115+
task: vi.fn(async () => {
1116+
await workerDone;
1117+
}),
1118+
};
1119+
});
1120+
1121+
try {
1122+
const session = createPollingSession({
1123+
abortSignal: abort.signal,
1124+
log,
1125+
setStatus,
1126+
isolatedIngress: {
1127+
enabled: true,
1128+
spoolDir: tempDir,
1129+
createWorker,
1130+
drainIntervalMs: 100,
1131+
},
1132+
});
1133+
1134+
const runPromise = session.runUntilAbort();
1135+
await vi.waitFor(() => expect(handleUpdate).toHaveBeenCalledTimes(1));
1136+
workerListeners[0]?.({
1137+
type: "poll-success",
1138+
offset: null,
1139+
count: 0,
1140+
finishedAt: Date.now(),
1141+
});
1142+
expect(statusPatches(setStatus).some((patch) => patch.connected === true)).toBe(true);
1143+
1144+
await vi.advanceTimersByTimeAsync(25 * 60_000 + 100);
1145+
1146+
await vi.waitFor(() =>
1147+
expect(log).toHaveBeenCalledWith(
1148+
expect.stringContaining("isolated polling spool backlog stalled"),
1149+
),
1150+
);
1151+
expect(
1152+
statusPatches(setStatus).some(
1153+
(patch) =>
1154+
patch.connected === false &&
1155+
String(patch.lastError).includes("isolated polling spool backlog stalled"),
1156+
),
1157+
).toBe(true);
1158+
workerListeners[0]?.({
1159+
type: "poll-success",
1160+
offset: null,
1161+
count: 0,
1162+
finishedAt: Date.now(),
1163+
});
1164+
expect(statusPatches(setStatus).at(-1)?.connected).toBe(false);
1165+
1166+
releaseRegularTurn?.();
1167+
await vi.advanceTimersByTimeAsync(1_000);
1168+
await vi.waitFor(async () =>
1169+
expect(
1170+
(await listTelegramSpooledUpdates({ spoolDir: tempDir })).map(
1171+
(update) => update.updateId,
1172+
),
1173+
).toEqual([]),
1174+
);
1175+
workerListeners[0]?.({
1176+
type: "poll-success",
1177+
offset: null,
1178+
count: 0,
1179+
finishedAt: Date.now(),
1180+
});
1181+
await vi.waitFor(() => expect(statusPatches(setStatus).at(-1)?.connected).toBe(true));
1182+
expect(createWorker).toHaveBeenCalledTimes(1);
1183+
1184+
abort.abort();
1185+
await vi.advanceTimersByTimeAsync(20_000);
1186+
await runPromise;
1187+
} finally {
1188+
releaseRegularTurn?.();
1189+
vi.useRealTimers();
1190+
await fs.rm(tempDir, { recursive: true, force: true });
1191+
}
1192+
});
1193+
8851194
it("forces a restart when polling stalls without getUpdates activity", async () => {
8861195
const abort = new AbortController();
8871196
const botStop = vi.fn(async () => undefined);

0 commit comments

Comments
 (0)