fix(telegram): recover stuck live-owned spool claims

Co-Authored-By: Claude <noreply@anthropic.com>
This commit is contained in:
mikasa0818
2026-06-20 17:54:07 +08:00
committed by Ayaan Zaidi
parent 6be98022da
commit 3be0fe722a
2 changed files with 119 additions and 0 deletions
@@ -509,6 +509,22 @@ async function failedUpdateIds(spoolDir: string): Promise<number[]> {
return rows.map((row) => Number(row.event_id));
}
async function failedUpdateReasons(
spoolDir: string,
): Promise<Array<{ id: number; reason: string }>> {
const { database, kysely } = openTelegramSpoolTestKysely(spoolDir);
const rows = executeSqliteQuerySync(
database.db,
kysely
.selectFrom("channel_ingress_events")
.select(["event_id", "failed_reason"])
.where("queue_name", "=", telegramTestQueueName(spoolDir))
.where("status", "=", "failed")
.orderBy("event_id", "asc"),
).rows;
return rows.map((row) => ({ id: Number(row.event_id), reason: String(row.failed_reason) }));
}
async function adoptClaimOwner(params: {
spoolDir: string;
updateId: number;
@@ -1960,6 +1976,58 @@ describe("TelegramPollingSession", () => {
});
});
it("fails timed-out live-owned claims before draining later same-lane updates", async () => {
await withTempSpool(async (tempDir) => {
const abort = new AbortController();
const log = vi.fn();
const events: string[] = [];
await writeSpooledTestUpdates(tempDir, [
topicUpdate(42, 10, "wedged topic 10 turn"),
topicUpdate(43, 10, "later topic 10 turn"),
]);
const interrupted = (await listTelegramSpooledUpdates({ spoolDir: tempDir })).find(
(update) => update.updateId === 42,
);
if (!interrupted) {
throw new Error("Expected interrupted update");
}
const claimed = await claimTelegramSpooledUpdate(interrupted);
if (!claimed) {
throw new Error("Expected claimed update");
}
await adoptClaimOwner({
spoolDir: tempDir,
updateId: 42,
ownerId: `${process.pid}:other-process`,
claimedAt: Date.now() - 101,
});
const { runPromise, stopWorker } = startIsolatedIngressSession({
abort,
spoolDir: tempDir,
log,
spooledUpdateHandlerTimeoutMs: 100,
handleUpdate: async (update) => {
events.push(`handled:${update.update_id}`);
abort.abort();
},
});
await vi.waitFor(() => expect(events).toEqual(["handled:43"]));
await runPromise;
expect(await failedUpdateReasons(tempDir)).toEqual([
{ id: 42, reason: "lane-released-on-stuck" },
]);
expect(await pendingUpdateIds(tempDir, "all")).toEqual([]);
expect(await listTelegramSpooledUpdateClaims({ spoolDir: tempDir })).toEqual([]);
expectLogIncludes(
log,
"spooled update 42 Telegram spooled update claim held by a live worker",
);
stopWorker();
});
});
it("scans past active-lane backlogs to start unrelated lanes", async () => {
await withTempSpool(async (tempDir) => {
const abort = new AbortController();
@@ -658,6 +658,53 @@ export class TelegramPollingSession {
return deferredSpooledUpdateClaimsByKey.has(buildDeferredSpooledUpdateClaimKey(update));
}
#isTimedOutSpooledUpdateClaim(update: ClaimedTelegramSpooledUpdate): boolean {
const claimedAt = update.claim?.claimedAt;
return claimedAt !== undefined && Date.now() - claimedAt >= this.#spooledUpdateHandlerTimeoutMs;
}
async #failTimedOutLiveOwnedSpooledUpdateClaims(params: {
activeLaneKeys: Set<string>;
spoolDir: string;
}): Promise<void> {
const claims = await listTelegramSpooledUpdateClaims({ spoolDir: params.spoolDir });
for (const claim of claims) {
if (this.#isDeferredSpooledUpdateClaim(claim)) {
continue;
}
if (params.activeLaneKeys.has(this.#spooledUpdateLaneKey(claim))) {
continue;
}
if (!this.#isTimedOutSpooledUpdateClaim(claim)) {
continue;
}
if (!isTelegramSpooledUpdateClaimOwnedByOtherLiveProcess(claim)) {
continue;
}
const claimedForMs = Date.now() - (claim.claim?.claimedAt ?? Date.now());
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.`;
try {
const failed = await failTelegramSpooledUpdateClaim({
update: claim,
reason: "lane-released-on-stuck",
message,
});
if (!failed) {
this.opts.log(
`[telegram][diag] spooled update ${claim.updateId} live-owned claim no longer had a processing marker to fail.`,
);
continue;
}
} catch (err) {
this.opts.log(
`[telegram][diag] spooled update ${claim.updateId} live-owned claim could not be marked failed: ${formatErrorMessage(err)}`,
);
continue;
}
this.opts.log(`[telegram][diag] spooled update ${claim.updateId} ${message}`);
}
}
async #failTimedOutDeferredSpooledUpdate(state: DeferredSpooledUpdateClaimState): Promise<void> {
const message =
state.timedOutMessage ??
@@ -781,6 +828,10 @@ export class TelegramPollingSession {
spoolDir: string;
}): Promise<SpooledUpdateDrainResult> {
const activeLaneKeys = this.#activeSpooledUpdateLaneKeysForSpool(params.spoolDir);
await this.#failTimedOutLiveOwnedSpooledUpdateClaims({
activeLaneKeys,
spoolDir: params.spoolDir,
});
await recoverStaleTelegramSpooledUpdateClaims({
spoolDir: params.spoolDir,
staleMs: 0,