@@ -74,6 +74,7 @@ let listTelegramSpooledUpdates: typeof import("./telegram-ingress-spool.js").lis
7474let recoverStaleTelegramSpooledUpdateClaims : typeof import ( "./telegram-ingress-spool.js" ) . recoverStaleTelegramSpooledUpdateClaims ;
7575let writeTelegramSpooledUpdate : typeof import ( "./telegram-ingress-spool.js" ) . writeTelegramSpooledUpdate ;
7676let createTelegramSpooledReplayDeferredParticipant : typeof import ( "./bot-processing-outcome.js" ) . createTelegramSpooledReplayDeferredParticipant ;
77+ let TelegramMessageDispatchReplayForgetError : typeof import ( "./message-dispatch-dedupe.js" ) . TelegramMessageDispatchReplayForgetError ;
7778type TelegramMessageProcessingResult =
7879 import ( "./bot-processing-outcome.js" ) . TelegramMessageProcessingResult ;
7980type TelegramSpooledReplayDeferredParticipant =
@@ -619,6 +620,7 @@ describe("TelegramPollingSession", () => {
619620 } = await import ( "./telegram-ingress-spool.js" ) ) ;
620621 ( { createTelegramSpooledReplayDeferredParticipant } =
621622 await import ( "./bot-processing-outcome.js" ) ) ;
623+ ( { TelegramMessageDispatchReplayForgetError } = await import ( "./message-dispatch-dedupe.js" ) ) ;
622624 ( {
623625 beginTelegramReplyFence,
624626 buildTelegramReplyFenceLaneKey,
@@ -1444,6 +1446,56 @@ describe("TelegramPollingSession", () => {
14441446 } ) ;
14451447 } ) ;
14461448
1449+ it ( "dead-letters buffered spooled claims when dispatch dedupe rollback fails" , async ( ) => {
1450+ await withTempSpool ( async ( tempDir ) => {
1451+ const abort = new AbortController ( ) ;
1452+ const log = vi . fn ( ) ;
1453+ const participants : TelegramSpooledReplayDeferredParticipant [ ] = [ ] ;
1454+ const events : string [ ] = [ ] ;
1455+ let attempts = 0 ;
1456+ await writeSpooledTestUpdates ( tempDir , [ topicUpdate ( 42 , 10 , "buffered rollback failure" ) ] ) ;
1457+
1458+ const { runPromise, stopWorker } = startIsolatedIngressSession ( {
1459+ abort,
1460+ spoolDir : tempDir ,
1461+ log,
1462+ drainIntervalMs : 10 ,
1463+ handleUpdate : async ( update ) => {
1464+ attempts += 1 ;
1465+ if ( attempts === 1 ) {
1466+ events . push ( `dispatch:${ update . update_id } ` ) ;
1467+ const participant = createTelegramSpooledReplayDeferredParticipant (
1468+ `test-buffer:${ update . update_id } ` ,
1469+ ) ;
1470+ if ( ! participant ) {
1471+ throw new Error ( "expected spooled replay participant" ) ;
1472+ }
1473+ participants . push ( participant ) ;
1474+ return ;
1475+ }
1476+ events . push ( `duplicate-skip:${ update . update_id } ` ) ;
1477+ } ,
1478+ } ) ;
1479+
1480+ await vi . waitFor ( ( ) => expect ( participants ) . toHaveLength ( 1 ) ) ;
1481+ participants [ 0 ] ?. settle ( {
1482+ kind : "failed-retryable" ,
1483+ error : new TelegramMessageDispatchReplayForgetError ( [ { key : "committed-dispatch-key" } ] ) ,
1484+ } ) ;
1485+
1486+ await vi . waitFor ( async ( ) => expect ( await failedUpdateIds ( tempDir ) ) . toEqual ( [ 42 ] ) ) ;
1487+ expect ( events ) . toEqual ( [ "dispatch:42" ] ) ;
1488+ expect ( await pendingUpdateIds ( tempDir , "all" ) ) . toEqual ( [ ] ) ;
1489+ expect ( await listTelegramSpooledUpdateClaims ( { spoolDir : tempDir } ) ) . toEqual ( [ ] ) ;
1490+ expectLogIncludes ( log , "non-retryable dispatch-dedupe-rollback-failed" ) ;
1491+ expectLogExcludes ( log , "spooled update 42 failed; keeping for retry" ) ;
1492+
1493+ abort . abort ( ) ;
1494+ stopWorker ( ) ;
1495+ await runPromise ;
1496+ } ) ;
1497+ } ) ;
1498+
14471499 it ( "fails buffered spooled claims instead of requeueing when deferred processing times out" , async ( ) => {
14481500 await withTempSpool ( async ( tempDir ) => {
14491501 const abort = new AbortController ( ) ;
0 commit comments