@@ -104,6 +104,7 @@ type WorkerPollSuccessListener = (message: {
104104type WorkerPollErrorListener = ( message : {
105105 type : "poll-error" ;
106106 message : string ;
107+ errorCode ?: number ;
107108 finishedAt : number ;
108109} ) => void ;
109110type WorkerMessageListener = ( message : TelegramIngressWorkerMessage ) => void ;
@@ -2356,6 +2357,89 @@ describe("TelegramPollingSession", () => {
23562357 }
23572358 } ) ;
23582359
2360+ it ( "restarts isolated ingress on a getUpdates conflict instead of crashing the account" , async ( ) => {
2361+ const abort = new AbortController ( ) ;
2362+ const tempDir = await fs . mkdtemp ( path . join ( os . tmpdir ( ) , "openclaw-telegram-spool-" ) ) ;
2363+ const log = vi . fn ( ) ;
2364+ const setStatus = vi . fn ( ) ;
2365+ // 409 conflicts are not "recoverable network errors"; the conflict branch
2366+ // must restart the cycle before that classifier is consulted.
2367+ isRecoverableTelegramNetworkErrorMock . mockReturnValue ( false ) ;
2368+ const deleteWebhook = vi . fn ( async ( ) => true ) ;
2369+ createTelegramBotMock . mockImplementation ( ( ) => ( {
2370+ api : {
2371+ deleteWebhook,
2372+ config : { use : vi . fn ( ) } ,
2373+ } ,
2374+ init : vi . fn ( async ( ) => undefined ) ,
2375+ handleUpdate : vi . fn ( async ( ) => undefined ) ,
2376+ stop : vi . fn ( async ( ) => undefined ) ,
2377+ } ) ) ;
2378+ const transport1 = makeTelegramTransport ( ) ;
2379+ const transport2 = makeTelegramTransport ( ) ;
2380+ const createTelegramTransport = vi
2381+ . fn < ( ) => ReturnType < typeof makeTelegramTransport > > ( )
2382+ . mockReturnValueOnce ( transport2 ) ;
2383+
2384+ let workerCycle = 0 ;
2385+ let listener : WorkerPollErrorListener | undefined ;
2386+ const createWorker = vi . fn ( ( ) => ( {
2387+ onMessage : vi . fn ( ( next : WorkerPollErrorListener ) => {
2388+ listener = next ;
2389+ return ( ) => undefined ;
2390+ } ) ,
2391+ stop : vi . fn ( async ( ) => undefined ) ,
2392+ task : vi . fn ( async ( ) => {
2393+ workerCycle += 1 ;
2394+ if ( workerCycle === 1 ) {
2395+ listener ?.( {
2396+ type : "poll-error" ,
2397+ message : "Conflict: terminated by other getUpdates request" ,
2398+ errorCode : 409 ,
2399+ finishedAt : Date . now ( ) ,
2400+ } ) ;
2401+ throw new Error ( "Telegram ingress worker exited with code 1" ) ;
2402+ }
2403+ abort . abort ( ) ;
2404+ } ) ,
2405+ } ) ) ;
2406+
2407+ try {
2408+ const session = createPollingSession ( {
2409+ abortSignal : abort . signal ,
2410+ log,
2411+ setStatus,
2412+ telegramTransport : transport1 ,
2413+ createTelegramTransport,
2414+ isolatedIngress : {
2415+ enabled : true ,
2416+ spoolDir : tempDir ,
2417+ createWorker,
2418+ drainIntervalMs : 100 ,
2419+ } ,
2420+ } ) ;
2421+
2422+ await session . runUntilAbort ( ) ;
2423+
2424+ expect ( createWorker ) . toHaveBeenCalledTimes ( 2 ) ;
2425+ // The conflict resets webhook cleanup so the next cycle re-runs deleteWebhook.
2426+ expect ( deleteWebhook ) . toHaveBeenCalledTimes ( 2 ) ;
2427+ // The conflict marks the transport dirty so the next cycle gets a fresh socket.
2428+ expect ( createTelegramTransport ) . toHaveBeenCalledTimes ( 1 ) ;
2429+ expectLogIncludes ( log , "Another OpenClaw gateway, script, or Telegram poller" ) ;
2430+ expect (
2431+ statusPatches ( setStatus ) . some (
2432+ ( patch ) =>
2433+ patch . connected === false &&
2434+ String ( patch . lastError ) . includes ( "Another OpenClaw gateway" ) ,
2435+ ) ,
2436+ ) . toBe ( true ) ;
2437+ } finally {
2438+ abort . abort ( ) ;
2439+ await fs . rm ( tempDir , { recursive : true , force : true } ) ;
2440+ }
2441+ } ) ;
2442+
23592443 it ( "keeps active spooled lanes blocked across account restarts" , async ( ) => {
23602444 vi . useFakeTimers ( { shouldAdvanceTime : true } ) ;
23612445 const firstAbort = new AbortController ( ) ;
@@ -3533,9 +3617,11 @@ describe("TelegramPollingSession", () => {
35333617 const watchdogHarness = installPollingStallWatchdogHarness ( ) ;
35343618
35353619 const log = vi . fn ( ) ;
3620+ const setStatus = vi . fn ( ) ;
35363621 const session = createPollingSession ( {
35373622 abortSignal : abort . signal ,
35383623 log,
3624+ setStatus,
35393625 } ) ;
35403626
35413627 try {
@@ -3567,6 +3653,14 @@ describe("TelegramPollingSession", () => {
35673653 abort . abort ( ) ;
35683654 resolveFirstTask ( ) ;
35693655 await runPromise ;
3656+
3657+ // The stall must reach channel status, not just the gateway log.
3658+ expect (
3659+ statusPatches ( setStatus ) . some (
3660+ ( patch ) =>
3661+ patch . connected === false && String ( patch . lastError ) . includes ( "Polling stall detected" ) ,
3662+ ) ,
3663+ ) . toBe ( true ) ;
35703664 } finally {
35713665 watchdogHarness . restore ( ) ;
35723666 }
@@ -3614,6 +3708,7 @@ describe("TelegramPollingSession", () => {
36143708 it ( "logs an actionable duplicate-poller hint for getUpdates conflicts" , async ( ) => {
36153709 const abort = new AbortController ( ) ;
36163710 const log = vi . fn ( ) ;
3711+ const setStatus = vi . fn ( ) ;
36173712 const conflictError = Object . assign (
36183713 new Error ( "Conflict: terminated by other getUpdates request" ) ,
36193714 {
@@ -3628,11 +3723,19 @@ describe("TelegramPollingSession", () => {
36283723 const session = createPollingSession ( {
36293724 abortSignal : abort . signal ,
36303725 log,
3726+ setStatus,
36313727 } ) ;
36323728
36333729 await session . runUntilAbort ( ) ;
36343730
36353731 expectLogIncludes ( log , "Another OpenClaw gateway, script, or Telegram poller" ) ;
3732+ // The hint must reach channel status, not just the gateway log.
3733+ expect (
3734+ statusPatches ( setStatus ) . some (
3735+ ( patch ) =>
3736+ patch . connected === false && String ( patch . lastError ) . includes ( "Another OpenClaw gateway" ) ,
3737+ ) ,
3738+ ) . toBe ( true ) ;
36363739 } ) ;
36373740
36383741 it ( "logs polling cycle start after a transport rebuild" , async ( ) => {
0 commit comments