Skip to content

Commit 0a20a63

Browse files
mikasa0818claude
andcommitted
fix(telegram): recover stuck live-owned spool claims
Co-Authored-By: Claude <[email protected]>
1 parent 7cd58cc commit 0a20a63

2 files changed

Lines changed: 119 additions & 0 deletions

File tree

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

Lines changed: 68 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -509,6 +509,22 @@ async function failedUpdateIds(spoolDir: string): Promise<number[]> {
509509
return rows.map((row) => Number(row.event_id));
510510
}
511511

512+
async function failedUpdateReasons(
513+
spoolDir: string,
514+
): Promise<Array<{ id: number; reason: string }>> {
515+
const { database, kysely } = openTelegramSpoolTestKysely(spoolDir);
516+
const rows = executeSqliteQuerySync(
517+
database.db,
518+
kysely
519+
.selectFrom("channel_ingress_events")
520+
.select(["event_id", "failed_reason"])
521+
.where("queue_name", "=", telegramTestQueueName(spoolDir))
522+
.where("status", "=", "failed")
523+
.orderBy("event_id", "asc"),
524+
).rows;
525+
return rows.map((row) => ({ id: Number(row.event_id), reason: String(row.failed_reason) }));
526+
}
527+
512528
async function adoptClaimOwner(params: {
513529
spoolDir: string;
514530
updateId: number;
@@ -1762,6 +1778,58 @@ describe("TelegramPollingSession", () => {
17621778
});
17631779
});
17641780

1781+
it("fails timed-out live-owned claims before draining later same-lane updates", async () => {
1782+
await withTempSpool(async (tempDir) => {
1783+
const abort = new AbortController();
1784+
const log = vi.fn();
1785+
const events: string[] = [];
1786+
await writeSpooledTestUpdates(tempDir, [
1787+
topicUpdate(42, 10, "wedged topic 10 turn"),
1788+
topicUpdate(43, 10, "later topic 10 turn"),
1789+
]);
1790+
const interrupted = (await listTelegramSpooledUpdates({ spoolDir: tempDir })).find(
1791+
(update) => update.updateId === 42,
1792+
);
1793+
if (!interrupted) {
1794+
throw new Error("Expected interrupted update");
1795+
}
1796+
const claimed = await claimTelegramSpooledUpdate(interrupted);
1797+
if (!claimed) {
1798+
throw new Error("Expected claimed update");
1799+
}
1800+
await adoptClaimOwner({
1801+
spoolDir: tempDir,
1802+
updateId: 42,
1803+
ownerId: `${process.pid}:other-process`,
1804+
claimedAt: Date.now() - 101,
1805+
});
1806+
1807+
const { runPromise, stopWorker } = startIsolatedIngressSession({
1808+
abort,
1809+
spoolDir: tempDir,
1810+
log,
1811+
spooledUpdateHandlerTimeoutMs: 100,
1812+
handleUpdate: async (update) => {
1813+
events.push(`handled:${update.update_id}`);
1814+
abort.abort();
1815+
},
1816+
});
1817+
1818+
await vi.waitFor(() => expect(events).toEqual(["handled:43"]));
1819+
await runPromise;
1820+
expect(await failedUpdateReasons(tempDir)).toEqual([
1821+
{ id: 42, reason: "lane-released-on-stuck" },
1822+
]);
1823+
expect(await pendingUpdateIds(tempDir, "all")).toEqual([]);
1824+
expect(await listTelegramSpooledUpdateClaims({ spoolDir: tempDir })).toEqual([]);
1825+
expectLogIncludes(
1826+
log,
1827+
"spooled update 42 Telegram spooled update claim held by a live worker",
1828+
);
1829+
stopWorker();
1830+
});
1831+
});
1832+
17651833
it("scans past active-lane backlogs to start unrelated lanes", async () => {
17661834
await withTempSpool(async (tempDir) => {
17671835
const abort = new AbortController();

extensions/telegram/src/polling-session.ts

Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -658,6 +658,53 @@ export class TelegramPollingSession {
658658
return deferredSpooledUpdateClaimsByKey.has(buildDeferredSpooledUpdateClaimKey(update));
659659
}
660660

661+
#isTimedOutSpooledUpdateClaim(update: ClaimedTelegramSpooledUpdate): boolean {
662+
const claimedAt = update.claim?.claimedAt;
663+
return claimedAt !== undefined && Date.now() - claimedAt >= this.#spooledUpdateHandlerTimeoutMs;
664+
}
665+
666+
async #failTimedOutLiveOwnedSpooledUpdateClaims(params: {
667+
activeLaneKeys: Set<string>;
668+
spoolDir: string;
669+
}): Promise<void> {
670+
const claims = await listTelegramSpooledUpdateClaims({ spoolDir: params.spoolDir });
671+
for (const claim of claims) {
672+
if (this.#isDeferredSpooledUpdateClaim(claim)) {
673+
continue;
674+
}
675+
if (params.activeLaneKeys.has(this.#spooledUpdateLaneKey(claim))) {
676+
continue;
677+
}
678+
if (!this.#isTimedOutSpooledUpdateClaim(claim)) {
679+
continue;
680+
}
681+
if (!isTelegramSpooledUpdateClaimOwnedByOtherLiveProcess(claim)) {
682+
continue;
683+
}
684+
const claimedForMs = Date.now() - (claim.claim?.claimedAt ?? Date.now());
685+
const message = `Telegram spooled update claim held by a live worker for ${formatDurationPrecise(claimedForMs)} without active handler state; marking failed so the lane can continue.`;
686+
try {
687+
const failed = await failTelegramSpooledUpdateClaim({
688+
update: claim,
689+
reason: "lane-released-on-stuck",
690+
message,
691+
});
692+
if (!failed) {
693+
this.opts.log(
694+
`[telegram][diag] spooled update ${claim.updateId} live-owned claim no longer had a processing marker to fail.`,
695+
);
696+
continue;
697+
}
698+
} catch (err) {
699+
this.opts.log(
700+
`[telegram][diag] spooled update ${claim.updateId} live-owned claim could not be marked failed: ${formatErrorMessage(err)}`,
701+
);
702+
continue;
703+
}
704+
this.opts.log(`[telegram][diag] spooled update ${claim.updateId} ${message}`);
705+
}
706+
}
707+
661708
async #failTimedOutDeferredSpooledUpdate(state: DeferredSpooledUpdateClaimState): Promise<void> {
662709
const message =
663710
state.timedOutMessage ??
@@ -781,6 +828,10 @@ export class TelegramPollingSession {
781828
spoolDir: string;
782829
}): Promise<SpooledUpdateDrainResult> {
783830
const activeLaneKeys = this.#activeSpooledUpdateLaneKeysForSpool(params.spoolDir);
831+
await this.#failTimedOutLiveOwnedSpooledUpdateClaims({
832+
activeLaneKeys,
833+
spoolDir: params.spoolDir,
834+
});
784835
await recoverStaleTelegramSpooledUpdateClaims({
785836
spoolDir: params.spoolDir,
786837
staleMs: 0,

0 commit comments

Comments
 (0)