Skip to content

Commit 769a993

Browse files
vincentkocRomneyDa
andcommitted
fix(telegram): recover stalled ingress spool claims (#97118)
* fix(telegram): drain ingress with native queue claims * fix(telegram): bound native claim drain snapshots * fix(telegram): recover pid-reused ingress claims * fix(channels): block claimed candidate lanes * fix(telegram): recover stalled ingress spool claims * test(telegram): cover native claimNext drain stalls --------- Co-authored-by: Dallin Romney <[email protected]> (cherry picked from commit b8e3de1)
1 parent c862a64 commit 769a993

6 files changed

Lines changed: 981 additions & 110 deletions

File tree

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

Lines changed: 222 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -72,6 +72,7 @@ let isTelegramSpooledUpdateClaimOwnedByOtherLiveProcess: typeof import("./telegr
7272
let listTelegramSpooledUpdateClaims: typeof import("./telegram-ingress-spool.js").listTelegramSpooledUpdateClaims;
7373
let listTelegramSpooledUpdates: typeof import("./telegram-ingress-spool.js").listTelegramSpooledUpdates;
7474
let recoverStaleTelegramSpooledUpdateClaims: typeof import("./telegram-ingress-spool.js").recoverStaleTelegramSpooledUpdateClaims;
75+
let telegramSpooledUpdateClaimLeaseMs: typeof import("./telegram-ingress-spool.js").TELEGRAM_SPOOLED_UPDATE_CLAIM_LEASE_MS;
7576
let writeTelegramSpooledUpdate: typeof import("./telegram-ingress-spool.js").writeTelegramSpooledUpdate;
7677
let createTelegramSpooledReplayDeferredParticipant: typeof import("./bot-processing-outcome.js").createTelegramSpooledReplayDeferredParticipant;
7778
let TelegramMessageDispatchReplayForgetError: typeof import("./message-dispatch-dedupe.js").TelegramMessageDispatchReplayForgetError;
@@ -476,6 +477,49 @@ async function pendingUpdateIds(spoolDir: string, limit: number | "all" = 100):
476477
return (await listTelegramSpooledUpdates({ spoolDir, limit })).map((update) => update.updateId);
477478
}
478479

480+
async function claimedAtForUpdate(spoolDir: string, updateId: number): Promise<number> {
481+
const claim = (await listTelegramSpooledUpdateClaims({ spoolDir })).find(
482+
(entry) => entry.updateId === updateId,
483+
);
484+
if (!claim?.claim) {
485+
throw new Error(`Expected claimed spooled update ${updateId}`);
486+
}
487+
return claim.claim.claimedAt;
488+
}
489+
490+
function installSpooledClaimRefreshHarness(): {
491+
restore: () => void;
492+
triggerRefresh: () => void;
493+
} {
494+
let refresh: (() => void) | undefined;
495+
const realSetInterval = globalThis.setInterval.bind(globalThis);
496+
const setIntervalSpy = vi.spyOn(globalThis, "setInterval").mockImplementation(((
497+
handler: Parameters<typeof setInterval>[0],
498+
timeout?: number,
499+
) => {
500+
if (timeout === pollingSessionTesting.spooledClaimRefreshIntervalMs) {
501+
refresh = () => {
502+
if (typeof handler === "function") {
503+
handler();
504+
}
505+
};
506+
const timer = realSetInterval(() => undefined, 2_147_483_647);
507+
timer.unref?.();
508+
return timer;
509+
}
510+
return realSetInterval(handler, timeout);
511+
}) as typeof setInterval);
512+
return {
513+
restore: () => setIntervalSpy.mockRestore(),
514+
triggerRefresh: () => {
515+
if (!refresh) {
516+
throw new Error("Expected spooled claim refresh interval to be registered");
517+
}
518+
refresh();
519+
},
520+
};
521+
}
522+
479523
function normalizeTelegramTestAccountId(spoolDir: string): string {
480524
const trimmed = path.basename(spoolDir).trim();
481525
return trimmed ? trimmed.replace(/[^a-z0-9._-]+/gi, "_") : "default";
@@ -632,6 +676,7 @@ describe("TelegramPollingSession", () => {
632676
listTelegramSpooledUpdateClaims,
633677
listTelegramSpooledUpdates,
634678
recoverStaleTelegramSpooledUpdateClaims,
679+
TELEGRAM_SPOOLED_UPDATE_CLAIM_LEASE_MS: telegramSpooledUpdateClaimLeaseMs,
635680
writeTelegramSpooledUpdate,
636681
} = await import("./telegram-ingress-spool.js"));
637682
({ createTelegramSpooledReplayDeferredParticipant } =
@@ -1014,13 +1059,13 @@ describe("TelegramPollingSession", () => {
10141059
const queue = createChannelIngressQueue({ ...options, channelId: "telegram" });
10151060
return {
10161061
...queue,
1017-
claim: async (...args: Parameters<typeof queue.claim>) => {
1018-
if (args[0] === "0000000000000001" && !blockedFirstClaim) {
1062+
claimNext: async (...args: Parameters<typeof queue.claimNext>) => {
1063+
if (!blockedFirstClaim) {
10191064
blockedFirstClaim = true;
10201065
firstClaimStarted?.();
10211066
await firstClaimGate;
10221067
}
1023-
return queue.claim(...args);
1068+
return queue.claimNext(...args);
10241069
},
10251070
};
10261071
},
@@ -1565,6 +1610,132 @@ describe("TelegramPollingSession", () => {
15651610
});
15661611
});
15671612

1613+
it("refreshes active spooled claims while the handler is still running", async () => {
1614+
const refreshHarness = installSpooledClaimRefreshHarness();
1615+
await withTempSpool(async (tempDir) => {
1616+
const abort = new AbortController();
1617+
const events: string[] = [];
1618+
let releaseHandler: (() => void) | undefined;
1619+
const handlerDone = new Promise<void>((resolve) => {
1620+
releaseHandler = resolve;
1621+
});
1622+
await writeSpooledTestUpdates(tempDir, [topicUpdate(42, 10, "long topic 10 turn")]);
1623+
1624+
const { runPromise, stopWorker } = startIsolatedIngressSession({
1625+
abort,
1626+
spoolDir: tempDir,
1627+
handleUpdate: async (update) => {
1628+
events.push(`topic10:${update.update_id}`);
1629+
await handlerDone;
1630+
},
1631+
});
1632+
1633+
try {
1634+
await vi.waitFor(() => expect(events).toEqual(["topic10:42"]));
1635+
const before = await claimedAtForUpdate(tempDir, 42);
1636+
1637+
await new Promise((resolve) => {
1638+
setTimeout(resolve, 2);
1639+
});
1640+
refreshHarness.triggerRefresh();
1641+
await vi.waitFor(async () =>
1642+
expect(await claimedAtForUpdate(tempDir, 42)).toBeGreaterThan(before),
1643+
);
1644+
1645+
releaseHandler?.();
1646+
await vi.waitFor(async () =>
1647+
expect(await listTelegramSpooledUpdateClaims({ spoolDir: tempDir })).toEqual([]),
1648+
);
1649+
} finally {
1650+
releaseHandler?.();
1651+
abort.abort();
1652+
stopWorker();
1653+
refreshHarness.restore();
1654+
await runPromise;
1655+
}
1656+
});
1657+
});
1658+
1659+
it("stops refreshing a claim when the drain loop is stalled", async () => {
1660+
vi.useFakeTimers({ now: 1_000 });
1661+
const refreshHarness = installSpooledClaimRefreshHarness();
1662+
await withTempSpool(async (tempDir) => {
1663+
let blockedSecondClaim = false;
1664+
let releaseSecondClaim: (() => void) | undefined;
1665+
const secondClaimStarted = new Promise<void>((resolve) => {
1666+
const gate = new Promise<void>((release) => {
1667+
releaseSecondClaim = release;
1668+
});
1669+
setTelegramRuntime({
1670+
state: {
1671+
resolveStateDir: () => tempDir,
1672+
openChannelIngressQueue: (
1673+
options?: Omit<Parameters<typeof createChannelIngressQueue>[0], "channelId">,
1674+
) => {
1675+
const queue = createChannelIngressQueue({ ...options, channelId: "telegram" });
1676+
return {
1677+
...queue,
1678+
claimNext: async (...args: Parameters<typeof queue.claimNext>) => {
1679+
const claimOptions = args[0];
1680+
const blockedLaneKeys = claimOptions?.blockedLaneKeys
1681+
? Array.from(claimOptions.blockedLaneKeys)
1682+
: [];
1683+
const candidateIds = claimOptions?.candidateIds
1684+
? Array.from(claimOptions.candidateIds)
1685+
: [];
1686+
if (
1687+
candidateIds.includes("0000000000000043") &&
1688+
blockedLaneKeys.length > 0 &&
1689+
!blockedSecondClaim
1690+
) {
1691+
blockedSecondClaim = true;
1692+
resolve();
1693+
await gate;
1694+
}
1695+
return queue.claimNext(...args);
1696+
},
1697+
};
1698+
},
1699+
},
1700+
} as TelegramRuntime);
1701+
});
1702+
const abort = new AbortController();
1703+
let releaseHandler: (() => void) | undefined;
1704+
const handlerDone = new Promise<void>((resolve) => {
1705+
releaseHandler = resolve;
1706+
});
1707+
await writeSpooledTestUpdates(tempDir, [
1708+
topicUpdate(42, 10, "first topic 10 turn"),
1709+
topicUpdate(43, 11, "blocked topic 11 turn"),
1710+
]);
1711+
1712+
const { runPromise, stopWorker } = startIsolatedIngressSession({
1713+
abort,
1714+
spoolDir: tempDir,
1715+
handleUpdate: async () => {
1716+
await handlerDone;
1717+
},
1718+
});
1719+
1720+
try {
1721+
await secondClaimStarted;
1722+
const before = await claimedAtForUpdate(tempDir, 42);
1723+
vi.setSystemTime(1_000 + pollingSessionTesting.spooledClaimRefreshIntervalMs * 2 + 1);
1724+
refreshHarness.triggerRefresh();
1725+
await Promise.resolve();
1726+
expect(await claimedAtForUpdate(tempDir, 42)).toBe(before);
1727+
} finally {
1728+
releaseSecondClaim?.();
1729+
releaseHandler?.();
1730+
abort.abort();
1731+
stopWorker();
1732+
refreshHarness.restore();
1733+
vi.useRealTimers();
1734+
await runPromise;
1735+
}
1736+
});
1737+
});
1738+
15681739
it("holds buffered spooled claims until deferred processing settles without blocking same-lane buffering", async () => {
15691740
await withTempSpool(async (tempDir) => {
15701741
const abort = new AbortController();
@@ -1953,10 +2124,11 @@ describe("TelegramPollingSession", () => {
19532124
if (!claimed) {
19542125
throw new Error("Expected claimed update");
19552126
}
2127+
const liveOwnerPid = process.ppid > 0 ? process.ppid : 1;
19562128
await adoptClaimOwner({
19572129
spoolDir: tempDir,
19582130
updateId: 42,
1959-
ownerId: `${process.pid}:other-process`,
2131+
ownerId: `${liveOwnerPid}:other-process`,
19602132
claimedAt: Date.now(),
19612133
});
19622134

@@ -1976,10 +2148,9 @@ describe("TelegramPollingSession", () => {
19762148
});
19772149
});
19782150

1979-
it("fails timed-out current-process claims before draining later same-lane updates", async () => {
2151+
it("releases pid-reused claims before draining later same-lane updates", async () => {
19802152
await withTempSpool(async (tempDir) => {
19812153
const abort = new AbortController();
1982-
const log = vi.fn();
19832154
const events: string[] = [];
19842155
await writeSpooledTestUpdates(tempDir, [
19852156
topicUpdate(42, 10, "wedged topic 10 turn"),
@@ -2005,26 +2176,62 @@ describe("TelegramPollingSession", () => {
20052176
const { runPromise, stopWorker } = startIsolatedIngressSession({
20062177
abort,
20072178
spoolDir: tempDir,
2008-
log,
20092179
spooledUpdateHandlerTimeoutMs: 100,
20102180
handleUpdate: async (update) => {
20112181
events.push(`handled:${update.update_id}`);
20122182
abort.abort();
20132183
},
20142184
});
20152185

2016-
await vi.waitFor(() => expect(events).toEqual(["handled:43"]));
2186+
await vi.waitFor(() => expect(events).toEqual(["handled:42"]));
20172187
await runPromise;
2018-
expect(await failedUpdateReasons(tempDir)).toEqual([
2019-
{ id: 42, reason: "lane-released-on-stuck" },
2020-
]);
2021-
expect(await pendingUpdateIds(tempDir, "all")).toEqual([]);
2188+
expect(await failedUpdateReasons(tempDir)).toEqual([]);
2189+
expect(await pendingUpdateIds(tempDir, "all")).toEqual([43]);
20222190
expect(await listTelegramSpooledUpdateClaims({ spoolDir: tempDir })).toEqual([]);
2023-
expectLogIncludes(
2024-
log,
2025-
"spooled update 42 Telegram spooled update claim owned by this process",
2191+
stopWorker();
2192+
});
2193+
});
2194+
2195+
it("reclaims an expired foreign claim so the lane can drain", async () => {
2196+
await withTempSpool(async (tempDir) => {
2197+
const abort = new AbortController();
2198+
const events: number[] = [];
2199+
await writeSpooledTestUpdates(tempDir, [
2200+
topicUpdate(42, 10, "expired foreign claim"),
2201+
topicUpdate(43, 10, "later topic 10 turn"),
2202+
]);
2203+
const interrupted = (await listTelegramSpooledUpdates({ spoolDir: tempDir })).find(
2204+
(update) => update.updateId === 42,
20262205
);
2206+
if (!interrupted) {
2207+
throw new Error("Expected interrupted update");
2208+
}
2209+
const claimed = await claimTelegramSpooledUpdate(interrupted);
2210+
if (!claimed) {
2211+
throw new Error("Expected claimed update");
2212+
}
2213+
await adoptClaimOwner({
2214+
spoolDir: tempDir,
2215+
updateId: 42,
2216+
ownerId: "1:other-process",
2217+
claimedAt: Date.now() - telegramSpooledUpdateClaimLeaseMs - 1,
2218+
});
2219+
2220+
const { runPromise, stopWorker } = startIsolatedIngressSession({
2221+
abort,
2222+
spoolDir: tempDir,
2223+
spooledUpdateHandlerTimeoutMs: 100,
2224+
handleUpdate: async (update) => {
2225+
events.push(update.update_id ?? -1);
2226+
if (events.length === 2) {
2227+
abort.abort();
2228+
}
2229+
},
2230+
});
2231+
2232+
await vi.waitFor(() => expect(events).toEqual([42, 43]));
20272233
stopWorker();
2234+
await runPromise;
20282235
});
20292236
});
20302237

0 commit comments

Comments
 (0)