@@ -9,6 +9,7 @@ const createTelegramBotMock = vi.hoisted(() => vi.fn());
99const isRecoverableTelegramNetworkErrorMock = vi . hoisted ( ( ) => vi . fn ( ( ) => true ) ) ;
1010const computeBackoffMock = vi . hoisted ( ( ) => vi . fn ( ( ) => 0 ) ) ;
1111const sleepWithAbortMock = vi . hoisted ( ( ) => vi . fn ( async ( ) => undefined ) ) ;
12+ const drainPendingDeliveriesMock = vi . hoisted ( ( ) => vi . fn ( async ( _opts : unknown ) => undefined ) ) ;
1213
1314vi . mock ( "@grammyjs/runner" , ( ) => ( {
1415 run : runMock ,
@@ -22,6 +23,10 @@ vi.mock("./network-errors.js", () => ({
2223 isRecoverableTelegramNetworkError : isRecoverableTelegramNetworkErrorMock ,
2324} ) ) ;
2425
26+ vi . mock ( "openclaw/plugin-sdk/delivery-queue-runtime" , ( ) => ( {
27+ drainPendingDeliveries : drainPendingDeliveriesMock ,
28+ } ) ) ;
29+
2530vi . mock ( "./api-logging.js" , ( ) => ( {
2631 withTelegramApiErrorLogging : async ( { fn } : { fn : ( ) => Promise < unknown > } ) => await fn ( ) ,
2732} ) ) ;
@@ -54,6 +59,18 @@ type TelegramApiMiddleware = (
5459 method : string ,
5560 payload : unknown ,
5661) => Promise < unknown > ;
62+ type DrainPendingDeliveriesCall = {
63+ drainKey : string ;
64+ logLabel : string ;
65+ selectEntry : (
66+ entry : {
67+ channel : string ;
68+ accountId ?: string ;
69+ lastError ?: string ;
70+ } ,
71+ now : number ,
72+ ) => { match : boolean ; bypassBackoff : boolean } ;
73+ } ;
5774type AsyncVoidFn = ( ) => Promise < void > ;
5875type MockCallSource = { mock : { calls : Array < Array < unknown > > } } ;
5976
@@ -164,6 +181,14 @@ function expectTelegramBotTransportSequence(firstTransport: unknown, secondTrans
164181 expect ( createTelegramBotMock . mock . calls . at ( 1 ) ?. [ 0 ] ?. telegramTransport ) . toBe ( secondTransport ) ;
165182}
166183
184+ function expectDrainPendingDeliveriesCall ( index = 0 ) : DrainPendingDeliveriesCall {
185+ const call = drainPendingDeliveriesMock . mock . calls [ index ] ?. [ 0 ] ;
186+ if ( ! call || typeof call !== "object" ) {
187+ throw new Error ( `Expected drainPendingDeliveries call ${ index } ` ) ;
188+ }
189+ return call as DrainPendingDeliveriesCall ;
190+ }
191+
167192function makeTelegramTransport ( ) {
168193 return {
169194 fetch : globalThis . fetch ,
@@ -292,6 +317,7 @@ describe("TelegramPollingSession", () => {
292317 isRecoverableTelegramNetworkErrorMock . mockReset ( ) . mockReturnValue ( true ) ;
293318 computeBackoffMock . mockReset ( ) . mockReturnValue ( 0 ) ;
294319 sleepWithAbortMock . mockReset ( ) . mockResolvedValue ( undefined ) ;
320+ drainPendingDeliveriesMock . mockReset ( ) . mockResolvedValue ( undefined ) ;
295321 } ) ;
296322
297323 it ( "uses backoff helpers for recoverable polling retries" , async ( ) => {
@@ -454,6 +480,58 @@ describe("TelegramPollingSession", () => {
454480 }
455481 } ) ;
456482
483+ it ( "drains Telegram delivery queue after isolated ingress reports poll success" , async ( ) => {
484+ const abort = new AbortController ( ) ;
485+ const init = vi . fn ( async ( ) => undefined ) ;
486+ const bot = {
487+ api : {
488+ deleteWebhook : vi . fn ( async ( ) => true ) ,
489+ config : { use : vi . fn ( ) } ,
490+ } ,
491+ init,
492+ handleUpdate : vi . fn ( async ( ) => undefined ) ,
493+ stop : vi . fn ( async ( ) => undefined ) ,
494+ } ;
495+ createTelegramBotMock . mockReturnValueOnce ( bot ) ;
496+ let onMessage :
497+ | ( ( message : { type : "poll-success" ; finishedAt : number ; count : number } ) => void )
498+ | undefined ;
499+ let stopWorker : ( ( ) => void ) | undefined ;
500+ const workerDone = new Promise < void > ( ( resolve ) => {
501+ stopWorker = resolve ;
502+ } ) ;
503+ const createWorker = vi . fn ( ( ) => ( {
504+ onMessage : vi . fn ( ( handler ) => {
505+ onMessage = handler ;
506+ return ( ) => undefined ;
507+ } ) ,
508+ stop : vi . fn ( async ( ) => {
509+ stopWorker ?.( ) ;
510+ } ) ,
511+ task : vi . fn ( async ( ) => {
512+ await workerDone ;
513+ } ) ,
514+ } ) ) ;
515+
516+ const session = createPollingSession ( {
517+ abortSignal : abort . signal ,
518+ isolatedIngress : {
519+ enabled : true ,
520+ createWorker,
521+ drainIntervalMs : 10 ,
522+ } ,
523+ } ) ;
524+
525+ const runPromise = session . runUntilAbort ( ) ;
526+ await vi . waitFor ( ( ) => expect ( init ) . toHaveBeenCalledTimes ( 1 ) ) ;
527+ onMessage ?.( { type : "poll-success" , finishedAt : Date . now ( ) , count : 0 } ) ;
528+
529+ await vi . waitFor ( ( ) => expect ( drainPendingDeliveriesMock ) . toHaveBeenCalledTimes ( 1 ) ) ;
530+
531+ abort . abort ( ) ;
532+ await runPromise ;
533+ } ) ;
534+
457535 it ( "lets isolated ingress drain interleave different Telegram topic lanes" , async ( ) => {
458536 const abort = new AbortController ( ) ;
459537 const tempDir = await fs . mkdtemp ( path . join ( os . tmpdir ( ) , "openclaw-telegram-spool-" ) ) ;
@@ -1175,6 +1253,87 @@ describe("TelegramPollingSession", () => {
11751253 } ) ;
11761254 } ) ;
11771255
1256+ it ( "drains Telegram delivery queue after getUpdates confirms polling reconnect" , async ( ) => {
1257+ const abort = new AbortController ( ) ;
1258+ const botStop = vi . fn ( async ( ) => undefined ) ;
1259+ const runnerStop = vi . fn ( async ( ) => undefined ) ;
1260+ const getApiMiddleware = mockBotCapturingApiMiddleware ( botStop ) ;
1261+ const resolveFirstTask = mockLongRunningPollingCycle ( runnerStop ) ;
1262+
1263+ const session = createPollingSession ( {
1264+ abortSignal : abort . signal ,
1265+ } ) ;
1266+
1267+ const runPromise = session . runUntilAbort ( ) ;
1268+ const apiMiddleware = await waitForApiMiddleware ( getApiMiddleware ) ;
1269+ await apiMiddleware (
1270+ vi . fn ( async ( ) => [ ] ) ,
1271+ "getUpdates" ,
1272+ { offset : 1 } ,
1273+ ) ;
1274+
1275+ await vi . waitFor ( ( ) => expect ( drainPendingDeliveriesMock ) . toHaveBeenCalledTimes ( 1 ) ) ;
1276+ const drain = expectDrainPendingDeliveriesCall ( ) ;
1277+ expect ( drain . drainKey ) . toBe ( "telegram:default" ) ;
1278+ expect ( drain . logLabel ) . toBe ( "Telegram reconnect drain" ) ;
1279+ expect ( drain . selectEntry ( { channel : "telegram" } , Date . now ( ) ) ) . toEqual ( {
1280+ match : true ,
1281+ bypassBackoff : false ,
1282+ } ) ;
1283+ expect (
1284+ drain . selectEntry (
1285+ {
1286+ channel : "telegram" ,
1287+ accountId : "default" ,
1288+ lastError : "Network request for 'sendMessage' failed!" ,
1289+ } ,
1290+ Date . now ( ) ,
1291+ ) ,
1292+ ) . toEqual ( {
1293+ match : true ,
1294+ bypassBackoff : true ,
1295+ } ) ;
1296+ expect ( drain . selectEntry ( { channel : "telegram" , accountId : "alerts" } , Date . now ( ) ) . match ) . toBe (
1297+ false ,
1298+ ) ;
1299+ expect ( drain . selectEntry ( { channel : "whatsapp" } , Date . now ( ) ) . match ) . toBe ( false ) ;
1300+
1301+ abort . abort ( ) ;
1302+ resolveFirstTask ( ) ;
1303+ await runPromise ;
1304+ } ) ;
1305+
1306+ it ( "drains Telegram delivery queue after each getUpdates success" , async ( ) => {
1307+ const abort = new AbortController ( ) ;
1308+ const botStop = vi . fn ( async ( ) => undefined ) ;
1309+ const runnerStop = vi . fn ( async ( ) => undefined ) ;
1310+ const getApiMiddleware = mockBotCapturingApiMiddleware ( botStop ) ;
1311+ const resolveFirstTask = mockLongRunningPollingCycle ( runnerStop ) ;
1312+
1313+ const session = createPollingSession ( {
1314+ abortSignal : abort . signal ,
1315+ } ) ;
1316+
1317+ const runPromise = session . runUntilAbort ( ) ;
1318+ const apiMiddleware = await waitForApiMiddleware ( getApiMiddleware ) ;
1319+ await apiMiddleware (
1320+ vi . fn ( async ( ) => [ ] ) ,
1321+ "getUpdates" ,
1322+ { offset : 1 } ,
1323+ ) ;
1324+ await apiMiddleware (
1325+ vi . fn ( async ( ) => [ ] ) ,
1326+ "getUpdates" ,
1327+ { offset : 2 } ,
1328+ ) ;
1329+
1330+ await vi . waitFor ( ( ) => expect ( drainPendingDeliveriesMock ) . toHaveBeenCalledTimes ( 2 ) ) ;
1331+
1332+ abort . abort ( ) ;
1333+ resolveFirstTask ( ) ;
1334+ await runPromise ;
1335+ } ) ;
1336+
11781337 it ( "keeps polling marked connected across recoverable restart cycles" , async ( ) => {
11791338 const abort = new AbortController ( ) ;
11801339 const recoverableError = new Error ( "recoverable polling error" ) ;
0 commit comments