@@ -3,7 +3,16 @@ import os from "node:os";
33import path from "node:path" ;
44import type { ChannelAccountSnapshot } from "openclaw/plugin-sdk/channel-contract" ;
55import { MAX_TIMER_TIMEOUT_MS } from "openclaw/plugin-sdk/number-runtime" ;
6- import { beforeAll , beforeEach , describe , expect , it , vi } from "vitest" ;
6+ import { afterEach , beforeAll , beforeEach , describe , expect , it , vi } from "vitest" ;
7+ import { createChannelIngressQueue } from "../../../src/channels/message/ingress-queue.js" ;
8+ import { executeSqliteQuerySync , getNodeSqliteKysely } from "../../../src/infra/kysely-sync.js" ;
9+ import type { DB as OpenClawStateKyselyDatabase } from "../../../src/state/openclaw-state-db.generated.js" ;
10+ import {
11+ closeOpenClawStateDatabaseForTest ,
12+ openOpenClawStateDatabase ,
13+ } from "../../../src/state/openclaw-state-db.js" ;
14+ import { clearTelegramRuntime , setTelegramRuntime } from "./runtime.js" ;
15+ import type { TelegramRuntime } from "./runtime.types.js" ;
716import type { TelegramIngressWorkerMessage } from "./telegram-ingress-worker.js" ;
817
918const runMock = vi . hoisted ( ( ) => vi . fn ( ) ) ;
@@ -96,9 +105,21 @@ type WorkerPollErrorListener = (message: {
96105type WorkerMessageListener = ( message : TelegramIngressWorkerMessage ) => void ;
97106type AsyncVoidFn = ( ) => Promise < void > ;
98107type MockCallSource = { mock : { calls : Array < Array < unknown > > } } ;
108+ type TelegramPollingTestDatabase = Pick < OpenClawStateKyselyDatabase , "channel_ingress_events" > ;
99109
100110const POLLING_TEST_WATCHDOG_INTERVAL_MS = 30_000 ;
101111
112+ function installTelegramIngressQueueRuntime ( resolveStateDir : ( ) => string ) : void {
113+ setTelegramRuntime ( {
114+ state : {
115+ resolveStateDir,
116+ openChannelIngressQueue : (
117+ options ?: Omit < Parameters < typeof createChannelIngressQueue > [ 0 ] , "channelId" > ,
118+ ) => createChannelIngressQueue ( { ...( options ?? { } ) , channelId : "telegram" } ) ,
119+ } ,
120+ } as TelegramRuntime ) ;
121+ }
122+
102123function mockObjectArg (
103124 source : MockCallSource ,
104125 label : string ,
@@ -411,24 +432,67 @@ async function pendingUpdateIds(spoolDir: string, limit: number | "all" = 100):
411432 return ( await listTelegramSpooledUpdates ( { spoolDir, limit } ) ) . map ( ( update ) => update . updateId ) ;
412433}
413434
414- async function failedUpdateIds ( spoolDir : string ) : Promise < number [ ] > {
415- const entries = await fs . readdir ( spoolDir ) . catch ( ( err ) => {
416- if ( ( err as { code ?: string } ) . code === "ENOENT" ) {
417- return [ ] ;
418- }
419- throw err ;
435+ function normalizeTelegramTestAccountId ( spoolDir : string ) : string {
436+ const trimmed = path . basename ( spoolDir ) . trim ( ) ;
437+ return trimmed ? trimmed . replace ( / [ ^ a - z 0 - 9 . _ - ] + / gi, "_" ) : "default" ;
438+ }
439+
440+ function telegramTestQueueName ( spoolDir : string ) : string {
441+ return JSON . stringify ( [ "telegram" , normalizeTelegramTestAccountId ( spoolDir ) ] ) ;
442+ }
443+
444+ function openTelegramSpoolTestKysely ( spoolDir : string ) {
445+ const database = openOpenClawStateDatabase ( {
446+ env : { ...process . env , OPENCLAW_STATE_DIR : spoolDir } ,
420447 } ) ;
421- return entries
422- . filter ( ( entry ) => entry . endsWith ( ".json.failed" ) )
423- . map ( ( entry ) => Number ( entry . slice ( 0 , 16 ) ) )
424- . toSorted ( ( a , b ) => a - b ) ;
448+ return {
449+ database,
450+ kysely : getNodeSqliteKysely < TelegramPollingTestDatabase > ( database . db ) ,
451+ } ;
452+ }
453+
454+ async function failedUpdateIds ( spoolDir : string ) : Promise < number [ ] > {
455+ const { database, kysely } = openTelegramSpoolTestKysely ( spoolDir ) ;
456+ const rows = executeSqliteQuerySync (
457+ database . db ,
458+ kysely
459+ . selectFrom ( "channel_ingress_events" )
460+ . select ( "event_id" )
461+ . where ( "queue_name" , "=" , telegramTestQueueName ( spoolDir ) )
462+ . where ( "status" , "=" , "failed" )
463+ . orderBy ( "event_id" , "asc" ) ,
464+ ) . rows ;
465+ return rows . map ( ( row ) => Number ( row . event_id ) ) ;
466+ }
467+
468+ async function adoptClaimOwner ( params : {
469+ spoolDir : string ;
470+ updateId : number ;
471+ ownerId : string ;
472+ claimedAt : number ;
473+ } ) : Promise < void > {
474+ const { database, kysely } = openTelegramSpoolTestKysely ( params . spoolDir ) ;
475+ executeSqliteQuerySync (
476+ database . db ,
477+ kysely
478+ . updateTable ( "channel_ingress_events" )
479+ . set ( {
480+ claim_owner : params . ownerId ,
481+ claimed_at : params . claimedAt ,
482+ updated_at : params . claimedAt ,
483+ } )
484+ . where ( "queue_name" , "=" , telegramTestQueueName ( params . spoolDir ) )
485+ . where ( "event_id" , "=" , String ( params . updateId ) . padStart ( 16 , "0" ) )
486+ . where ( "status" , "=" , "claimed" ) ,
487+ ) ;
425488}
426489
427490async function withTempSpool < T > ( fn : ( spoolDir : string ) => Promise < T > ) : Promise < T > {
428491 const spoolDir = await fs . mkdtemp ( path . join ( os . tmpdir ( ) , "openclaw-telegram-spool-" ) ) ;
429492 try {
430493 return await fn ( spoolDir ) ;
431494 } finally {
495+ closeOpenClawStateDatabaseForTest ( ) ;
432496 await fs . rm ( spoolDir , { recursive : true , force : true } ) ;
433497 }
434498}
@@ -526,6 +590,14 @@ describe("TelegramPollingSession", () => {
526590 sleepWithAbortMock . mockReset ( ) . mockResolvedValue ( undefined ) ;
527591 drainPendingDeliveriesMock . mockReset ( ) . mockResolvedValue ( undefined ) ;
528592 resetTelegramReplyFenceForTests ( ) ;
593+ installTelegramIngressQueueRuntime ( ( ) =>
594+ path . join ( os . tmpdir ( ) , "openclaw-telegram-test-state" ) ,
595+ ) ;
596+ } ) ;
597+
598+ afterEach ( ( ) => {
599+ clearTelegramRuntime ( ) ;
600+ closeOpenClawStateDatabaseForTest ( ) ;
529601 } ) ;
530602
531603 it ( "uses backoff helpers for recoverable polling retries" , async ( ) => {
@@ -667,7 +739,14 @@ describe("TelegramPollingSession", () => {
667739
668740 const runPromise = session . runUntilAbort ( ) ;
669741 await vi . waitFor ( ( ) => expect ( handleUpdate ) . toHaveBeenCalledTimes ( 1 ) ) ;
670- await vi . waitFor ( async ( ) => expect ( await fs . readdir ( tempDir ) ) . toEqual ( [ ] ) ) ;
742+ await vi . waitFor ( async ( ) => expect ( await pendingUpdateIds ( tempDir , "all" ) ) . toEqual ( [ ] ) ) ;
743+ await vi . waitFor ( async ( ) =>
744+ expect (
745+ await listTelegramSpooledUpdateClaims ( {
746+ spoolDir : tempDir ,
747+ } ) ,
748+ ) . toEqual ( [ ] ) ,
749+ ) ;
671750 abort . abort ( ) ;
672751 await runPromise ;
673752
@@ -686,6 +765,76 @@ describe("TelegramPollingSession", () => {
686765 expect ( init ) . toHaveBeenCalledBefore ( handleUpdate ) ;
687766 expect ( handleUpdate ) . toHaveBeenCalledWith ( { update_id : 42 , message : { text : "hello" } } ) ;
688767 } finally {
768+ abort . abort ( ) ;
769+ await fs . rm ( tempDir , { recursive : true , force : true } ) ;
770+ }
771+ } ) ;
772+
773+ it ( "writes isolated worker updates through the main runtime queue" , async ( ) => {
774+ const abort = new AbortController ( ) ;
775+ const tempDir = await fs . mkdtemp ( path . join ( os . tmpdir ( ) , "openclaw-telegram-spool-" ) ) ;
776+ const handleUpdate = vi . fn ( async ( ) => undefined ) ;
777+ const bot = {
778+ api : {
779+ deleteWebhook : vi . fn ( async ( ) => true ) ,
780+ config : { use : vi . fn ( ) } ,
781+ } ,
782+ init : vi . fn ( async ( ) => undefined ) ,
783+ handleUpdate,
784+ stop : vi . fn ( async ( ) => undefined ) ,
785+ } ;
786+ createTelegramBotMock . mockReturnValueOnce ( bot ) ;
787+ let onMessage : WorkerMessageListener | undefined ;
788+ let stopWorker : ( ( ) => void ) | undefined ;
789+ const workerDone = new Promise < void > ( ( resolve ) => {
790+ stopWorker = resolve ;
791+ } ) ;
792+ const ackSpooledUpdate = vi . fn ( ) ;
793+ const createWorker = vi . fn ( ( ) => ( {
794+ onMessage : vi . fn ( ( listener : WorkerMessageListener ) => {
795+ onMessage = listener ;
796+ return ( ) => undefined ;
797+ } ) ,
798+ ackSpooledUpdate,
799+ stop : vi . fn ( async ( ) => {
800+ stopWorker ?.( ) ;
801+ } ) ,
802+ task : vi . fn ( async ( ) => {
803+ await workerDone ;
804+ } ) ,
805+ } ) ) ;
806+
807+ try {
808+ const session = createPollingSession ( {
809+ abortSignal : abort . signal ,
810+ isolatedIngress : {
811+ enabled : true ,
812+ spoolDir : tempDir ,
813+ createWorker,
814+ drainIntervalMs : 10 ,
815+ } ,
816+ } ) ;
817+
818+ const runPromise = session . runUntilAbort ( ) ;
819+ await vi . waitFor ( ( ) => expect ( onMessage ) . toBeDefined ( ) ) ;
820+ onMessage ?.( {
821+ type : "update" ,
822+ requestId : "write-1" ,
823+ update : { update_id : 42 , message : { text : "hello" } } ,
824+ queued : 1 ,
825+ } ) ;
826+
827+ await vi . waitFor ( ( ) =>
828+ expect ( ackSpooledUpdate ) . toHaveBeenCalledWith ( "write-1" , { ok : true , updateId : 42 } ) ,
829+ ) ;
830+ await vi . waitFor ( ( ) =>
831+ expect ( handleUpdate ) . toHaveBeenCalledWith ( { update_id : 42 , message : { text : "hello" } } ) ,
832+ ) ;
833+ await vi . waitFor ( async ( ) => expect ( await pendingUpdateIds ( tempDir , "all" ) ) . toEqual ( [ ] ) ) ;
834+ abort . abort ( ) ;
835+ await runPromise ;
836+ } finally {
837+ abort . abort ( ) ;
689838 await fs . rm ( tempDir , { recursive : true , force : true } ) ;
690839 }
691840 } ) ;
@@ -735,7 +884,14 @@ describe("TelegramPollingSession", () => {
735884
736885 const runPromise = session . runUntilAbort ( ) ;
737886 await vi . waitFor ( ( ) => expect ( handleUpdate ) . toHaveBeenCalledTimes ( 1 ) ) ;
738- await vi . waitFor ( async ( ) => expect ( await fs . readdir ( tempDir ) ) . toEqual ( [ ] ) ) ;
887+ await vi . waitFor ( async ( ) => expect ( await pendingUpdateIds ( tempDir , "all" ) ) . toEqual ( [ ] ) ) ;
888+ await vi . waitFor ( async ( ) =>
889+ expect (
890+ await listTelegramSpooledUpdateClaims ( {
891+ spoolDir : tempDir ,
892+ } ) ,
893+ ) . toEqual ( [ ] ) ,
894+ ) ;
739895 abort . abort ( ) ;
740896 await runPromise ;
741897
@@ -1128,7 +1284,7 @@ describe("TelegramPollingSession", () => {
11281284 await runPromise ;
11291285 expect ( events ) . toEqual ( [ "handled:42" , "handled:44" ] ) ;
11301286 expect ( await pendingUpdateIds ( tempDir ) ) . toEqual ( [ 43 ] ) ;
1131- expect ( ( await fs . readdir ( tempDir ) ) . toSorted ( ) ) . toEqual ( [ "0000000000000043.json" ] ) ;
1287+ expect ( await listTelegramSpooledUpdateClaims ( { spoolDir : tempDir } ) ) . toEqual ( [ ] ) ;
11321288 stopWorker ( ) ;
11331289 } ) ;
11341290 } ) ;
@@ -1189,21 +1345,12 @@ describe("TelegramPollingSession", () => {
11891345 if ( ! claimed ) {
11901346 throw new Error ( "Expected claimed update" ) ;
11911347 }
1192- await fs . writeFile (
1193- claimed . path ,
1194- `${ JSON . stringify ( {
1195- version : 1 ,
1196- updateId : 42 ,
1197- receivedAt : interrupted . receivedAt ,
1198- update : interruptedUpdate ,
1199- claim : {
1200- processId : "other-process" ,
1201- processPid : process . pid ,
1202- claimedAt : Date . now ( ) ,
1203- } ,
1204- } ) } \n`,
1205- { mode : 0o600 } ,
1206- ) ;
1348+ await adoptClaimOwner ( {
1349+ spoolDir : tempDir ,
1350+ updateId : 42 ,
1351+ ownerId : `${ process . pid } :other-process` ,
1352+ claimedAt : Date . now ( ) ,
1353+ } ) ;
12071354
12081355 const recovered = await recoverStaleTelegramSpooledUpdateClaims ( {
12091356 spoolDir : tempDir ,
@@ -1213,10 +1360,11 @@ describe("TelegramPollingSession", () => {
12131360
12141361 expect ( recovered ) . toBe ( 0 ) ;
12151362 expect ( await pendingUpdateIds ( tempDir ) ) . toEqual ( [ 43 ] ) ;
1216- expect ( ( await fs . readdir ( tempDir ) ) . toSorted ( ) ) . toEqual ( [
1217- "0000000000000042.json.processing" ,
1218- "0000000000000043.json" ,
1219- ] ) ;
1363+ expect (
1364+ ( await listTelegramSpooledUpdateClaims ( { spoolDir : tempDir } ) ) . map (
1365+ ( claim ) => claim . updateId ,
1366+ ) ,
1367+ ) . toEqual ( [ 42 ] ) ;
12201368 } ) ;
12211369 } ) ;
12221370
@@ -2360,7 +2508,7 @@ describe("TelegramPollingSession", () => {
23602508 }
23612509 } ) ;
23622510
2363- it ( "keeps a timed-out lane guarded when its failed tombstone cannot be written" , async ( ) => {
2511+ it ( "keeps a timed-out lane guarded when its failed state cannot be written" , async ( ) => {
23642512 vi . useFakeTimers ( { shouldAdvanceTime : true } ) ;
23652513 const abort = new AbortController ( ) ;
23662514 const tempDir = await fs . mkdtemp ( path . join ( os . tmpdir ( ) , "openclaw-telegram-spool-" ) ) ;
@@ -2371,15 +2519,10 @@ describe("TelegramPollingSession", () => {
23712519 const regularTurnDone = new Promise < void > ( ( resolve ) => {
23722520 releaseRegularTurn = resolve ;
23732521 } ) ;
2374- const originalWriteFile = fs . writeFile . bind ( fs ) ;
2375- const writeFileSpy = vi
2376- . spyOn ( fs , "writeFile" )
2377- . mockImplementation ( async ( ...args : Parameters < typeof fs . writeFile > ) => {
2378- if ( typeof args [ 0 ] === "string" && args [ 0 ] . includes ( ".json.failed." ) ) {
2379- throw new Error ( "disk full" ) ;
2380- }
2381- return await originalWriteFile ( ...args ) ;
2382- } ) ;
2522+ const spoolModule = await import ( "./telegram-ingress-spool.js" ) ;
2523+ const failSpy = vi
2524+ . spyOn ( spoolModule , "failTelegramSpooledUpdateClaim" )
2525+ . mockRejectedValueOnce ( new Error ( "disk full" ) ) ;
23832526 createTelegramBotMock . mockReturnValueOnce ( {
23842527 api : {
23852528 deleteWebhook : vi . fn ( async ( ) => true ) ,
@@ -2460,7 +2603,7 @@ describe("TelegramPollingSession", () => {
24602603 await vi . advanceTimersByTimeAsync ( 20_000 ) ;
24612604 await runPromise ;
24622605 } finally {
2463- writeFileSpy . mockRestore ( ) ;
2606+ failSpy . mockRestore ( ) ;
24642607 releaseRegularTurn ?.( ) ;
24652608 abort . abort ( ) ;
24662609 stopWorker ?.( ) ;
0 commit comments