@@ -72,6 +72,7 @@ let isTelegramSpooledUpdateClaimOwnedByOtherLiveProcess: typeof import("./telegr
7272let listTelegramSpooledUpdateClaims : typeof import ( "./telegram-ingress-spool.js" ) . listTelegramSpooledUpdateClaims ;
7373let listTelegramSpooledUpdates : typeof import ( "./telegram-ingress-spool.js" ) . listTelegramSpooledUpdates ;
7474let recoverStaleTelegramSpooledUpdateClaims : typeof import ( "./telegram-ingress-spool.js" ) . recoverStaleTelegramSpooledUpdateClaims ;
75+ let telegramSpooledUpdateClaimLeaseMs : typeof import ( "./telegram-ingress-spool.js" ) . TELEGRAM_SPOOLED_UPDATE_CLAIM_LEASE_MS ;
7576let writeTelegramSpooledUpdate : typeof import ( "./telegram-ingress-spool.js" ) . writeTelegramSpooledUpdate ;
7677let createTelegramSpooledReplayDeferredParticipant : typeof import ( "./bot-processing-outcome.js" ) . createTelegramSpooledReplayDeferredParticipant ;
7778let TelegramMessageDispatchReplayForgetError : typeof import ( "./message-dispatch-dedupe.js" ) . TelegramMessageDispatchReplayForgetError ;
@@ -685,6 +686,7 @@ describe("TelegramPollingSession", () => {
685686 listTelegramSpooledUpdateClaims,
686687 listTelegramSpooledUpdates,
687688 recoverStaleTelegramSpooledUpdateClaims,
689+ TELEGRAM_SPOOLED_UPDATE_CLAIM_LEASE_MS : telegramSpooledUpdateClaimLeaseMs ,
688690 writeTelegramSpooledUpdate,
689691 } = await import ( "./telegram-ingress-spool.js" ) ) ;
690692 ( { createTelegramSpooledReplayDeferredParticipant } =
@@ -1067,13 +1069,13 @@ describe("TelegramPollingSession", () => {
10671069 const queue = createChannelIngressQueue ( { ...options , channelId : "telegram" } ) ;
10681070 return {
10691071 ...queue ,
1070- claim : async ( ...args : Parameters < typeof queue . claim > ) => {
1071- if ( args [ 0 ] === "0000000000000001" && ! blockedFirstClaim ) {
1072+ claimNext : async ( ...args : Parameters < typeof queue . claimNext > ) => {
1073+ if ( ! blockedFirstClaim ) {
10721074 blockedFirstClaim = true ;
10731075 firstClaimStarted ?.( ) ;
10741076 await firstClaimGate ;
10751077 }
1076- return queue . claim ( ...args ) ;
1078+ return queue . claimNext ( ...args ) ;
10771079 } ,
10781080 } ;
10791081 } ,
@@ -1661,6 +1663,86 @@ describe("TelegramPollingSession", () => {
16611663 } ) ;
16621664 } ) ;
16631665
1666+ it ( "stops refreshing a claim when the drain loop is stalled" , async ( ) => {
1667+ vi . useFakeTimers ( { now : 1_000 } ) ;
1668+ const refreshHarness = installSpooledClaimRefreshHarness ( ) ;
1669+ await withTempSpool ( async ( tempDir ) => {
1670+ let blockedSecondClaim = false ;
1671+ let releaseSecondClaim : ( ( ) => void ) | undefined ;
1672+ const secondClaimStarted = new Promise < void > ( ( resolve ) => {
1673+ const gate = new Promise < void > ( ( release ) => {
1674+ releaseSecondClaim = release ;
1675+ } ) ;
1676+ setTelegramRuntime ( {
1677+ state : {
1678+ resolveStateDir : ( ) => tempDir ,
1679+ openChannelIngressQueue : (
1680+ options ?: Omit < Parameters < typeof createChannelIngressQueue > [ 0 ] , "channelId" > ,
1681+ ) => {
1682+ const queue = createChannelIngressQueue ( { ...options , channelId : "telegram" } ) ;
1683+ return {
1684+ ...queue ,
1685+ claimNext : async ( ...args : Parameters < typeof queue . claimNext > ) => {
1686+ const claimOptions = args [ 0 ] ;
1687+ const blockedLaneKeys = claimOptions ?. blockedLaneKeys
1688+ ? Array . from ( claimOptions . blockedLaneKeys )
1689+ : [ ] ;
1690+ const candidateIds = claimOptions ?. candidateIds
1691+ ? Array . from ( claimOptions . candidateIds )
1692+ : [ ] ;
1693+ if (
1694+ candidateIds . includes ( "0000000000000043" ) &&
1695+ blockedLaneKeys . length > 0 &&
1696+ ! blockedSecondClaim
1697+ ) {
1698+ blockedSecondClaim = true ;
1699+ resolve ( ) ;
1700+ await gate ;
1701+ }
1702+ return queue . claimNext ( ...args ) ;
1703+ } ,
1704+ } ;
1705+ } ,
1706+ } ,
1707+ } as TelegramRuntime ) ;
1708+ } ) ;
1709+ const abort = new AbortController ( ) ;
1710+ let releaseHandler : ( ( ) => void ) | undefined ;
1711+ const handlerDone = new Promise < void > ( ( resolve ) => {
1712+ releaseHandler = resolve ;
1713+ } ) ;
1714+ await writeSpooledTestUpdates ( tempDir , [
1715+ topicUpdate ( 42 , 10 , "first topic 10 turn" ) ,
1716+ topicUpdate ( 43 , 11 , "blocked topic 11 turn" ) ,
1717+ ] ) ;
1718+
1719+ const { runPromise, stopWorker } = startIsolatedIngressSession ( {
1720+ abort,
1721+ spoolDir : tempDir ,
1722+ handleUpdate : async ( ) => {
1723+ await handlerDone ;
1724+ } ,
1725+ } ) ;
1726+
1727+ try {
1728+ await secondClaimStarted ;
1729+ const before = await claimedAtForUpdate ( tempDir , 42 ) ;
1730+ vi . setSystemTime ( 1_000 + pollingSessionTesting . spooledClaimRefreshIntervalMs * 2 + 1 ) ;
1731+ refreshHarness . triggerRefresh ( ) ;
1732+ await Promise . resolve ( ) ;
1733+ expect ( await claimedAtForUpdate ( tempDir , 42 ) ) . toBe ( before ) ;
1734+ } finally {
1735+ releaseSecondClaim ?.( ) ;
1736+ releaseHandler ?.( ) ;
1737+ abort . abort ( ) ;
1738+ stopWorker ( ) ;
1739+ refreshHarness . restore ( ) ;
1740+ vi . useRealTimers ( ) ;
1741+ await runPromise ;
1742+ }
1743+ } ) ;
1744+ } ) ;
1745+
16641746 it ( "holds buffered spooled claims until deferred processing settles without blocking same-lane buffering" , async ( ) => {
16651747 await withTempSpool ( async ( tempDir ) => {
16661748 const abort = new AbortController ( ) ;
@@ -2180,10 +2262,11 @@ describe("TelegramPollingSession", () => {
21802262 if ( ! claimed ) {
21812263 throw new Error ( "Expected claimed update" ) ;
21822264 }
2265+ const liveOwnerPid = process . ppid > 0 ? process . ppid : 1 ;
21832266 await adoptClaimOwner ( {
21842267 spoolDir : tempDir ,
21852268 updateId : 42 ,
2186- ownerId : `${ process . pid } :other-process` ,
2269+ ownerId : `${ liveOwnerPid } :other-process` ,
21872270 claimedAt : Date . now ( ) ,
21882271 } ) ;
21892272
@@ -2203,10 +2286,9 @@ describe("TelegramPollingSession", () => {
22032286 } ) ;
22042287 } ) ;
22052288
2206- it ( "fails timed-out current-process claims before draining later same-lane updates" , async ( ) => {
2289+ it ( "releases pid-reused claims before draining later same-lane updates" , async ( ) => {
22072290 await withTempSpool ( async ( tempDir ) => {
22082291 const abort = new AbortController ( ) ;
2209- const log = vi . fn ( ) ;
22102292 const events : string [ ] = [ ] ;
22112293 await writeSpooledTestUpdates ( tempDir , [
22122294 topicUpdate ( 42 , 10 , "wedged topic 10 turn" ) ,
@@ -2232,26 +2314,62 @@ describe("TelegramPollingSession", () => {
22322314 const { runPromise, stopWorker } = startIsolatedIngressSession ( {
22332315 abort,
22342316 spoolDir : tempDir ,
2235- log,
22362317 spooledUpdateHandlerTimeoutMs : 100 ,
22372318 handleUpdate : async ( update ) => {
22382319 events . push ( `handled:${ update . update_id } ` ) ;
22392320 abort . abort ( ) ;
22402321 } ,
22412322 } ) ;
22422323
2243- await vi . waitFor ( ( ) => expect ( events ) . toEqual ( [ "handled:43 " ] ) ) ;
2324+ await vi . waitFor ( ( ) => expect ( events ) . toEqual ( [ "handled:42 " ] ) ) ;
22442325 await runPromise ;
2245- expect ( await failedUpdateReasons ( tempDir ) ) . toEqual ( [
2246- { id : 42 , reason : "lane-released-on-stuck" } ,
2247- ] ) ;
2248- expect ( await pendingUpdateIds ( tempDir , "all" ) ) . toEqual ( [ ] ) ;
2326+ expect ( await failedUpdateReasons ( tempDir ) ) . toEqual ( [ ] ) ;
2327+ expect ( await pendingUpdateIds ( tempDir , "all" ) ) . toEqual ( [ 43 ] ) ;
22492328 expect ( await listTelegramSpooledUpdateClaims ( { spoolDir : tempDir } ) ) . toEqual ( [ ] ) ;
2250- expectLogIncludes (
2251- log ,
2252- "spooled update 42 Telegram spooled update claim owned by this process" ,
2329+ stopWorker ( ) ;
2330+ } ) ;
2331+ } ) ;
2332+
2333+ it ( "reclaims an expired foreign claim so the lane can drain" , async ( ) => {
2334+ await withTempSpool ( async ( tempDir ) => {
2335+ const abort = new AbortController ( ) ;
2336+ const events : number [ ] = [ ] ;
2337+ await writeSpooledTestUpdates ( tempDir , [
2338+ topicUpdate ( 42 , 10 , "expired foreign claim" ) ,
2339+ topicUpdate ( 43 , 10 , "later topic 10 turn" ) ,
2340+ ] ) ;
2341+ const interrupted = ( await listTelegramSpooledUpdates ( { spoolDir : tempDir } ) ) . find (
2342+ ( update ) => update . updateId === 42 ,
22532343 ) ;
2344+ if ( ! interrupted ) {
2345+ throw new Error ( "Expected interrupted update" ) ;
2346+ }
2347+ const claimed = await claimTelegramSpooledUpdate ( interrupted ) ;
2348+ if ( ! claimed ) {
2349+ throw new Error ( "Expected claimed update" ) ;
2350+ }
2351+ await adoptClaimOwner ( {
2352+ spoolDir : tempDir ,
2353+ updateId : 42 ,
2354+ ownerId : "1:other-process" ,
2355+ claimedAt : Date . now ( ) - telegramSpooledUpdateClaimLeaseMs - 1 ,
2356+ } ) ;
2357+
2358+ const { runPromise, stopWorker } = startIsolatedIngressSession ( {
2359+ abort,
2360+ spoolDir : tempDir ,
2361+ spooledUpdateHandlerTimeoutMs : 100 ,
2362+ handleUpdate : async ( update ) => {
2363+ events . push ( update . update_id ?? - 1 ) ;
2364+ if ( events . length === 2 ) {
2365+ abort . abort ( ) ;
2366+ }
2367+ } ,
2368+ } ) ;
2369+
2370+ await vi . waitFor ( ( ) => expect ( events ) . toEqual ( [ 42 , 43 ] ) ) ;
22542371 stopWorker ( ) ;
2372+ await runPromise ;
22552373 } ) ;
22562374 } ) ;
22572375
0 commit comments