diff --git a/extensions/telegram/src/polling-session.test.ts b/extensions/telegram/src/polling-session.test.ts index 3590493f74d2..1dd53c2bb439 100644 --- a/extensions/telegram/src/polling-session.test.ts +++ b/extensions/telegram/src/polling-session.test.ts @@ -509,6 +509,22 @@ async function failedUpdateIds(spoolDir: string): Promise { return rows.map((row) => Number(row.event_id)); } +async function failedUpdateReasons( + spoolDir: string, +): Promise> { + 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(); diff --git a/extensions/telegram/src/polling-session.ts b/extensions/telegram/src/polling-session.ts index 5a4f72b44395..1796ca2b6ce0 100644 --- a/extensions/telegram/src/polling-session.ts +++ b/extensions/telegram/src/polling-session.ts @@ -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; + spoolDir: string; + }): Promise { + 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 { const message = state.timedOutMessage ?? @@ -781,6 +828,10 @@ export class TelegramPollingSession { spoolDir: string; }): Promise { const activeLaneKeys = this.#activeSpooledUpdateLaneKeysForSpool(params.spoolDir); + await this.#failTimedOutLiveOwnedSpooledUpdateClaims({ + activeLaneKeys, + spoolDir: params.spoolDir, + }); await recoverStaleTelegramSpooledUpdateClaims({ spoolDir: params.spoolDir, staleMs: 0,