
























@@ -509,6 +509,22 @@ async function failedUpdateIds(spoolDir: string): Promise<number[]> {
509509return rows.map((row) => Number(row.event_id));
510510}
511511512+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+512528async function adoptClaimOwner(params: {
513529spoolDir: string;
514530updateId: number;
@@ -1960,6 +1976,58 @@ describe("TelegramPollingSession", () => {
19601976});
19611977});
196219781979+it("fails timed-out live-owned claims before draining later same-lane updates", async () => {
1980+await withTempSpool(async (tempDir) => {
1981+const abort = new AbortController();
1982+const log = vi.fn();
1983+const events: string[] = [];
1984+await writeSpooledTestUpdates(tempDir, [
1985+topicUpdate(42, 10, "wedged topic 10 turn"),
1986+topicUpdate(43, 10, "later topic 10 turn"),
1987+]);
1988+const interrupted = (await listTelegramSpooledUpdates({ spoolDir: tempDir })).find(
1989+(update) => update.updateId === 42,
1990+);
1991+if (!interrupted) {
1992+throw new Error("Expected interrupted update");
1993+}
1994+const claimed = await claimTelegramSpooledUpdate(interrupted);
1995+if (!claimed) {
1996+throw new Error("Expected claimed update");
1997+}
1998+await adoptClaimOwner({
1999+spoolDir: tempDir,
2000+updateId: 42,
2001+ownerId: `${process.pid}:other-process`,
2002+claimedAt: Date.now() - 101,
2003+});
2004+2005+const { runPromise, stopWorker } = startIsolatedIngressSession({
2006+ abort,
2007+spoolDir: tempDir,
2008+ log,
2009+spooledUpdateHandlerTimeoutMs: 100,
2010+handleUpdate: async (update) => {
2011+events.push(`handled:${update.update_id}`);
2012+abort.abort();
2013+},
2014+});
2015+2016+await vi.waitFor(() => expect(events).toEqual(["handled:43"]));
2017+await runPromise;
2018+expect(await failedUpdateReasons(tempDir)).toEqual([
2019+{ id: 42, reason: "lane-released-on-stuck" },
2020+]);
2021+expect(await pendingUpdateIds(tempDir, "all")).toEqual([]);
2022+expect(await listTelegramSpooledUpdateClaims({ spoolDir: tempDir })).toEqual([]);
2023+expectLogIncludes(
2024+log,
2025+"spooled update 42 Telegram spooled update claim held by a live worker",
2026+);
2027+stopWorker();
2028+});
2029+});
2030+19632031it("scans past active-lane backlogs to start unrelated lanes", async () => {
19642032await withTempSpool(async (tempDir) => {
19652033const abort = new AbortController();
此内容由惯性聚合(RSS阅读器)自动聚合整理,仅供阅读参考。 原文来自 — 版权归原作者所有。