@@ -71,6 +71,12 @@ type DrainPendingDeliveriesCall = {
7171 now : number ,
7272 ) => { match : boolean ; bypassBackoff : boolean } ;
7373} ;
74+ type WorkerPollSuccessListener = ( message : {
75+ type : "poll-success" ;
76+ offset : null ;
77+ count : number ;
78+ finishedAt : number ;
79+ } ) => void ;
7480type AsyncVoidFn = ( ) => Promise < void > ;
7581type MockCallSource = { mock : { calls : Array < Array < unknown > > } } ;
7682
@@ -882,6 +888,309 @@ describe("TelegramPollingSession", () => {
882888 }
883889 } ) ;
884890
891+ it ( "keeps active spooled lanes blocked across isolated ingress restarts" , async ( ) => {
892+ vi . useFakeTimers ( { shouldAdvanceTime : true } ) ;
893+ const abort = new AbortController ( ) ;
894+ const tempDir = await fs . mkdtemp ( path . join ( os . tmpdir ( ) , "openclaw-telegram-spool-" ) ) ;
895+ let releaseRegularTurn : ( ( ) => void ) | undefined ;
896+ const regularTurnDone = new Promise < void > ( ( resolve ) => {
897+ releaseRegularTurn = resolve ;
898+ } ) ;
899+ const handleUpdate = vi . fn ( async ( ) => {
900+ await regularTurnDone ;
901+ } ) ;
902+ createTelegramBotMock . mockImplementation ( ( ) => ( {
903+ api : {
904+ deleteWebhook : vi . fn ( async ( ) => true ) ,
905+ config : { use : vi . fn ( ) } ,
906+ } ,
907+ init : vi . fn ( async ( ) => undefined ) ,
908+ handleUpdate,
909+ stop : vi . fn ( async ( ) => undefined ) ,
910+ } ) ) ;
911+ await writeTelegramSpooledUpdate ( {
912+ spoolDir : tempDir ,
913+ update : {
914+ update_id : 42 ,
915+ message : { text : "summarize this" , chat : { id : - 100 , type : "supergroup" } } ,
916+ } ,
917+ } ) ;
918+
919+ let workerTaskCalls = 0 ;
920+ let stopWorker : ( ( ) => void ) | undefined ;
921+ const workerDone = new Promise < void > ( ( resolve ) => {
922+ stopWorker = resolve ;
923+ } ) ;
924+ const createWorker = vi . fn ( ( ) => ( {
925+ onMessage : vi . fn ( ( ) => ( ) => undefined ) ,
926+ stop : vi . fn ( async ( ) => {
927+ stopWorker ?.( ) ;
928+ } ) ,
929+ task : vi . fn ( async ( ) => {
930+ workerTaskCalls += 1 ;
931+ if ( workerTaskCalls === 1 ) {
932+ return ;
933+ }
934+ await workerDone ;
935+ } ) ,
936+ } ) ) ;
937+
938+ try {
939+ const session = createPollingSession ( {
940+ abortSignal : abort . signal ,
941+ isolatedIngress : {
942+ enabled : true ,
943+ spoolDir : tempDir ,
944+ createWorker,
945+ drainIntervalMs : 100 ,
946+ } ,
947+ } ) ;
948+
949+ const runPromise = session . runUntilAbort ( ) ;
950+ await vi . waitFor ( ( ) => expect ( handleUpdate ) . toHaveBeenCalledTimes ( 1 ) ) ;
951+ await vi . advanceTimersByTimeAsync ( 16_000 ) ;
952+ await vi . waitFor ( ( ) => expect ( createWorker ) . toHaveBeenCalledTimes ( 2 ) ) ;
953+ expect ( handleUpdate ) . toHaveBeenCalledTimes ( 1 ) ;
954+
955+ releaseRegularTurn ?.( ) ;
956+ await vi . advanceTimersByTimeAsync ( 1_000 ) ;
957+ await vi . waitFor ( async ( ) =>
958+ expect (
959+ ( await listTelegramSpooledUpdates ( { spoolDir : tempDir } ) ) . map (
960+ ( update ) => update . updateId ,
961+ ) ,
962+ ) . toEqual ( [ ] ) ,
963+ ) ;
964+ abort . abort ( ) ;
965+ await vi . advanceTimersByTimeAsync ( 20_000 ) ;
966+ await runPromise ;
967+ } finally {
968+ releaseRegularTurn ?.( ) ;
969+ vi . useRealTimers ( ) ;
970+ await fs . rm ( tempDir , { recursive : true , force : true } ) ;
971+ }
972+ } ) ;
973+
974+ it ( "keeps active spooled lanes blocked across account restarts" , async ( ) => {
975+ vi . useFakeTimers ( { shouldAdvanceTime : true } ) ;
976+ const firstAbort = new AbortController ( ) ;
977+ const secondAbort = new AbortController ( ) ;
978+ const tempDir = await fs . mkdtemp ( path . join ( os . tmpdir ( ) , "openclaw-telegram-spool-" ) ) ;
979+ let releaseRegularTurn : ( ( ) => void ) | undefined ;
980+ const regularTurnDone = new Promise < void > ( ( resolve ) => {
981+ releaseRegularTurn = resolve ;
982+ } ) ;
983+ const handleUpdate = vi . fn ( async ( ) => {
984+ await regularTurnDone ;
985+ } ) ;
986+ createTelegramBotMock . mockImplementation ( ( ) => ( {
987+ api : {
988+ deleteWebhook : vi . fn ( async ( ) => true ) ,
989+ config : { use : vi . fn ( ) } ,
990+ } ,
991+ init : vi . fn ( async ( ) => undefined ) ,
992+ handleUpdate,
993+ stop : vi . fn ( async ( ) => undefined ) ,
994+ } ) ) ;
995+ await writeTelegramSpooledUpdate ( {
996+ spoolDir : tempDir ,
997+ update : {
998+ update_id : 42 ,
999+ message : { text : "summarize this" , chat : { id : - 100 , type : "supergroup" } } ,
1000+ } ,
1001+ } ) ;
1002+
1003+ const createWorker = vi . fn ( ( ) => {
1004+ let stopWorker : ( ( ) => void ) | undefined ;
1005+ const workerDone = new Promise < void > ( ( resolve ) => {
1006+ stopWorker = resolve ;
1007+ } ) ;
1008+ return {
1009+ onMessage : vi . fn ( ( ) => ( ) => undefined ) ,
1010+ stop : vi . fn ( async ( ) => {
1011+ stopWorker ?.( ) ;
1012+ } ) ,
1013+ task : vi . fn ( async ( ) => {
1014+ await workerDone ;
1015+ } ) ,
1016+ } ;
1017+ } ) ;
1018+
1019+ try {
1020+ const firstSession = createPollingSession ( {
1021+ abortSignal : firstAbort . signal ,
1022+ isolatedIngress : {
1023+ enabled : true ,
1024+ spoolDir : tempDir ,
1025+ createWorker,
1026+ drainIntervalMs : 100 ,
1027+ } ,
1028+ } ) ;
1029+
1030+ const firstRunPromise = firstSession . runUntilAbort ( ) ;
1031+ await vi . waitFor ( ( ) => expect ( handleUpdate ) . toHaveBeenCalledTimes ( 1 ) ) ;
1032+ firstAbort . abort ( ) ;
1033+ await vi . advanceTimersByTimeAsync ( 16_000 ) ;
1034+ await firstRunPromise ;
1035+
1036+ const secondSession = createPollingSession ( {
1037+ abortSignal : secondAbort . signal ,
1038+ isolatedIngress : {
1039+ enabled : true ,
1040+ spoolDir : tempDir ,
1041+ createWorker,
1042+ drainIntervalMs : 100 ,
1043+ } ,
1044+ } ) ;
1045+ const secondRunPromise = secondSession . runUntilAbort ( ) ;
1046+ await vi . waitFor ( ( ) => expect ( createWorker ) . toHaveBeenCalledTimes ( 2 ) ) ;
1047+ await vi . advanceTimersByTimeAsync ( 1_000 ) ;
1048+ expect ( handleUpdate ) . toHaveBeenCalledTimes ( 1 ) ;
1049+
1050+ releaseRegularTurn ?.( ) ;
1051+ await vi . advanceTimersByTimeAsync ( 1_000 ) ;
1052+ await vi . waitFor ( async ( ) =>
1053+ expect (
1054+ ( await listTelegramSpooledUpdates ( { spoolDir : tempDir } ) ) . map (
1055+ ( update ) => update . updateId ,
1056+ ) ,
1057+ ) . toEqual ( [ ] ) ,
1058+ ) ;
1059+ secondAbort . abort ( ) ;
1060+ await vi . advanceTimersByTimeAsync ( 20_000 ) ;
1061+ await secondRunPromise ;
1062+ } finally {
1063+ releaseRegularTurn ?.( ) ;
1064+ vi . useRealTimers ( ) ;
1065+ await fs . rm ( tempDir , { recursive : true , force : true } ) ;
1066+ }
1067+ } ) ;
1068+
1069+ it ( "marks isolated ingress unhealthy when a spooled backlog wedges while polling stays live" , async ( ) => {
1070+ vi . useFakeTimers ( { shouldAdvanceTime : true } ) ;
1071+ const abort = new AbortController ( ) ;
1072+ const tempDir = await fs . mkdtemp ( path . join ( os . tmpdir ( ) , "openclaw-telegram-spool-" ) ) ;
1073+ const log = vi . fn ( ) ;
1074+ const setStatus = vi . fn ( ) ;
1075+ let releaseRegularTurn : ( ( ) => void ) | undefined ;
1076+ const regularTurnDone = new Promise < void > ( ( resolve ) => {
1077+ releaseRegularTurn = resolve ;
1078+ } ) ;
1079+ const handleUpdate = vi . fn ( async ( ) => {
1080+ await regularTurnDone ;
1081+ } ) ;
1082+ createTelegramBotMock . mockImplementation ( ( ) => ( {
1083+ api : {
1084+ deleteWebhook : vi . fn ( async ( ) => true ) ,
1085+ config : { use : vi . fn ( ) } ,
1086+ } ,
1087+ init : vi . fn ( async ( ) => undefined ) ,
1088+ handleUpdate,
1089+ stop : vi . fn ( async ( ) => undefined ) ,
1090+ } ) ) ;
1091+ for ( const updateId of [ 42 , 43 ] ) {
1092+ await writeTelegramSpooledUpdate ( {
1093+ spoolDir : tempDir ,
1094+ update : {
1095+ update_id : updateId ,
1096+ message : { text : `dm ${ updateId } ` , chat : { id : 123 , type : "private" } } ,
1097+ } ,
1098+ } ) ;
1099+ }
1100+
1101+ const workerListeners : WorkerPollSuccessListener [ ] = [ ] ;
1102+ const createWorker = vi . fn ( ( ) => {
1103+ let stopWorker : ( ( ) => void ) | undefined ;
1104+ const workerDone = new Promise < void > ( ( resolve ) => {
1105+ stopWorker = resolve ;
1106+ } ) ;
1107+ return {
1108+ onMessage : vi . fn ( ( listener : WorkerPollSuccessListener ) => {
1109+ workerListeners . push ( listener ) ;
1110+ return ( ) => undefined ;
1111+ } ) ,
1112+ stop : vi . fn ( async ( ) => {
1113+ stopWorker ?.( ) ;
1114+ } ) ,
1115+ task : vi . fn ( async ( ) => {
1116+ await workerDone ;
1117+ } ) ,
1118+ } ;
1119+ } ) ;
1120+
1121+ try {
1122+ const session = createPollingSession ( {
1123+ abortSignal : abort . signal ,
1124+ log,
1125+ setStatus,
1126+ isolatedIngress : {
1127+ enabled : true ,
1128+ spoolDir : tempDir ,
1129+ createWorker,
1130+ drainIntervalMs : 100 ,
1131+ } ,
1132+ } ) ;
1133+
1134+ const runPromise = session . runUntilAbort ( ) ;
1135+ await vi . waitFor ( ( ) => expect ( handleUpdate ) . toHaveBeenCalledTimes ( 1 ) ) ;
1136+ workerListeners [ 0 ] ?.( {
1137+ type : "poll-success" ,
1138+ offset : null ,
1139+ count : 0 ,
1140+ finishedAt : Date . now ( ) ,
1141+ } ) ;
1142+ expect ( statusPatches ( setStatus ) . some ( ( patch ) => patch . connected === true ) ) . toBe ( true ) ;
1143+
1144+ await vi . advanceTimersByTimeAsync ( 25 * 60_000 + 100 ) ;
1145+
1146+ await vi . waitFor ( ( ) =>
1147+ expect ( log ) . toHaveBeenCalledWith (
1148+ expect . stringContaining ( "isolated polling spool backlog stalled" ) ,
1149+ ) ,
1150+ ) ;
1151+ expect (
1152+ statusPatches ( setStatus ) . some (
1153+ ( patch ) =>
1154+ patch . connected === false &&
1155+ String ( patch . lastError ) . includes ( "isolated polling spool backlog stalled" ) ,
1156+ ) ,
1157+ ) . toBe ( true ) ;
1158+ workerListeners [ 0 ] ?.( {
1159+ type : "poll-success" ,
1160+ offset : null ,
1161+ count : 0 ,
1162+ finishedAt : Date . now ( ) ,
1163+ } ) ;
1164+ expect ( statusPatches ( setStatus ) . at ( - 1 ) ?. connected ) . toBe ( false ) ;
1165+
1166+ releaseRegularTurn ?.( ) ;
1167+ await vi . advanceTimersByTimeAsync ( 1_000 ) ;
1168+ await vi . waitFor ( async ( ) =>
1169+ expect (
1170+ ( await listTelegramSpooledUpdates ( { spoolDir : tempDir } ) ) . map (
1171+ ( update ) => update . updateId ,
1172+ ) ,
1173+ ) . toEqual ( [ ] ) ,
1174+ ) ;
1175+ workerListeners [ 0 ] ?.( {
1176+ type : "poll-success" ,
1177+ offset : null ,
1178+ count : 0 ,
1179+ finishedAt : Date . now ( ) ,
1180+ } ) ;
1181+ await vi . waitFor ( ( ) => expect ( statusPatches ( setStatus ) . at ( - 1 ) ?. connected ) . toBe ( true ) ) ;
1182+ expect ( createWorker ) . toHaveBeenCalledTimes ( 1 ) ;
1183+
1184+ abort . abort ( ) ;
1185+ await vi . advanceTimersByTimeAsync ( 20_000 ) ;
1186+ await runPromise ;
1187+ } finally {
1188+ releaseRegularTurn ?.( ) ;
1189+ vi . useRealTimers ( ) ;
1190+ await fs . rm ( tempDir , { recursive : true , force : true } ) ;
1191+ }
1192+ } ) ;
1193+
8851194 it ( "forces a restart when polling stalls without getUpdates activity" , async ( ) => {
8861195 const abort = new AbortController ( ) ;
8871196 const botStop = vi . fn ( async ( ) => undefined ) ;
0 commit comments