@@ -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 ;
@@ -476,6 +477,49 @@ async function pendingUpdateIds(spoolDir: string, limit: number | "all" = 100):
476477 return ( await listTelegramSpooledUpdates ( { spoolDir, limit } ) ) . map ( ( update ) => update . updateId ) ;
477478}
478479
480+ async function claimedAtForUpdate ( spoolDir : string , updateId : number ) : Promise < number > {
481+ const claim = ( await listTelegramSpooledUpdateClaims ( { spoolDir } ) ) . find (
482+ ( entry ) => entry . updateId === updateId ,
483+ ) ;
484+ if ( ! claim ?. claim ) {
485+ throw new Error ( `Expected claimed spooled update ${ updateId } ` ) ;
486+ }
487+ return claim . claim . claimedAt ;
488+ }
489+
490+ function installSpooledClaimRefreshHarness ( ) : {
491+ restore : ( ) => void ;
492+ triggerRefresh : ( ) => void ;
493+ } {
494+ let refresh : ( ( ) => void ) | undefined ;
495+ const realSetInterval = globalThis . setInterval . bind ( globalThis ) ;
496+ const setIntervalSpy = vi . spyOn ( globalThis , "setInterval" ) . mockImplementation ( ( (
497+ handler : Parameters < typeof setInterval > [ 0 ] ,
498+ timeout ?: number ,
499+ ) => {
500+ if ( timeout === pollingSessionTesting . spooledClaimRefreshIntervalMs ) {
501+ refresh = ( ) => {
502+ if ( typeof handler === "function" ) {
503+ handler ( ) ;
504+ }
505+ } ;
506+ const timer = realSetInterval ( ( ) => undefined , 2_147_483_647 ) ;
507+ timer . unref ?.( ) ;
508+ return timer ;
509+ }
510+ return realSetInterval ( handler , timeout ) ;
511+ } ) as typeof setInterval ) ;
512+ return {
513+ restore : ( ) => setIntervalSpy . mockRestore ( ) ,
514+ triggerRefresh : ( ) => {
515+ if ( ! refresh ) {
516+ throw new Error ( "Expected spooled claim refresh interval to be registered" ) ;
517+ }
518+ refresh ( ) ;
519+ } ,
520+ } ;
521+ }
522+
479523function normalizeTelegramTestAccountId ( spoolDir : string ) : string {
480524 const trimmed = path . basename ( spoolDir ) . trim ( ) ;
481525 return trimmed ? trimmed . replace ( / [ ^ a - z 0 - 9 . _ - ] + / gi, "_" ) : "default" ;
@@ -632,6 +676,7 @@ describe("TelegramPollingSession", () => {
632676 listTelegramSpooledUpdateClaims,
633677 listTelegramSpooledUpdates,
634678 recoverStaleTelegramSpooledUpdateClaims,
679+ TELEGRAM_SPOOLED_UPDATE_CLAIM_LEASE_MS : telegramSpooledUpdateClaimLeaseMs ,
635680 writeTelegramSpooledUpdate,
636681 } = await import ( "./telegram-ingress-spool.js" ) ) ;
637682 ( { createTelegramSpooledReplayDeferredParticipant } =
@@ -1014,13 +1059,13 @@ describe("TelegramPollingSession", () => {
10141059 const queue = createChannelIngressQueue ( { ...options , channelId : "telegram" } ) ;
10151060 return {
10161061 ...queue ,
1017- claim : async ( ...args : Parameters < typeof queue . claim > ) => {
1018- if ( args [ 0 ] === "0000000000000001" && ! blockedFirstClaim ) {
1062+ claimNext : async ( ...args : Parameters < typeof queue . claimNext > ) => {
1063+ if ( ! blockedFirstClaim ) {
10191064 blockedFirstClaim = true ;
10201065 firstClaimStarted ?.( ) ;
10211066 await firstClaimGate ;
10221067 }
1023- return queue . claim ( ...args ) ;
1068+ return queue . claimNext ( ...args ) ;
10241069 } ,
10251070 } ;
10261071 } ,
@@ -1565,6 +1610,132 @@ describe("TelegramPollingSession", () => {
15651610 } ) ;
15661611 } ) ;
15671612
1613+ it ( "refreshes active spooled claims while the handler is still running" , async ( ) => {
1614+ const refreshHarness = installSpooledClaimRefreshHarness ( ) ;
1615+ await withTempSpool ( async ( tempDir ) => {
1616+ const abort = new AbortController ( ) ;
1617+ const events : string [ ] = [ ] ;
1618+ let releaseHandler : ( ( ) => void ) | undefined ;
1619+ const handlerDone = new Promise < void > ( ( resolve ) => {
1620+ releaseHandler = resolve ;
1621+ } ) ;
1622+ await writeSpooledTestUpdates ( tempDir , [ topicUpdate ( 42 , 10 , "long topic 10 turn" ) ] ) ;
1623+
1624+ const { runPromise, stopWorker } = startIsolatedIngressSession ( {
1625+ abort,
1626+ spoolDir : tempDir ,
1627+ handleUpdate : async ( update ) => {
1628+ events . push ( `topic10:${ update . update_id } ` ) ;
1629+ await handlerDone ;
1630+ } ,
1631+ } ) ;
1632+
1633+ try {
1634+ await vi . waitFor ( ( ) => expect ( events ) . toEqual ( [ "topic10:42" ] ) ) ;
1635+ const before = await claimedAtForUpdate ( tempDir , 42 ) ;
1636+
1637+ await new Promise ( ( resolve ) => {
1638+ setTimeout ( resolve , 2 ) ;
1639+ } ) ;
1640+ refreshHarness . triggerRefresh ( ) ;
1641+ await vi . waitFor ( async ( ) =>
1642+ expect ( await claimedAtForUpdate ( tempDir , 42 ) ) . toBeGreaterThan ( before ) ,
1643+ ) ;
1644+
1645+ releaseHandler ?.( ) ;
1646+ await vi . waitFor ( async ( ) =>
1647+ expect ( await listTelegramSpooledUpdateClaims ( { spoolDir : tempDir } ) ) . toEqual ( [ ] ) ,
1648+ ) ;
1649+ } finally {
1650+ releaseHandler ?.( ) ;
1651+ abort . abort ( ) ;
1652+ stopWorker ( ) ;
1653+ refreshHarness . restore ( ) ;
1654+ await runPromise ;
1655+ }
1656+ } ) ;
1657+ } ) ;
1658+
1659+ it ( "stops refreshing a claim when the drain loop is stalled" , async ( ) => {
1660+ vi . useFakeTimers ( { now : 1_000 } ) ;
1661+ const refreshHarness = installSpooledClaimRefreshHarness ( ) ;
1662+ await withTempSpool ( async ( tempDir ) => {
1663+ let blockedSecondClaim = false ;
1664+ let releaseSecondClaim : ( ( ) => void ) | undefined ;
1665+ const secondClaimStarted = new Promise < void > ( ( resolve ) => {
1666+ const gate = new Promise < void > ( ( release ) => {
1667+ releaseSecondClaim = release ;
1668+ } ) ;
1669+ setTelegramRuntime ( {
1670+ state : {
1671+ resolveStateDir : ( ) => tempDir ,
1672+ openChannelIngressQueue : (
1673+ options ?: Omit < Parameters < typeof createChannelIngressQueue > [ 0 ] , "channelId" > ,
1674+ ) => {
1675+ const queue = createChannelIngressQueue ( { ...options , channelId : "telegram" } ) ;
1676+ return {
1677+ ...queue ,
1678+ claimNext : async ( ...args : Parameters < typeof queue . claimNext > ) => {
1679+ const claimOptions = args [ 0 ] ;
1680+ const blockedLaneKeys = claimOptions ?. blockedLaneKeys
1681+ ? Array . from ( claimOptions . blockedLaneKeys )
1682+ : [ ] ;
1683+ const candidateIds = claimOptions ?. candidateIds
1684+ ? Array . from ( claimOptions . candidateIds )
1685+ : [ ] ;
1686+ if (
1687+ candidateIds . includes ( "0000000000000043" ) &&
1688+ blockedLaneKeys . length > 0 &&
1689+ ! blockedSecondClaim
1690+ ) {
1691+ blockedSecondClaim = true ;
1692+ resolve ( ) ;
1693+ await gate ;
1694+ }
1695+ return queue . claimNext ( ...args ) ;
1696+ } ,
1697+ } ;
1698+ } ,
1699+ } ,
1700+ } as TelegramRuntime ) ;
1701+ } ) ;
1702+ const abort = new AbortController ( ) ;
1703+ let releaseHandler : ( ( ) => void ) | undefined ;
1704+ const handlerDone = new Promise < void > ( ( resolve ) => {
1705+ releaseHandler = resolve ;
1706+ } ) ;
1707+ await writeSpooledTestUpdates ( tempDir , [
1708+ topicUpdate ( 42 , 10 , "first topic 10 turn" ) ,
1709+ topicUpdate ( 43 , 11 , "blocked topic 11 turn" ) ,
1710+ ] ) ;
1711+
1712+ const { runPromise, stopWorker } = startIsolatedIngressSession ( {
1713+ abort,
1714+ spoolDir : tempDir ,
1715+ handleUpdate : async ( ) => {
1716+ await handlerDone ;
1717+ } ,
1718+ } ) ;
1719+
1720+ try {
1721+ await secondClaimStarted ;
1722+ const before = await claimedAtForUpdate ( tempDir , 42 ) ;
1723+ vi . setSystemTime ( 1_000 + pollingSessionTesting . spooledClaimRefreshIntervalMs * 2 + 1 ) ;
1724+ refreshHarness . triggerRefresh ( ) ;
1725+ await Promise . resolve ( ) ;
1726+ expect ( await claimedAtForUpdate ( tempDir , 42 ) ) . toBe ( before ) ;
1727+ } finally {
1728+ releaseSecondClaim ?.( ) ;
1729+ releaseHandler ?.( ) ;
1730+ abort . abort ( ) ;
1731+ stopWorker ( ) ;
1732+ refreshHarness . restore ( ) ;
1733+ vi . useRealTimers ( ) ;
1734+ await runPromise ;
1735+ }
1736+ } ) ;
1737+ } ) ;
1738+
15681739 it ( "holds buffered spooled claims until deferred processing settles without blocking same-lane buffering" , async ( ) => {
15691740 await withTempSpool ( async ( tempDir ) => {
15701741 const abort = new AbortController ( ) ;
@@ -1953,10 +2124,11 @@ describe("TelegramPollingSession", () => {
19532124 if ( ! claimed ) {
19542125 throw new Error ( "Expected claimed update" ) ;
19552126 }
2127+ const liveOwnerPid = process . ppid > 0 ? process . ppid : 1 ;
19562128 await adoptClaimOwner ( {
19572129 spoolDir : tempDir ,
19582130 updateId : 42 ,
1959- ownerId : `${ process . pid } :other-process` ,
2131+ ownerId : `${ liveOwnerPid } :other-process` ,
19602132 claimedAt : Date . now ( ) ,
19612133 } ) ;
19622134
@@ -1976,10 +2148,9 @@ describe("TelegramPollingSession", () => {
19762148 } ) ;
19772149 } ) ;
19782150
1979- it ( "fails timed-out current-process claims before draining later same-lane updates" , async ( ) => {
2151+ it ( "releases pid-reused claims before draining later same-lane updates" , async ( ) => {
19802152 await withTempSpool ( async ( tempDir ) => {
19812153 const abort = new AbortController ( ) ;
1982- const log = vi . fn ( ) ;
19832154 const events : string [ ] = [ ] ;
19842155 await writeSpooledTestUpdates ( tempDir , [
19852156 topicUpdate ( 42 , 10 , "wedged topic 10 turn" ) ,
@@ -2005,26 +2176,62 @@ describe("TelegramPollingSession", () => {
20052176 const { runPromise, stopWorker } = startIsolatedIngressSession ( {
20062177 abort,
20072178 spoolDir : tempDir ,
2008- log,
20092179 spooledUpdateHandlerTimeoutMs : 100 ,
20102180 handleUpdate : async ( update ) => {
20112181 events . push ( `handled:${ update . update_id } ` ) ;
20122182 abort . abort ( ) ;
20132183 } ,
20142184 } ) ;
20152185
2016- await vi . waitFor ( ( ) => expect ( events ) . toEqual ( [ "handled:43 " ] ) ) ;
2186+ await vi . waitFor ( ( ) => expect ( events ) . toEqual ( [ "handled:42 " ] ) ) ;
20172187 await runPromise ;
2018- expect ( await failedUpdateReasons ( tempDir ) ) . toEqual ( [
2019- { id : 42 , reason : "lane-released-on-stuck" } ,
2020- ] ) ;
2021- expect ( await pendingUpdateIds ( tempDir , "all" ) ) . toEqual ( [ ] ) ;
2188+ expect ( await failedUpdateReasons ( tempDir ) ) . toEqual ( [ ] ) ;
2189+ expect ( await pendingUpdateIds ( tempDir , "all" ) ) . toEqual ( [ 43 ] ) ;
20222190 expect ( await listTelegramSpooledUpdateClaims ( { spoolDir : tempDir } ) ) . toEqual ( [ ] ) ;
2023- expectLogIncludes (
2024- log ,
2025- "spooled update 42 Telegram spooled update claim owned by this process" ,
2191+ stopWorker ( ) ;
2192+ } ) ;
2193+ } ) ;
2194+
2195+ it ( "reclaims an expired foreign claim so the lane can drain" , async ( ) => {
2196+ await withTempSpool ( async ( tempDir ) => {
2197+ const abort = new AbortController ( ) ;
2198+ const events : number [ ] = [ ] ;
2199+ await writeSpooledTestUpdates ( tempDir , [
2200+ topicUpdate ( 42 , 10 , "expired foreign claim" ) ,
2201+ topicUpdate ( 43 , 10 , "later topic 10 turn" ) ,
2202+ ] ) ;
2203+ const interrupted = ( await listTelegramSpooledUpdates ( { spoolDir : tempDir } ) ) . find (
2204+ ( update ) => update . updateId === 42 ,
20262205 ) ;
2206+ if ( ! interrupted ) {
2207+ throw new Error ( "Expected interrupted update" ) ;
2208+ }
2209+ const claimed = await claimTelegramSpooledUpdate ( interrupted ) ;
2210+ if ( ! claimed ) {
2211+ throw new Error ( "Expected claimed update" ) ;
2212+ }
2213+ await adoptClaimOwner ( {
2214+ spoolDir : tempDir ,
2215+ updateId : 42 ,
2216+ ownerId : "1:other-process" ,
2217+ claimedAt : Date . now ( ) - telegramSpooledUpdateClaimLeaseMs - 1 ,
2218+ } ) ;
2219+
2220+ const { runPromise, stopWorker } = startIsolatedIngressSession ( {
2221+ abort,
2222+ spoolDir : tempDir ,
2223+ spooledUpdateHandlerTimeoutMs : 100 ,
2224+ handleUpdate : async ( update ) => {
2225+ events . push ( update . update_id ?? - 1 ) ;
2226+ if ( events . length === 2 ) {
2227+ abort . abort ( ) ;
2228+ }
2229+ } ,
2230+ } ) ;
2231+
2232+ await vi . waitFor ( ( ) => expect ( events ) . toEqual ( [ 42 , 43 ] ) ) ;
20272233 stopWorker ( ) ;
2234+ await runPromise ;
20282235 } ) ;
20292236 } ) ;
20302237
0 commit comments