@@ -486,6 +486,49 @@ async function pendingUpdateIds(spoolDir: string, limit: number | "all" = 100):
486486 return ( await listTelegramSpooledUpdates ( { spoolDir, limit } ) ) . map ( ( update ) => update . updateId ) ;
487487}
488488
489+ async function claimedAtForUpdate ( spoolDir : string , updateId : number ) : Promise < number > {
490+ const claim = ( await listTelegramSpooledUpdateClaims ( { spoolDir } ) ) . find (
491+ ( entry ) => entry . updateId === updateId ,
492+ ) ;
493+ if ( ! claim ?. claim ) {
494+ throw new Error ( `Expected claimed spooled update ${ updateId } ` ) ;
495+ }
496+ return claim . claim . claimedAt ;
497+ }
498+
499+ function installSpooledClaimRefreshHarness ( ) : {
500+ restore : ( ) => void ;
501+ triggerRefresh : ( ) => void ;
502+ } {
503+ let refresh : ( ( ) => void ) | undefined ;
504+ const realSetInterval = globalThis . setInterval . bind ( globalThis ) ;
505+ const setIntervalSpy = vi . spyOn ( globalThis , "setInterval" ) . mockImplementation ( ( (
506+ handler : Parameters < typeof setInterval > [ 0 ] ,
507+ timeout ?: number ,
508+ ) => {
509+ if ( timeout === pollingSessionTesting . spooledClaimRefreshIntervalMs ) {
510+ refresh = ( ) => {
511+ if ( typeof handler === "function" ) {
512+ handler ( ) ;
513+ }
514+ } ;
515+ const timer = realSetInterval ( ( ) => undefined , 2_147_483_647 ) ;
516+ timer . unref ?.( ) ;
517+ return timer ;
518+ }
519+ return realSetInterval ( handler , timeout ) ;
520+ } ) as typeof setInterval ) ;
521+ return {
522+ restore : ( ) => setIntervalSpy . mockRestore ( ) ,
523+ triggerRefresh : ( ) => {
524+ if ( ! refresh ) {
525+ throw new Error ( "Expected spooled claim refresh interval to be registered" ) ;
526+ }
527+ refresh ( ) ;
528+ } ,
529+ } ;
530+ }
531+
489532function normalizeTelegramTestAccountId ( spoolDir : string ) : string {
490533 const trimmed = path . basename ( spoolDir ) . trim ( ) ;
491534 return trimmed ? trimmed . replace ( / [ ^ a - z 0 - 9 . _ - ] + / gi, "_" ) : "default" ;
@@ -1575,6 +1618,49 @@ describe("TelegramPollingSession", () => {
15751618 } ) ;
15761619 } ) ;
15771620
1621+ it ( "refreshes active spooled claims while the handler is still running" , async ( ) => {
1622+ const refreshHarness = installSpooledClaimRefreshHarness ( ) ;
1623+ await withTempSpool ( async ( tempDir ) => {
1624+ const abort = new AbortController ( ) ;
1625+ const events : string [ ] = [ ] ;
1626+ let releaseHandler : ( ( ) => void ) | undefined ;
1627+ const handlerDone = new Promise < void > ( ( resolve ) => {
1628+ releaseHandler = resolve ;
1629+ } ) ;
1630+ await writeSpooledTestUpdates ( tempDir , [ topicUpdate ( 42 , 10 , "long topic 10 turn" ) ] ) ;
1631+
1632+ const { runPromise, stopWorker } = startIsolatedIngressSession ( {
1633+ abort,
1634+ spoolDir : tempDir ,
1635+ handleUpdate : async ( update ) => {
1636+ events . push ( `topic10:${ update . update_id } ` ) ;
1637+ await handlerDone ;
1638+ } ,
1639+ } ) ;
1640+
1641+ try {
1642+ await vi . waitFor ( ( ) => expect ( events ) . toEqual ( [ "topic10:42" ] ) ) ;
1643+ const before = await claimedAtForUpdate ( tempDir , 42 ) ;
1644+
1645+ refreshHarness . triggerRefresh ( ) ;
1646+ await vi . waitFor ( async ( ) =>
1647+ expect ( await claimedAtForUpdate ( tempDir , 42 ) ) . toBeGreaterThan ( before ) ,
1648+ ) ;
1649+
1650+ releaseHandler ?.( ) ;
1651+ await vi . waitFor ( async ( ) =>
1652+ expect ( await listTelegramSpooledUpdateClaims ( { spoolDir : tempDir } ) ) . toEqual ( [ ] ) ,
1653+ ) ;
1654+ } finally {
1655+ releaseHandler ?.( ) ;
1656+ abort . abort ( ) ;
1657+ stopWorker ( ) ;
1658+ refreshHarness . restore ( ) ;
1659+ await runPromise ;
1660+ }
1661+ } ) ;
1662+ } ) ;
1663+
15781664 it ( "holds buffered spooled claims until deferred processing settles without blocking same-lane buffering" , async ( ) => {
15791665 await withTempSpool ( async ( tempDir ) => {
15801666 const abort = new AbortController ( ) ;
@@ -1625,6 +1711,50 @@ describe("TelegramPollingSession", () => {
16251711 } ) ;
16261712 } ) ;
16271713
1714+ it ( "refreshes deferred spooled claims after the active handler hands off" , async ( ) => {
1715+ const refreshHarness = installSpooledClaimRefreshHarness ( ) ;
1716+ await withTempSpool ( async ( tempDir ) => {
1717+ const abort = new AbortController ( ) ;
1718+ const participants : TelegramSpooledReplayDeferredParticipant [ ] = [ ] ;
1719+ await writeSpooledTestUpdates ( tempDir , [ topicUpdate ( 42 , 10 , "buffered topic 10 turn" ) ] ) ;
1720+
1721+ const { runPromise, stopWorker } = startIsolatedIngressSession ( {
1722+ abort,
1723+ spoolDir : tempDir ,
1724+ handleUpdate : async ( update ) => {
1725+ const participant = createTelegramSpooledReplayDeferredParticipant (
1726+ `test-buffer:${ update . update_id } ` ,
1727+ ) ;
1728+ if ( ! participant ) {
1729+ throw new Error ( "expected spooled replay participant" ) ;
1730+ }
1731+ participants . push ( participant ) ;
1732+ } ,
1733+ } ) ;
1734+
1735+ try {
1736+ await vi . waitFor ( ( ) => expect ( participants ) . toHaveLength ( 1 ) ) ;
1737+ const before = await claimedAtForUpdate ( tempDir , 42 ) ;
1738+
1739+ refreshHarness . triggerRefresh ( ) ;
1740+ await vi . waitFor ( async ( ) =>
1741+ expect ( await claimedAtForUpdate ( tempDir , 42 ) ) . toBeGreaterThan ( before ) ,
1742+ ) ;
1743+
1744+ participants [ 0 ] ?. settle ( { kind : "completed" } ) ;
1745+ await vi . waitFor ( async ( ) =>
1746+ expect ( await listTelegramSpooledUpdateClaims ( { spoolDir : tempDir } ) ) . toEqual ( [ ] ) ,
1747+ ) ;
1748+ } finally {
1749+ participants [ 0 ] ?. settle ( { kind : "completed" } ) ;
1750+ abort . abort ( ) ;
1751+ stopWorker ( ) ;
1752+ refreshHarness . restore ( ) ;
1753+ await runPromise ;
1754+ }
1755+ } ) ;
1756+ } ) ;
1757+
16281758 it ( "releases buffered spooled claims for retry when deferred processing fails" , async ( ) => {
16291759 await withTempSpool ( async ( tempDir ) => {
16301760 const abort = new AbortController ( ) ;
@@ -3585,6 +3715,106 @@ describe("TelegramPollingSession", () => {
35853715 }
35863716 } ) ;
35873717
3718+ it ( "marks isolated ingress unhealthy when a spooled backlog stalls before handler timeout" , async ( ) => {
3719+ vi . useFakeTimers ( { now : 1_000 , shouldAdvanceTime : true } ) ;
3720+ const abort = new AbortController ( ) ;
3721+ const tempDir = await fs . mkdtemp ( path . join ( os . tmpdir ( ) , "openclaw-telegram-spool-" ) ) ;
3722+ const setStatus = vi . fn ( ) ;
3723+ let releaseRegularTurn : ( ( ) => void ) | undefined ;
3724+ const regularTurnDone = new Promise < void > ( ( resolve ) => {
3725+ releaseRegularTurn = resolve ;
3726+ } ) ;
3727+ const handleUpdate = vi . fn ( async ( ) => {
3728+ await regularTurnDone ;
3729+ } ) ;
3730+ createTelegramBotMock . mockReturnValueOnce ( {
3731+ api : {
3732+ deleteWebhook : vi . fn ( async ( ) => true ) ,
3733+ config : { use : vi . fn ( ) } ,
3734+ } ,
3735+ init : vi . fn ( async ( ) => undefined ) ,
3736+ handleUpdate,
3737+ stop : vi . fn ( async ( ) => undefined ) ,
3738+ } ) ;
3739+ await writeSpooledTestUpdates ( tempDir , [
3740+ topicUpdate ( 42 , 10 , "active topic 10 turn" ) ,
3741+ topicUpdate ( 43 , 10 , "later topic 10 turn" ) ,
3742+ ] ) ;
3743+
3744+ const workerListeners : WorkerMessageListener [ ] = [ ] ;
3745+ let stopWorker : ( ( ) => void ) | undefined ;
3746+ const workerDone = new Promise < void > ( ( resolve ) => {
3747+ stopWorker = resolve ;
3748+ } ) ;
3749+ const createWorker = vi . fn ( ( ) => ( {
3750+ onMessage : vi . fn ( ( listener : WorkerMessageListener ) => {
3751+ workerListeners . push ( listener ) ;
3752+ return ( ) => undefined ;
3753+ } ) ,
3754+ stop : vi . fn ( async ( ) => {
3755+ stopWorker ?.( ) ;
3756+ } ) ,
3757+ task : vi . fn ( async ( ) => {
3758+ await workerDone ;
3759+ } ) ,
3760+ } ) ) ;
3761+
3762+ try {
3763+ const session = createPollingSession ( {
3764+ abortSignal : abort . signal ,
3765+ setStatus,
3766+ isolatedIngress : {
3767+ enabled : true ,
3768+ spoolDir : tempDir ,
3769+ createWorker,
3770+ drainIntervalMs : pollingSessionTesting . isolatedIngressBacklogStallMs * 2 ,
3771+ spooledUpdateHandlerTimeoutMs : pollingSessionTesting . isolatedIngressBacklogStallMs * 2 ,
3772+ } ,
3773+ } ) ;
3774+
3775+ const runPromise = session . runUntilAbort ( ) ;
3776+ await vi . waitFor ( ( ) => expect ( handleUpdate ) . toHaveBeenCalledTimes ( 1 ) ) ;
3777+ workerListeners [ 0 ] ?.( {
3778+ type : "poll-success" ,
3779+ offset : null ,
3780+ count : 0 ,
3781+ finishedAt : Date . now ( ) ,
3782+ } ) ;
3783+ expect ( statusPatches ( setStatus ) . some ( ( patch ) => patch . connected === true ) ) . toBe ( true ) ;
3784+
3785+ vi . setSystemTime ( 1_000 + pollingSessionTesting . isolatedIngressBacklogStallMs + 1 ) ;
3786+ workerListeners [ 0 ] ?.( { type : "spooled" , updateId : 43 , queued : 1 } ) ;
3787+ await vi . waitFor ( ( ) =>
3788+ expect (
3789+ statusPatches ( setStatus ) . some (
3790+ ( patch ) =>
3791+ patch . connected === false &&
3792+ String ( patch . lastError ) . includes ( "isolated polling spool backlog stalled" ) ,
3793+ ) ,
3794+ ) . toBe ( true ) ,
3795+ ) ;
3796+ expect ( await failedUpdateIds ( tempDir ) ) . toEqual ( [ ] ) ;
3797+ expect ( await pendingUpdateIds ( tempDir , "all" ) ) . toEqual ( [ 43 ] ) ;
3798+ expect (
3799+ ( await listTelegramSpooledUpdateClaims ( { spoolDir : tempDir } ) ) . map (
3800+ ( claim ) => claim . updateId ,
3801+ ) ,
3802+ ) . toEqual ( [ 42 ] ) ;
3803+
3804+ releaseRegularTurn ?.( ) ;
3805+ abort . abort ( ) ;
3806+ stopWorker ?.( ) ;
3807+ await vi . advanceTimersByTimeAsync ( 20_000 ) ;
3808+ await runPromise ;
3809+ } finally {
3810+ releaseRegularTurn ?.( ) ;
3811+ abort . abort ( ) ;
3812+ stopWorker ?.( ) ;
3813+ vi . useRealTimers ( ) ;
3814+ await fs . rm ( tempDir , { recursive : true , force : true } ) ;
3815+ }
3816+ } ) ;
3817+
35883818 it ( "marks isolated ingress unhealthy when a spooled backlog handler times out" , async ( ) => {
35893819 vi . useFakeTimers ( { shouldAdvanceTime : true } ) ;
35903820 const abort = new AbortController ( ) ;
0 commit comments