From a1097d21471277824cb8cd544ac5fb8503ab2bad Mon Sep 17 00:00:00 2001 From: yiliang114 Date: Fri, 31 Jul 2026 16:37:06 +0800 Subject: [PATCH 1/9] fix(channel): recover ACP bridge after wake --- .../channels/base/src/ChannelBase.test.ts | 199 ++++++++++++++++- packages/channels/base/src/ChannelBase.ts | 21 ++ .../cli/src/commands/channel/start.test.ts | 200 ++++++++++++++++- packages/cli/src/commands/channel/start.ts | 208 +++++++++++------- .../core/src/telemetry/event-loop-lag.test.ts | 169 ++++++++++++++ packages/core/src/telemetry/event-loop-lag.ts | 84 ++++++- 6 files changed, 785 insertions(+), 96 deletions(-) diff --git a/packages/channels/base/src/ChannelBase.test.ts b/packages/channels/base/src/ChannelBase.test.ts index ef9d4be4358..c613b74d23f 100644 --- a/packages/channels/base/src/ChannelBase.test.ts +++ b/packages/channels/base/src/ChannelBase.test.ts @@ -9318,7 +9318,9 @@ describe('ChannelBase', () => { const firstCancel = ch.handleInbound(envelope({ text: '/cancel' })); const secondCancel = ch.handleInbound(envelope({ text: '/cancel' })); - expect(bridge.cancelSession).toHaveBeenCalledTimes(1); + await vi.waitFor(() => + expect(bridge.cancelSession).toHaveBeenCalledTimes(1), + ); resolveCancel(); await Promise.all([firstCancel, secondCancel]); resolvePrompt('late response'); @@ -15910,7 +15912,7 @@ describe('ChannelBase', () => { title: 'CI failed again', }); secondRun.catch(() => undefined); - await Promise.resolve(); + await new Promise((resolve) => setImmediate(resolve)); ( ch as unknown as { sessionGenerations: Map; @@ -16012,6 +16014,197 @@ describe('ChannelBase', () => { .calls[1][1] as string; expect(collectedPrompt).toContain('follow-up while webhook runs'); }); + + it('rechecks bridge recovery before a webhook prompt starts', async () => { + const recoveryState: { current?: Promise } = {}; + let releaseRecovery: (() => void) | undefined; + let releaseMemoryRead: (() => void) | undefined; + const memoryRead = new Promise((resolve) => { + releaseMemoryRead = () => resolve(''); + }); + const channelMemory = createChannelMemory(); + channelMemory.readChannelMemory.mockReturnValue(memoryRead); + const ch = createChannel( + { approvalMode: 'yolo', webhooks }, + { bridgeRecovery: () => recoveryState.current, channelMemory }, + ); + ch.proactiveSupported = true; + + const run = ch.runWebhookTask(webhookTask); + await vi.waitFor(() => + expect(channelMemory.readChannelMemory).toHaveBeenCalledTimes(1), + ); + recoveryState.current = new Promise((resolve) => { + releaseRecovery = resolve; + }); + releaseMemoryRead!(); + await Promise.resolve(); + + expect(bridge.prompt).not.toHaveBeenCalled(); + + releaseRecovery!(); + await expect(run).resolves.toBe('agent response'); + + expect(bridge.prompt).toHaveBeenCalledTimes(1); + }); + }); + + it('waits for bridge recovery before resolving an inbound session', async () => { + let releaseBridge: (() => void) | undefined; + const bridgeReady = new Promise((resolve) => { + releaseBridge = resolve; + }); + const ch = createChannel({}, { bridgeRecovery: () => bridgeReady }); + + const inbound = ch.handleInbound(envelope()); + await Promise.resolve(); + expect(bridge.newSession).not.toHaveBeenCalled(); + expect(bridge.prompt).not.toHaveBeenCalled(); + + releaseBridge!(); + await inbound; + + expect(bridge.prompt).toHaveBeenCalledTimes(1); + }); + + it('waits for bridge recovery after adapter-specific preflight', async () => { + let releaseBridge: (() => void) | undefined; + const bridgeReady = new Promise((resolve) => { + releaseBridge = resolve; + }); + const ch = createChannel({}, { bridgeRecovery: () => bridgeReady }); + + const inbound = ch.processAfterAdapterPreflight(envelope()); + await Promise.resolve(); + expect(bridge.newSession).not.toHaveBeenCalled(); + expect(bridge.prompt).not.toHaveBeenCalled(); + + releaseBridge!(); + await inbound; + + expect(bridge.prompt).toHaveBeenCalledTimes(1); + }); + + it('rechecks bridge recovery after inbound preprocessing has started', async () => { + const recoveryState: { current?: Promise } = {}; + let releaseRecovery: (() => void) | undefined; + let releaseContactRecord: (() => void) | undefined; + const contactRecord = new Promise((resolve) => { + releaseContactRecord = resolve; + }); + const ch = createChannel( + {}, + { + bridgeRecovery: () => recoveryState.current, + observedContacts: { + record: vi.fn().mockImplementation(() => contactRecord), + }, + }, + ); + + const inbound = ch.handleInbound(envelope()); + await vi.waitFor(() => expect(bridge.newSession).not.toHaveBeenCalled()); + recoveryState.current = new Promise((resolve) => { + releaseRecovery = resolve; + }); + releaseContactRecord!(); + await Promise.resolve(); + + expect(bridge.newSession).not.toHaveBeenCalled(); + expect(bridge.prompt).not.toHaveBeenCalled(); + + releaseRecovery!(); + await inbound; + + expect(bridge.prompt).toHaveBeenCalledTimes(1); + }); + + it('waits for bridge recovery before running a loop prompt', async () => { + let releaseBridge: (() => void) | undefined; + const bridgeReady = new Promise((resolve) => { + releaseBridge = resolve; + }); + const ch = createChannel({}, { bridgeRecovery: () => bridgeReady }); + ch.proactiveSupported = true; + + const loopRun = ch.runLoopPrompt({ + id: 'job-recovery-gate', + channelName: 'test-chan', + target: { + channelName: 'test-chan', + senderId: 'alice', + chatId: 'chat-1', + isGroup: false, + }, + cwd: '/tmp', + cron: '0 9 * * *', + prompt: 'post summary', + recurring: true, + enabled: true, + createdBy: 'Alice', + createdAt: '2026-06-30T01:00:00.000Z', + consecutiveFailures: 0, + runCount: 0, + }); + await Promise.resolve(); + expect(bridge.newSession).not.toHaveBeenCalled(); + expect(bridge.prompt).not.toHaveBeenCalled(); + + releaseBridge!(); + await expect(loopRun).resolves.toBe('agent response'); + + expect(bridge.prompt).toHaveBeenCalledTimes(1); + }); + + it('rechecks bridge recovery before a loop prompt starts', async () => { + const recoveryState: { current?: Promise } = {}; + let releaseRecovery: (() => void) | undefined; + let releaseMemoryRead: (() => void) | undefined; + const memoryRead = new Promise((resolve) => { + releaseMemoryRead = () => resolve(''); + }); + const channelMemory = createChannelMemory(); + channelMemory.readChannelMemory.mockReturnValue(memoryRead); + const ch = createChannel( + {}, + { bridgeRecovery: () => recoveryState.current, channelMemory }, + ); + ch.proactiveSupported = true; + + const loopRun = ch.runLoopPrompt({ + id: 'job-late-recovery-gate', + channelName: 'test-chan', + target: { + channelName: 'test-chan', + senderId: 'alice', + chatId: 'chat-1', + isGroup: false, + }, + cwd: '/tmp', + cron: '0 9 * * *', + prompt: 'post summary', + recurring: true, + enabled: true, + createdBy: 'Alice', + createdAt: '2026-06-30T01:00:00.000Z', + consecutiveFailures: 0, + runCount: 0, + }); + await vi.waitFor(() => + expect(channelMemory.readChannelMemory).toHaveBeenCalledTimes(1), + ); + recoveryState.current = new Promise((resolve) => { + releaseRecovery = resolve; + }); + releaseMemoryRead!(); + await Promise.resolve(); + + expect(bridge.prompt).not.toHaveBeenCalled(); + + releaseRecovery!(); + await expect(loopRun).resolves.toBe('agent response'); + + expect(bridge.prompt).toHaveBeenCalledTimes(1); }); it('runs a loop prompt as a follow-up and pushes the result proactively', async () => { @@ -17253,7 +17446,7 @@ describe('ChannelBase', () => { consecutiveFailures: 0, runCount: 0, }); - await Promise.resolve(); + await new Promise((resolve) => setImmediate(resolve)); expect(bridge.prompt).toHaveBeenCalledOnce(); ( ch as unknown as { sessionGenerations: Map } diff --git a/packages/channels/base/src/ChannelBase.ts b/packages/channels/base/src/ChannelBase.ts index e11eebcab1f..651c9d62ebd 100644 --- a/packages/channels/base/src/ChannelBase.ts +++ b/packages/channels/base/src/ChannelBase.ts @@ -218,6 +218,8 @@ export interface ChannelBaseOptions { * events directly. */ registerBridgeEvents?: boolean; + /** Return the active bridge recovery barrier, if recovery is in progress. */ + bridgeRecovery?: () => Promise | undefined; groupHistoryPath?: string; loopController?: ChannelLoopController; observedContacts?: { @@ -379,6 +381,7 @@ export abstract class ChannelBase { /** Per-session promise chain to serialize prompt + send (followup mode). */ private sessionQueues: Map> = new Map(); private readonly registerBridgeEvents: boolean; + private readonly bridgeRecovery?: () => Promise | undefined; /** * Per-session generation, bumped by /clear. A queued followup turn captures the * generation when it enqueues and bails if /clear bumped it before the turn ran, @@ -804,6 +807,7 @@ export abstract class ChannelBase { ); this.loopController = options?.loopController; this.observedContacts = options?.observedContacts; + this.bridgeRecovery = options?.bridgeRecovery; this.groupGate = new GroupGate(config.groupPolicy, config.groups); this.dmGate = new DmGate(config.dmPolicy); @@ -1461,6 +1465,7 @@ export abstract class ChannelBase { throw new Error(`Loop ${job.id} target is no longer authorized.`); } + await this.waitForBridgeRecovery(); const sessionId = await this.router.resolve( this.name, job.target.senderId, @@ -1611,6 +1616,7 @@ export abstract class ChannelBase { heldChunks.length = 0; this.onResponseBoundary(job.target.chatId, sessionId); }; + await this.waitForBridgeRecovery(); const promptBridge = this.bridge; promptBridge.on('textChunk', onChunk); promptBridge.on('responseBoundary', onResponseBoundary); @@ -1778,6 +1784,7 @@ export abstract class ChannelBase { ): Promise { const target = this.resolveWebhookTaskTarget(task); + await this.waitForBridgeRecovery(); const sessionId = await this.router.resolve( this.name, target.senderId, @@ -1906,6 +1913,7 @@ export abstract class ChannelBase { releaseHeldChunks(); } }; + await this.waitForBridgeRecovery(); const promptBridge = this.bridge; promptBridge.on('textChunk', onChunk); @@ -4904,6 +4912,12 @@ export abstract class ChannelBase { this.preflightedEnvelopes.add(envelope); } + /** Wait until the currently active bridge recovery, if any, has completed. */ + private async waitForBridgeRecovery(): Promise { + const bridgeRecovery = this.bridgeRecovery?.(); + if (bridgeRecovery) await bridgeRecovery; + } + /** * Process an inbound message after preflight gates have passed. * @@ -4912,6 +4926,7 @@ export abstract class ChannelBase { * already preflighted, such as during collect-buffer drain. */ protected async processInbound(envelope: Envelope): Promise { + await this.waitForBridgeRecovery(); if (!this.preflightedEnvelopes.delete(envelope)) { throw new Error( 'processInbound called without a successful preflightInbound check.', @@ -5002,6 +5017,9 @@ export abstract class ChannelBase { } } + // Preprocessing above can await memory/command hooks; recovery may have + // started since the entry check. Recheck immediately before session routing. + await this.waitForBridgeRecovery(); const sessionId = await this.router.resolve( this.name, envelope.senderId, @@ -5487,6 +5505,9 @@ export abstract class ChannelBase { ); streamer?.stop(); }; + // Queue wait and memory recall can outlive a bridge crash. Capture the + // bridge only after the latest recovery has restored session routing. + await this.waitForBridgeRecovery(); const promptBridge = this.bridge; promptBridge.on('textChunk', onChunk); promptBridge.on('responseBoundary', onResponseBoundary); diff --git a/packages/cli/src/commands/channel/start.test.ts b/packages/cli/src/commands/channel/start.test.ts index 3673184baa5..0929ba093a8 100644 --- a/packages/cli/src/commands/channel/start.test.ts +++ b/packages/cli/src/commands/channel/start.test.ts @@ -775,7 +775,7 @@ describe('startCommand.handler', () => { const restartedBridge = mockAcpBridge.mock.results[1]!.value; expect(mockRouterSetBridge).toHaveBeenCalledWith(restartedBridge); expect(mockChannelSetBridge).toHaveBeenCalledWith(restartedBridge); - expect(mockChannelConnect).toHaveBeenCalledTimes(2); + expect(mockChannelConnect).toHaveBeenCalledTimes(1); const sessionDiedCalls = mockBridgeOn.mock.calls.filter( ([eventName]) => eventName === 'sessionDied', @@ -787,7 +787,7 @@ describe('startCommand.handler', () => { expect(mockBridgeOn.mock.invocationCallOrder.at(-2)).toBeLessThan( mockRouterRestoreSessions.mock.invocationCallOrder[0]!, ); - expect(mockChannelConnect.mock.invocationCallOrder[1]).toBeLessThan( + expect(mockChannelConnect.mock.invocationCallOrder[0]).toBeLessThan( mockRouterRestoreSessions.mock.invocationCallOrder[0]!, ); @@ -802,6 +802,145 @@ describe('startCommand.handler', () => { } }); + it('recovers a standalone bridge without reconnecting the channel adapter', async () => { + mockChannelConnect.mockResolvedValue(undefined); + const channels = { telegram: { type: 'telegram' } }; + mockLoadSettings.mockReturnValue({ merged: { channels } }); + const processOnSpy = vi + .spyOn(process, 'on') + .mockImplementation(() => process); + + try { + void invokeStartHandler({ name: 'telegram' }); + await new Promise((resolve) => setImmediate(resolve)); + + const disconnectedListener = mockBridgeOn.mock.calls.find( + ([eventName]) => eventName === 'disconnected', + )?.[1] as (() => Promise) | undefined; + expect(disconnectedListener).toBeDefined(); + + vi.useFakeTimers(); + const restart = disconnectedListener!(); + await vi.advanceTimersByTimeAsync(3000); + await restart; + + expect(mockChannelConnect).toHaveBeenCalledTimes(1); + expect(mockChannelDisconnect).not.toHaveBeenCalled(); + expect(mockChannelLoopSchedulerStop).not.toHaveBeenCalled(); + } finally { + processOnSpy.mockRestore(); + vi.useRealTimers(); + } + }); + + it('coalesces duplicate standalone disconnect events into one recovery', async () => { + mockChannelConnect.mockResolvedValue(undefined); + const channels = { telegram: { type: 'telegram' } }; + mockLoadSettings.mockReturnValue({ merged: { channels } }); + const processOnSpy = vi + .spyOn(process, 'on') + .mockImplementation(() => process); + + try { + void invokeStartHandler({ name: 'telegram' }); + await new Promise((resolve) => setImmediate(resolve)); + + const disconnectedListener = mockBridgeOn.mock.calls.find( + ([eventName]) => eventName === 'disconnected', + )?.[1] as (() => Promise) | undefined; + expect(disconnectedListener).toBeDefined(); + + vi.useFakeTimers(); + const firstRestart = disconnectedListener!(); + const secondRestart = disconnectedListener!(); + await vi.advanceTimersByTimeAsync(3000); + await Promise.all([firstRestart, secondRestart]); + + expect(mockAcpBridge).toHaveBeenCalledTimes(2); + expect(mockRouterRestoreSessions).toHaveBeenCalledTimes(1); + } finally { + processOnSpy.mockRestore(); + vi.useRealTimers(); + } + }); + + it('keeps the readiness gate blocked and coalesces replacement disconnects', async () => { + mockChannelConnect.mockResolvedValue(undefined); + let resolveFirstRestore: + | ((value: { failed: number; restored: number }) => void) + | undefined; + let resolveSecondRestore: + | ((value: { failed: number; restored: number }) => void) + | undefined; + mockRouterRestoreSessions + .mockImplementationOnce( + () => + new Promise((resolve) => { + resolveFirstRestore = resolve; + }), + ) + .mockImplementationOnce( + () => + new Promise((resolve) => { + resolveSecondRestore = resolve; + }), + ); + const channels = { telegram: { type: 'telegram' } }; + mockLoadSettings.mockReturnValue({ merged: { channels } }); + const processOnSpy = vi + .spyOn(process, 'on') + .mockImplementation(() => process); + + try { + void invokeStartHandler({ name: 'telegram' }); + await new Promise((resolve) => setImmediate(resolve)); + const options = mockCreateChannel.mock.calls[0]?.[3] as + | ChannelBaseOptions + | undefined; + const firstDisconnect = mockBridgeOn.mock.calls.find( + ([eventName]) => eventName === 'disconnected', + )?.[1] as (() => void) | undefined; + expect(firstDisconnect).toBeDefined(); + + vi.useFakeTimers(); + firstDisconnect!(); + const recoveryGate = options?.bridgeRecovery?.(); + let gateReleased = false; + void recoveryGate?.then(() => { + gateReleased = true; + }); + await vi.advanceTimersByTimeAsync(3000); + await vi.waitFor(() => + expect(mockRouterRestoreSessions).toHaveBeenCalledTimes(1), + ); + + const disconnectListeners = mockBridgeOn.mock.calls.filter( + ([eventName]) => eventName === 'disconnected', + ); + expect(disconnectListeners).toHaveLength(2); + const replacementDisconnect = disconnectListeners[1]![1] as () => void; + replacementDisconnect(); + replacementDisconnect(); + + resolveFirstRestore!({ failed: 0, restored: 0 }); + await vi.advanceTimersByTimeAsync(3000); + await vi.waitFor(() => + expect(mockRouterRestoreSessions).toHaveBeenCalledTimes(2), + ); + + expect(gateReleased).toBe(false); + expect(mockAcpBridge).toHaveBeenCalledTimes(3); + + resolveSecondRestore!({ failed: 0, restored: 0 }); + await vi.waitFor(() => expect(gateReleased).toBe(true)); + + expect(mockRouterRestoreSessions).toHaveBeenCalledTimes(2); + } finally { + processOnSpy.mockRestore(); + vi.useRealTimers(); + } + }); + it('starts all channels with one shared bridge and router', async () => { const channels = { first: { type: 'telegram' }, @@ -1064,4 +1203,61 @@ describe('startCommand.handler', () => { vi.useRealTimers(); } }); + it('recovers a shared bridge without reconnecting channel adapters', async () => { + const channels = { + first: { type: 'telegram' }, + second: { type: 'telegram' }, + }; + const firstChannel = { + connect: vi.fn().mockResolvedValue(undefined), + disconnect: vi.fn(), + onSessionDied: vi.fn(), + onToolCall: vi.fn(), + setBridge: vi.fn(), + }; + const secondChannel = { + connect: vi.fn().mockResolvedValue(undefined), + disconnect: vi.fn(), + onSessionDied: vi.fn(), + onToolCall: vi.fn(), + setBridge: vi.fn(), + }; + mockLoadSettings.mockReturnValue({ merged: { channels } }); + mockParseChannelConfig.mockImplementation(async (name: string) => ({ + ...mockParsedChannelConfig, + cwd: `/tmp/${name}`, + model: 'shared-model', + sessionScope: 'user', + })); + mockCreateChannel + .mockReturnValueOnce(firstChannel) + .mockReturnValueOnce(secondChannel); + const processOnSpy = vi + .spyOn(process, 'on') + .mockImplementation(() => process); + + try { + void invokeStartHandler({}); + await new Promise((resolve) => setImmediate(resolve)); + + const disconnectedListener = mockBridgeOn.mock.calls.find( + ([eventName]) => eventName === 'disconnected', + )?.[1] as (() => Promise) | undefined; + expect(disconnectedListener).toBeDefined(); + + vi.useFakeTimers(); + const restart = disconnectedListener!(); + await vi.advanceTimersByTimeAsync(3000); + await restart; + + expect(firstChannel.connect).toHaveBeenCalledTimes(1); + expect(secondChannel.connect).toHaveBeenCalledTimes(1); + expect(firstChannel.disconnect).not.toHaveBeenCalled(); + expect(secondChannel.disconnect).not.toHaveBeenCalled(); + expect(mockChannelLoopSchedulerStop).not.toHaveBeenCalled(); + } finally { + processOnSpy.mockRestore(); + vi.useRealTimers(); + } + }); }); diff --git a/packages/cli/src/commands/channel/start.ts b/packages/cli/src/commands/channel/start.ts index c733f2643ac..cdbcde55332 100644 --- a/packages/cli/src/commands/channel/start.ts +++ b/packages/cli/src/commands/channel/start.ts @@ -124,6 +124,30 @@ function cleanupStartedChannels( } } +function createBridgeReadinessGate(): { + current: () => Promise | undefined; + block: () => void; + release: () => void; +} { + let pending: Promise | undefined; + let releasePending: (() => void) | undefined; + return { + current: () => pending, + block: () => { + if (pending) return; + pending = new Promise((resolve) => { + releasePending = resolve; + }); + }, + release: () => { + const release = releasePending; + pending = undefined; + releasePending = undefined; + release?.(); + }, + }; +} + /** Check for duplicate instance and abort if one is already running. */ function checkDuplicateInstance(): void { const existing = readServiceInfo(); @@ -179,6 +203,9 @@ async function startSingle( const cliEntryPath = findCliEntryPath(); let shuttingDown = false; const crashTimestamps: number[] = []; + let recoveryTask: Promise | undefined; + let recoveryRequested = false; + const bridgeReadiness = createBridgeReadinessGate(); const bridgeOpts = { cliEntryPath, cwd: config.cwd, model: config.model }; let bridge = new AcpBridge(bridgeOpts); @@ -203,6 +230,7 @@ async function startSingle( proxy, ...channelMemoryOptions(() => bridge, config.cwd), ...(loopController ? { loopController } : {}), + bridgeRecovery: bridgeReadiness.current, }); channels.set(name, channel); const scheduler = loopStore @@ -232,62 +260,78 @@ async function startSingle( scheduler?.start(); writeStdoutLine(`[Channel] "${name}" is running. Press Ctrl+C to stop.`); - const attachDisconnectHandler = (b: AcpBridge): void => { - b.on('disconnected', async () => { + const attachDisconnectHandler = (failedBridge: AcpBridge): void => { + failedBridge.on('disconnected', () => { if (shuttingDown) return; + if (recoveryTask) { + if (failedBridge === bridge) recoveryRequested = true; + return; + } + recoverBridge(); + }); + }; - const now = Date.now(); - crashTimestamps.push(now); - // Only count crashes within the recent window - const recentCrashes = crashTimestamps.filter( - (ts) => now - ts < CRASH_WINDOW_MS, - ); + const recoverBridge = (): void => { + bridgeReadiness.block(); + const task = (async () => { + do { + recoveryRequested = false; + const now = Date.now(); + crashTimestamps.push(now); + // Only count crashes within the recent window + const recentCrashes = crashTimestamps.filter( + (timestamp) => now - timestamp < CRASH_WINDOW_MS, + ); + + if (recentCrashes.length > MAX_CRASH_RESTARTS) { + writeStderrLine( + `[Channel] Bridge crashed ${recentCrashes.length} times in ${CRASH_WINDOW_MS / 1000}s. Giving up.`, + ); + scheduler?.stop(); + channel.disconnect(); + router.clearAll(); + removeServiceInfo(); + process.exit(1); + } - if (recentCrashes.length > MAX_CRASH_RESTARTS) { writeStderrLine( - `[Channel] Bridge crashed ${recentCrashes.length} times in ${CRASH_WINDOW_MS / 1000}s. Giving up.`, + `[Channel] Bridge crashed (${recentCrashes.length}/${MAX_CRASH_RESTARTS} in window). Restarting in ${RESTART_DELAY_MS / 1000}s...`, ); - scheduler?.stop(); - channel.disconnect(); - router.clearAll(); - removeServiceInfo(); - process.exit(1); - } + await new Promise((resolve) => setTimeout(resolve, RESTART_DELAY_MS)); - writeStderrLine( - `[Channel] Bridge crashed (${recentCrashes.length}/${MAX_CRASH_RESTARTS} in window). Restarting in ${RESTART_DELAY_MS / 1000}s...`, - ); - scheduler?.stop(); - await new Promise((r) => setTimeout(r, RESTART_DELAY_MS)); - - try { bridge = new AcpBridge(bridgeOpts); + attachDisconnectHandler(bridge); await bridge.start(); router.setBridge(bridge); channel.setBridge(bridge); - channel.disconnect(); - await channel.connect(); registerToolCallDispatch(bridge, router, channels); registerBackgroundResponseRelay(bridge, router, channels); registerPermissionRelay(bridge, router, channels); registerSessionCleanup(bridge, router, channels); - attachDisconnectHandler(bridge); const result = await router.restoreSessions(); - scheduler?.start(); writeStdoutLine( `[Channel] Bridge restarted. Sessions restored: ${result.restored}, failed: ${result.failed}`, ); - } catch (err) { + } while (recoveryRequested && !shuttingDown); + })() + .catch((err) => { writeStderrLine( `[Channel] Failed to restart bridge: ${err instanceof Error ? err.message : String(err)}`, ); + scheduler?.stop(); channel.disconnect(); router.clearAll(); removeServiceInfo(); process.exit(1); - } - }); + }) + .finally(() => { + if (recoveryTask === task) { + recoveryTask = undefined; + bridgeReadiness.release(); + } + }); + recoveryTask = task; }; attachDisconnectHandler(bridge); @@ -340,6 +384,9 @@ async function startAll( const defaultCwd = process.cwd(); let shuttingDown = false; const crashTimestamps: number[] = []; + let recoveryTask: Promise | undefined; + let recoveryRequested = false; + const bridgeReadiness = createBridgeReadinessGate(); const bridgeOpts = { cliEntryPath, @@ -374,6 +421,7 @@ async function startAll( proxy, ...channelMemoryOptions(() => bridge, config.cwd), ...(loopController ? { loopController } : {}), + bridgeRecovery: bridgeReadiness.current, }), ); } @@ -419,79 +467,73 @@ async function startAll( `[Channel] Running ${connectedCount} channel(s). Press Ctrl+C to stop.`, ); - const attachDisconnectHandler = (b: AcpBridge): void => { - b.on('disconnected', async () => { + const attachDisconnectHandler = (failedBridge: AcpBridge): void => { + failedBridge.on('disconnected', () => { if (shuttingDown) return; + if (recoveryTask) { + if (failedBridge === bridge) recoveryRequested = true; + return; + } + recoverBridge(); + }); + }; - const now = Date.now(); - crashTimestamps.push(now); - const recentCrashes = crashTimestamps.filter( - (ts) => now - ts < CRASH_WINDOW_MS, - ); - - if (recentCrashes.length > MAX_CRASH_RESTARTS) { - writeStderrLine( - `[Channel] Bridge crashed ${recentCrashes.length} times in ${CRASH_WINDOW_MS / 1000}s. Giving up.`, + const recoverBridge = (): void => { + bridgeReadiness.block(); + const task = (async () => { + do { + recoveryRequested = false; + const now = Date.now(); + crashTimestamps.push(now); + const recentCrashes = crashTimestamps.filter( + (timestamp) => now - timestamp < CRASH_WINDOW_MS, ); - scheduler?.stop(); - for (const channel of channels.values()) { - try { - channel.disconnect(); - } catch { - // best-effort + + if (recentCrashes.length > MAX_CRASH_RESTARTS) { + writeStderrLine( + `[Channel] Bridge crashed ${recentCrashes.length} times in ${CRASH_WINDOW_MS / 1000}s. Giving up.`, + ); + scheduler?.stop(); + for (const channel of channels.values()) { + try { + channel.disconnect(); + } catch { + // best-effort + } } + router.clearAll(); + removeServiceInfo(); + process.exit(1); } - router.clearAll(); - removeServiceInfo(); - process.exit(1); - } - writeStderrLine( - `[Channel] Bridge crashed (${recentCrashes.length}/${MAX_CRASH_RESTARTS} in window). Restarting in ${RESTART_DELAY_MS / 1000}s...`, - ); - scheduler?.stop(); - await new Promise((r) => setTimeout(r, RESTART_DELAY_MS)); + writeStderrLine( + `[Channel] Bridge crashed (${recentCrashes.length}/${MAX_CRASH_RESTARTS} in window). Restarting in ${RESTART_DELAY_MS / 1000}s...`, + ); + await new Promise((resolve) => setTimeout(resolve, RESTART_DELAY_MS)); - try { bridge = new AcpBridge(bridgeOpts); + attachDisconnectHandler(bridge); await bridge.start(); router.setBridge(bridge); for (const channel of channels.values()) { channel.setBridge(bridge); } - for (const [name, channel] of connectedChannels) { - try { - channel.disconnect(); - await channel.connect(); - } catch (err) { - writeStderrLine( - `[Channel] "${name}" failed to reconnect: ${err instanceof Error ? err.message : String(err)}`, - ); - connectedChannels.delete(name); - } - } - if (connectedChannels.size === 0) { - writeStderrLine('[Channel] No channels reconnected. Exiting.'); - bridge.stop(); - router.clearAll(); - removeServiceInfo(); - process.exit(1); - } registerToolCallDispatch(bridge, router, channels); registerBackgroundResponseRelay(bridge, router, channels); registerPermissionRelay(bridge, router, channels); registerSessionCleanup(bridge, router, channels); - attachDisconnectHandler(bridge); const result = await router.restoreSessions(); - scheduler?.start(); writeStdoutLine( `[Channel] Bridge restarted. Sessions restored: ${result.restored}, failed: ${result.failed}`, ); - } catch (err) { + } while (recoveryRequested && !shuttingDown); + })() + .catch((err) => { writeStderrLine( `[Channel] Failed to restart bridge: ${err instanceof Error ? err.message : String(err)}`, ); + scheduler?.stop(); for (const channel of channels.values()) { try { channel.disconnect(); @@ -502,8 +544,14 @@ async function startAll( router.clearAll(); removeServiceInfo(); process.exit(1); - } - }); + }) + .finally(() => { + if (recoveryTask === task) { + recoveryTask = undefined; + bridgeReadiness.release(); + } + }); + recoveryTask = task; }; attachDisconnectHandler(bridge); diff --git a/packages/core/src/telemetry/event-loop-lag.test.ts b/packages/core/src/telemetry/event-loop-lag.test.ts index 5e04f2d0bc1..5fc5123dcea 100644 --- a/packages/core/src/telemetry/event-loop-lag.test.ts +++ b/packages/core/src/telemetry/event-loop-lag.test.ts @@ -9,6 +9,7 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; const histogram = vi.hoisted(() => ({ enable: vi.fn(), disable: vi.fn(), + reset: vi.fn(), mean: Number.NaN, max: Number.NaN, percentile: vi.fn((_percentile: number) => Number.NaN), @@ -18,6 +19,8 @@ vi.mock('node:perf_hooks', () => ({ monitorEventLoopDelay: vi.fn(() => histogram), })); +const cpuUsage = vi.hoisted(() => vi.fn()); + describe('startEventLoopLagMonitor', () => { let startEventLoopLagMonitor: typeof import('./event-loop-lag.js').startEventLoopLagMonitor; @@ -25,6 +28,11 @@ describe('startEventLoopLagMonitor', () => { vi.resetModules(); vi.clearAllMocks(); vi.useFakeTimers(); + cpuUsage.mockReset(); + cpuUsage + .mockReturnValueOnce({ user: 0, system: 0 }) + .mockReturnValue({ user: 0, system: 0 }); + vi.spyOn(process, 'cpuUsage').mockImplementation(cpuUsage); histogram.mean = Number.NaN; histogram.max = Number.NaN; histogram.percentile.mockReturnValue(Number.NaN); @@ -108,6 +116,167 @@ describe('startEventLoopLagMonitor', () => { monitor.dispose(); }); + it('resets suspended samples without reporting them as stalls', async () => { + const onNewMaxStall = vi.fn(); + histogram.max = 15_000_000_000; + const monitor = startEventLoopLagMonitor({ + resolutionMs: 10, + stallThresholdMs: 1_000, + suspendThresholdMs: 10_000, + onNewMaxStall, + }); + + await vi.advanceTimersByTimeAsync(10); + + expect(histogram.reset).toHaveBeenCalledTimes(1); + expect(onNewMaxStall).not.toHaveBeenCalled(); + + monitor.dispose(); + }); + + it('resets suspended samples without a stall callback', () => { + histogram.mean = 15_000_000_000; + histogram.max = 15_000_000_000; + histogram.percentile.mockReturnValue(15_000_000_000); + histogram.reset.mockImplementation(() => { + histogram.mean = Number.NaN; + histogram.max = Number.NaN; + histogram.percentile.mockReturnValue(Number.NaN); + }); + const monitor = startEventLoopLagMonitor({ + resolutionMs: 10, + suspendThresholdMs: 10_000, + }); + + vi.setSystemTime(Date.now() + 10_000); + expect(monitor.snapshot()).toEqual({ + meanMs: 0, + p50Ms: 0, + p99Ms: 0, + maxMs: 0, + }); + expect(histogram.reset).toHaveBeenCalledTimes(1); + + monitor.dispose(); + }); + + it('reports lower stalls after resetting a suspended sample', async () => { + const onNewMaxStall = vi.fn(); + histogram.max = 15_000_000_000; + histogram.reset.mockImplementation(() => { + histogram.max = Number.NaN; + }); + const monitor = startEventLoopLagMonitor({ + resolutionMs: 10, + stallThresholdMs: 1_000, + suspendThresholdMs: 10_000, + onNewMaxStall, + }); + + await vi.advanceTimersByTimeAsync(10); + histogram.max = 5_000_000_000; + await vi.advanceTimersByTimeAsync(10); + + expect(histogram.reset).toHaveBeenCalledTimes(1); + expect(onNewMaxStall).toHaveBeenCalledOnce(); + expect(onNewMaxStall).toHaveBeenCalledWith(5_000); + + monitor.dispose(); + }); + + it('reports a real stall below the suspend threshold', async () => { + const onNewMaxStall = vi.fn(); + histogram.max = 5_000_000_000; + const monitor = startEventLoopLagMonitor({ + resolutionMs: 10, + stallThresholdMs: 1_000, + suspendThresholdMs: 10_000, + onNewMaxStall, + }); + + await vi.advanceTimersByTimeAsync(10); + + expect(histogram.reset).not.toHaveBeenCalled(); + expect(onNewMaxStall).toHaveBeenCalledWith(5_000); + + monitor.dispose(); + }); + + it('reports a stall just below the default suspend threshold', async () => { + const onNewMaxStall = vi.fn(); + histogram.max = 599_000_000_000; + const monitor = startEventLoopLagMonitor({ + resolutionMs: 10, + stallThresholdMs: 1_000, + onNewMaxStall, + }); + + await vi.advanceTimersByTimeAsync(10); + + expect(histogram.reset).not.toHaveBeenCalled(); + expect(onNewMaxStall).toHaveBeenCalledWith(599_000); + + monitor.dispose(); + }); + + it('resets a low-CPU sample at the default suspend threshold', async () => { + const onNewMaxStall = vi.fn(); + histogram.max = 600_000_000_000; + const monitor = startEventLoopLagMonitor({ + resolutionMs: 10, + stallThresholdMs: 1_000, + onNewMaxStall, + }); + + await vi.advanceTimersByTimeAsync(10); + + expect(histogram.reset).toHaveBeenCalledTimes(1); + expect(onNewMaxStall).not.toHaveBeenCalled(); + + monitor.dispose(); + }); + + it('reports a long active stall when CPU time advanced', async () => { + const onNewMaxStall = vi.fn(); + cpuUsage + .mockReset() + .mockReturnValueOnce({ user: 0, system: 0 }) + .mockReturnValue({ user: 20_000_000, system: 0 }); + histogram.max = 600_000_000_000; + const monitor = startEventLoopLagMonitor({ + resolutionMs: 10, + stallThresholdMs: 1_000, + onNewMaxStall, + }); + + await vi.advanceTimersByTimeAsync(10); + + expect(histogram.reset).not.toHaveBeenCalled(); + expect(onNewMaxStall).toHaveBeenCalledWith(600_000); + + monitor.dispose(); + }); + + it('reports rather than suppressing when CPU usage is unavailable', async () => { + const onNewMaxStall = vi.fn(); + cpuUsage.mockImplementation(() => { + throw new Error('cpu accounting unavailable'); + }); + histogram.max = 600_000_000_000; + const monitor = startEventLoopLagMonitor({ + resolutionMs: 10, + stallThresholdMs: 1_000, + onNewMaxStall, + }); + + await vi.advanceTimersByTimeAsync(10); + + expect(histogram.reset).not.toHaveBeenCalled(); + expect(onNewMaxStall).toHaveBeenCalledWith(600_000); + + monitor.dispose(); + }); + it('enables and disables the underlying histogram', () => { const monitor = startEventLoopLagMonitor(); diff --git a/packages/core/src/telemetry/event-loop-lag.ts b/packages/core/src/telemetry/event-loop-lag.ts index 2aacdaf84cf..9a36a372407 100644 --- a/packages/core/src/telemetry/event-loop-lag.ts +++ b/packages/core/src/telemetry/event-loop-lag.ts @@ -21,11 +21,20 @@ export interface EventLoopLagMonitor { export interface EventLoopLagMonitorOptions { resolutionMs?: number; stallThresholdMs?: number; + /** + * Consider a long scheduling gap to be host suspension only when process CPU + * time stayed below this fraction of the gap. Default: 1%. + */ + suspendCpuRatio?: number; + /** Minimum event-loop gap eligible for host-suspension filtering. */ + suspendThresholdMs?: number; onNewMaxStall?: (maxMs: number) => void; } const DEFAULT_RESOLUTION_MS = 20; const DEFAULT_STALL_THRESHOLD_MS = 1_000; +const DEFAULT_SUSPEND_THRESHOLD_MS = 10 * 60 * 1_000; +const DEFAULT_SUSPEND_CPU_RATIO = 0.01; const NS_PER_MS = 1_000_000; export function startEventLoopLagMonitor( @@ -39,16 +48,45 @@ export function startEventLoopLagMonitor( options.stallThresholdMs, DEFAULT_STALL_THRESHOLD_MS, ); + const suspendThresholdMs = positiveFiniteOrDefault( + options.suspendThresholdMs, + DEFAULT_SUSPEND_THRESHOLD_MS, + ); + const suspendCpuRatio = fractionOrDefault( + options.suspendCpuRatio, + DEFAULT_SUSPEND_CPU_RATIO, + ); const histogram = monitorEventLoopDelay({ resolution: resolutionMs }); histogram.enable(); let disposed = false; let lastReportedMaxMs = 0; + let lastCheckTimeMs = Date.now(); + let lastCpuUsage = safeCpuUsage(); const readMaxMs = () => nsToMs(histogram.max); - const checkForNewMaxStall = () => { - if (disposed || !options.onNewMaxStall) return; + const checkHistogram = () => { + if (disposed) return; + const nowMs = Date.now(); + const cpuUsage = safeCpuUsage(); + const elapsedMs = Math.max(0, nowMs - lastCheckTimeMs); const maxMs = readMaxMs(); - if (maxMs >= stallThresholdMs && maxMs > lastReportedMaxMs) { + const cpuRatio = calculateCpuRatio(lastCpuUsage, cpuUsage, elapsedMs); + lastCheckTimeMs = nowMs; + if (cpuUsage) lastCpuUsage = cpuUsage; + if ( + maxMs >= suspendThresholdMs && + cpuRatio !== undefined && + cpuRatio <= suspendCpuRatio + ) { + histogram.reset(); + lastReportedMaxMs = 0; + return; + } + if ( + options.onNewMaxStall && + maxMs >= stallThresholdMs && + maxMs > lastReportedMaxMs + ) { lastReportedMaxMs = maxMs; try { options.onNewMaxStall(maxMs); @@ -57,14 +95,12 @@ export function startEventLoopLagMonitor( } } }; - const interval = - options.onNewMaxStall !== undefined - ? setInterval(checkForNewMaxStall, resolutionMs) - : undefined; - interval?.unref(); + const interval = setInterval(checkHistogram, resolutionMs); + interval.unref(); return { snapshot(): EventLoopLagSnapshot { + checkHistogram(); return { meanMs: nsToMs(histogram.mean), p50Ms: nsToMs(histogram.percentile(50)), @@ -75,14 +111,31 @@ export function startEventLoopLagMonitor( dispose(): void { if (disposed) return; disposed = true; - if (interval) { - clearInterval(interval); - } + clearInterval(interval); histogram.disable(); }, }; } +function safeCpuUsage(): NodeJS.CpuUsage | undefined { + try { + return process.cpuUsage(); + } catch { + return undefined; + } +} + +function calculateCpuRatio( + previous: NodeJS.CpuUsage | undefined, + current: NodeJS.CpuUsage | undefined, + elapsedMs: number, +): number | undefined { + if (!previous || !current || elapsedMs <= 0) return undefined; + const cpuMicroseconds = + current.user - previous.user + (current.system - previous.system); + return Math.max(0, cpuMicroseconds / (elapsedMs * 1_000)); +} + function nsToMs(value: number): number { return Number.isFinite(value) ? value / NS_PER_MS : 0; } @@ -92,3 +145,12 @@ function positiveFiniteOrDefault(value: number | undefined, fallback: number) { ? value : fallback; } + +function fractionOrDefault(value: number | undefined, fallback: number) { + return value !== undefined && + Number.isFinite(value) && + value >= 0 && + value <= 1 + ? value + : fallback; +} From 6d2498da5ade4ac61ff50f70470f2f7bd81c6990 Mon Sep 17 00:00:00 2001 From: yiliang114 Date: Fri, 31 Jul 2026 17:48:08 +0800 Subject: [PATCH 2/9] test(channel): cover bridge recovery gates --- .../channels/base/src/ChannelBase.test.ts | 60 +++++++++++++++++++ 1 file changed, 60 insertions(+) diff --git a/packages/channels/base/src/ChannelBase.test.ts b/packages/channels/base/src/ChannelBase.test.ts index c613b74d23f..6b62452a853 100644 --- a/packages/channels/base/src/ChannelBase.test.ts +++ b/packages/channels/base/src/ChannelBase.test.ts @@ -14963,6 +14963,44 @@ describe('ChannelBase', () => { ]); }); + it('rechecks bridge recovery before a queued followup prompt starts', async () => { + const recoveryState: { current?: Promise } = {}; + let releaseRecovery: (() => void) | undefined; + let resolveFirst!: (v: string) => void; + const firstPrompt = new Promise((r) => { + resolveFirst = r; + }); + let callCount = 0; + (bridge.prompt as ReturnType).mockImplementation(() => { + callCount++; + if (callCount === 1) return firstPrompt; + return Promise.resolve(`response-${callCount}`); + }); + const ch = createChannel( + { dispatchMode: 'followup' }, + { bridgeRecovery: () => recoveryState.current }, + ); + + const first = ch.handleInbound(envelope({ text: 'task one' })); + await vi.waitFor(() => expect(bridge.prompt).toHaveBeenCalledTimes(1)); + + const second = ch.handleInbound(envelope({ text: 'task two' })); + await Promise.resolve(); + recoveryState.current = new Promise((resolve) => { + releaseRecovery = resolve; + }); + resolveFirst('response-1'); + await first; + await Promise.resolve(); + + expect(bridge.prompt).toHaveBeenCalledTimes(1); + + releaseRecovery!(); + await second; + + expect(bridge.prompt).toHaveBeenCalledTimes(2); + }); + it('steer is the default mode when dispatchMode not set', async () => { let resolveFirst!: (v: string) => void; const firstPrompt = new Promise((r) => { @@ -16015,6 +16053,28 @@ describe('ChannelBase', () => { expect(collectedPrompt).toContain('follow-up while webhook runs'); }); + it('waits for bridge recovery before resolving a webhook session', async () => { + let releaseBridge: (() => void) | undefined; + const bridgeReady = new Promise((resolve) => { + releaseBridge = resolve; + }); + const ch = createChannel( + { approvalMode: 'yolo', webhooks }, + { bridgeRecovery: () => bridgeReady }, + ); + ch.proactiveSupported = true; + + const run = ch.runWebhookTask(webhookTask); + await Promise.resolve(); + expect(bridge.newSession).not.toHaveBeenCalled(); + expect(bridge.prompt).not.toHaveBeenCalled(); + + releaseBridge!(); + await expect(run).resolves.toBe('agent response'); + + expect(bridge.prompt).toHaveBeenCalledTimes(1); + }); + it('rechecks bridge recovery before a webhook prompt starts', async () => { const recoveryState: { current?: Promise } = {}; let releaseRecovery: (() => void) | undefined; From 929ed183d13434c9a9a1f3f8ac4755ee5ae738a9 Mon Sep 17 00:00:00 2001 From: "jinjing.zzj" Date: Fri, 31 Jul 2026 19:31:16 +0800 Subject: [PATCH 3/9] fix(channel): close the 5-10min wake misclassification window and dedupe recovery - Lower the default event-loop suspend threshold to 5 minutes so any low-CPU sleep gap is filtered before it reaches the AcpBridge stall-kill threshold; pin the invariant with a cross-package test. - Extract the duplicated startSingle/startAll bridge recovery into a shared createBridgeRecovery helper. --- packages/channels/base/src/index.ts | 2 +- .../acp-integration/stall-thresholds.test.ts | 22 ++ packages/cli/src/commands/channel/start.ts | 318 +++++++++--------- .../core/src/telemetry/event-loop-lag.test.ts | 6 +- packages/core/src/telemetry/event-loop-lag.ts | 10 +- packages/core/src/telemetry/index.ts | 1 + 6 files changed, 187 insertions(+), 172 deletions(-) create mode 100644 packages/cli/src/acp-integration/stall-thresholds.test.ts diff --git a/packages/channels/base/src/index.ts b/packages/channels/base/src/index.ts index f4d0477e92a..74d62854403 100644 --- a/packages/channels/base/src/index.ts +++ b/packages/channels/base/src/index.ts @@ -1,6 +1,6 @@ export { getGlobalQwenDir, resolvePath } from './paths.js'; export { PollingChannelBase } from './PollingChannelBase.js'; -export { AcpBridge } from './AcpBridge.js'; +export { ACP_EVENT_LOOP_STALL_RESTART_MS, AcpBridge } from './AcpBridge.js'; export type { AvailableCommand, BridgeSessionInfo, diff --git a/packages/cli/src/acp-integration/stall-thresholds.test.ts b/packages/cli/src/acp-integration/stall-thresholds.test.ts new file mode 100644 index 00000000000..bad3917677b --- /dev/null +++ b/packages/cli/src/acp-integration/stall-thresholds.test.ts @@ -0,0 +1,22 @@ +/** + * @license + * Copyright 2026 Qwen Team + * SPDX-License-Identifier: Apache-2.0 + */ + +import { describe, expect, it } from 'vitest'; +import { ACP_EVENT_LOOP_STALL_RESTART_MS } from '@qwen-code/channel-base'; +import { DEFAULT_EVENT_LOOP_SUSPEND_THRESHOLD_MS } from '@qwen-code/qwen-code-core'; + +describe('acp stall thresholds', () => { + // The ACP agent starts its event-loop lag monitor with default options + // (acpAgent.ts), while AcpBridge kills the child once a reported stall + // reaches ACP_EVENT_LOOP_STALL_RESTART_MS. If the suspend threshold ever + // exceeded the kill threshold, a host sleep in between the two values would + // be reported as a stall and kill a healthy child on wake. + it('keeps host-suspension filtering at or below the bridge stall-kill threshold', () => { + expect(DEFAULT_EVENT_LOOP_SUSPEND_THRESHOLD_MS).toBeLessThanOrEqual( + ACP_EVENT_LOOP_STALL_RESTART_MS, + ); + }); +}); diff --git a/packages/cli/src/commands/channel/start.ts b/packages/cli/src/commands/channel/start.ts index cdbcde55332..4631b900969 100644 --- a/packages/cli/src/commands/channel/start.ts +++ b/packages/cli/src/commands/channel/start.ts @@ -18,7 +18,11 @@ import { ChannelLoopStore, SessionRouter, } from '@qwen-code/channel-base'; -import type { ChannelBase, ChannelBaseOptions } from '@qwen-code/channel-base'; +import type { + AcpBridgeOptions, + ChannelBase, + ChannelBaseOptions, +} from '@qwen-code/channel-base'; import { findCliEntryPath, parseChannelConfig } from './config-utils.js'; import { resolveProxy } from './proxy.js'; import { @@ -148,6 +152,129 @@ function createBridgeReadinessGate(): { }; } +interface BridgeRecoveryOptions { + bridgeOpts: AcpBridgeOptions; + router: SessionRouter; + channels: Map; + scheduler: ChannelLoopScheduler | undefined; + bridgeReadiness: ReturnType; + isShuttingDown: () => boolean; + getBridge: () => AcpBridge; + setBridge: (bridge: AcpBridge) => void; +} + +/** + * Rebuild the ACP bridge after a disconnect while keeping channel adapters + * connected. Shared by the standalone and all-channel start paths; the only + * per-path state comes in through the accessors. + */ +function createBridgeRecovery(options: BridgeRecoveryOptions): { + attachDisconnectHandler: (bridge: AcpBridge) => void; +} { + const { + bridgeOpts, + router, + channels, + scheduler, + bridgeReadiness, + isShuttingDown, + getBridge, + setBridge, + } = options; + const crashTimestamps: number[] = []; + let recoveryTask: Promise | undefined; + let recoveryRequested = false; + + const attachDisconnectHandler = (failedBridge: AcpBridge): void => { + failedBridge.on('disconnected', () => { + if (isShuttingDown()) return; + if (recoveryTask) { + if (failedBridge === getBridge()) recoveryRequested = true; + return; + } + recoverBridge(); + }); + }; + + const disconnectAllChannels = (): void => { + for (const channel of channels.values()) { + try { + channel.disconnect(); + } catch { + // best-effort + } + } + }; + + const recoverBridge = (): void => { + bridgeReadiness.block(); + const task = (async () => { + do { + recoveryRequested = false; + const now = Date.now(); + crashTimestamps.push(now); + // Only count crashes within the recent window + const recentCrashes = crashTimestamps.filter( + (timestamp) => now - timestamp < CRASH_WINDOW_MS, + ); + + if (recentCrashes.length > MAX_CRASH_RESTARTS) { + writeStderrLine( + `[Channel] Bridge crashed ${recentCrashes.length} times in ${CRASH_WINDOW_MS / 1000}s. Giving up.`, + ); + scheduler?.stop(); + disconnectAllChannels(); + router.clearAll(); + removeServiceInfo(); + process.exit(1); + } + + writeStderrLine( + `[Channel] Bridge crashed (${recentCrashes.length}/${MAX_CRASH_RESTARTS} in window). Restarting in ${RESTART_DELAY_MS / 1000}s...`, + ); + await new Promise((resolve) => setTimeout(resolve, RESTART_DELAY_MS)); + + const bridge = new AcpBridge(bridgeOpts); + setBridge(bridge); + attachDisconnectHandler(bridge); + await bridge.start(); + router.setBridge(bridge); + for (const channel of channels.values()) { + channel.setBridge(bridge); + } + registerToolCallDispatch(bridge, router, channels); + registerBackgroundResponseRelay(bridge, router, channels); + registerPermissionRelay(bridge, router, channels); + registerSessionCleanup(bridge, router, channels); + + const result = await router.restoreSessions(); + writeStdoutLine( + `[Channel] Bridge restarted. Sessions restored: ${result.restored}, failed: ${result.failed}`, + ); + } while (recoveryRequested && !isShuttingDown()); + })() + .catch((err) => { + writeStderrLine( + `[Channel] Failed to restart bridge: ${err instanceof Error ? err.message : String(err)}`, + ); + scheduler?.stop(); + disconnectAllChannels(); + router.clearAll(); + removeServiceInfo(); + process.exit(1); + }) + .finally(() => { + if (recoveryTask === task) { + recoveryTask = undefined; + bridgeReadiness.release(); + } + }); + recoveryTask = task; + }; + + return { attachDisconnectHandler }; +} + /** Check for duplicate instance and abort if one is already running. */ function checkDuplicateInstance(): void { const existing = readServiceInfo(); @@ -202,9 +329,6 @@ async function startSingle( const cliEntryPath = findCliEntryPath(); let shuttingDown = false; - const crashTimestamps: number[] = []; - let recoveryTask: Promise | undefined; - let recoveryRequested = false; const bridgeReadiness = createBridgeReadinessGate(); const bridgeOpts = { cliEntryPath, cwd: config.cwd, model: config.model }; @@ -260,79 +384,18 @@ async function startSingle( scheduler?.start(); writeStdoutLine(`[Channel] "${name}" is running. Press Ctrl+C to stop.`); - const attachDisconnectHandler = (failedBridge: AcpBridge): void => { - failedBridge.on('disconnected', () => { - if (shuttingDown) return; - if (recoveryTask) { - if (failedBridge === bridge) recoveryRequested = true; - return; - } - recoverBridge(); - }); - }; - - const recoverBridge = (): void => { - bridgeReadiness.block(); - const task = (async () => { - do { - recoveryRequested = false; - const now = Date.now(); - crashTimestamps.push(now); - // Only count crashes within the recent window - const recentCrashes = crashTimestamps.filter( - (timestamp) => now - timestamp < CRASH_WINDOW_MS, - ); - - if (recentCrashes.length > MAX_CRASH_RESTARTS) { - writeStderrLine( - `[Channel] Bridge crashed ${recentCrashes.length} times in ${CRASH_WINDOW_MS / 1000}s. Giving up.`, - ); - scheduler?.stop(); - channel.disconnect(); - router.clearAll(); - removeServiceInfo(); - process.exit(1); - } - - writeStderrLine( - `[Channel] Bridge crashed (${recentCrashes.length}/${MAX_CRASH_RESTARTS} in window). Restarting in ${RESTART_DELAY_MS / 1000}s...`, - ); - await new Promise((resolve) => setTimeout(resolve, RESTART_DELAY_MS)); - - bridge = new AcpBridge(bridgeOpts); - attachDisconnectHandler(bridge); - await bridge.start(); - router.setBridge(bridge); - channel.setBridge(bridge); - registerToolCallDispatch(bridge, router, channels); - registerBackgroundResponseRelay(bridge, router, channels); - registerPermissionRelay(bridge, router, channels); - registerSessionCleanup(bridge, router, channels); - - const result = await router.restoreSessions(); - writeStdoutLine( - `[Channel] Bridge restarted. Sessions restored: ${result.restored}, failed: ${result.failed}`, - ); - } while (recoveryRequested && !shuttingDown); - })() - .catch((err) => { - writeStderrLine( - `[Channel] Failed to restart bridge: ${err instanceof Error ? err.message : String(err)}`, - ); - scheduler?.stop(); - channel.disconnect(); - router.clearAll(); - removeServiceInfo(); - process.exit(1); - }) - .finally(() => { - if (recoveryTask === task) { - recoveryTask = undefined; - bridgeReadiness.release(); - } - }); - recoveryTask = task; - }; + const { attachDisconnectHandler } = createBridgeRecovery({ + bridgeOpts, + router, + channels, + scheduler, + bridgeReadiness, + isShuttingDown: () => shuttingDown, + getBridge: () => bridge, + setBridge: (next) => { + bridge = next; + }, + }); attachDisconnectHandler(bridge); const shutdown = () => { @@ -383,9 +446,6 @@ async function startAll( const cliEntryPath = findCliEntryPath(); const defaultCwd = process.cwd(); let shuttingDown = false; - const crashTimestamps: number[] = []; - let recoveryTask: Promise | undefined; - let recoveryRequested = false; const bridgeReadiness = createBridgeReadinessGate(); const bridgeOpts = { @@ -467,92 +527,18 @@ async function startAll( `[Channel] Running ${connectedCount} channel(s). Press Ctrl+C to stop.`, ); - const attachDisconnectHandler = (failedBridge: AcpBridge): void => { - failedBridge.on('disconnected', () => { - if (shuttingDown) return; - if (recoveryTask) { - if (failedBridge === bridge) recoveryRequested = true; - return; - } - recoverBridge(); - }); - }; - - const recoverBridge = (): void => { - bridgeReadiness.block(); - const task = (async () => { - do { - recoveryRequested = false; - const now = Date.now(); - crashTimestamps.push(now); - const recentCrashes = crashTimestamps.filter( - (timestamp) => now - timestamp < CRASH_WINDOW_MS, - ); - - if (recentCrashes.length > MAX_CRASH_RESTARTS) { - writeStderrLine( - `[Channel] Bridge crashed ${recentCrashes.length} times in ${CRASH_WINDOW_MS / 1000}s. Giving up.`, - ); - scheduler?.stop(); - for (const channel of channels.values()) { - try { - channel.disconnect(); - } catch { - // best-effort - } - } - router.clearAll(); - removeServiceInfo(); - process.exit(1); - } - - writeStderrLine( - `[Channel] Bridge crashed (${recentCrashes.length}/${MAX_CRASH_RESTARTS} in window). Restarting in ${RESTART_DELAY_MS / 1000}s...`, - ); - await new Promise((resolve) => setTimeout(resolve, RESTART_DELAY_MS)); - - bridge = new AcpBridge(bridgeOpts); - attachDisconnectHandler(bridge); - await bridge.start(); - router.setBridge(bridge); - for (const channel of channels.values()) { - channel.setBridge(bridge); - } - registerToolCallDispatch(bridge, router, channels); - registerBackgroundResponseRelay(bridge, router, channels); - registerPermissionRelay(bridge, router, channels); - registerSessionCleanup(bridge, router, channels); - - const result = await router.restoreSessions(); - writeStdoutLine( - `[Channel] Bridge restarted. Sessions restored: ${result.restored}, failed: ${result.failed}`, - ); - } while (recoveryRequested && !shuttingDown); - })() - .catch((err) => { - writeStderrLine( - `[Channel] Failed to restart bridge: ${err instanceof Error ? err.message : String(err)}`, - ); - scheduler?.stop(); - for (const channel of channels.values()) { - try { - channel.disconnect(); - } catch { - // best-effort - } - } - router.clearAll(); - removeServiceInfo(); - process.exit(1); - }) - .finally(() => { - if (recoveryTask === task) { - recoveryTask = undefined; - bridgeReadiness.release(); - } - }); - recoveryTask = task; - }; + const { attachDisconnectHandler } = createBridgeRecovery({ + bridgeOpts, + router, + channels, + scheduler, + bridgeReadiness, + isShuttingDown: () => shuttingDown, + getBridge: () => bridge, + setBridge: (next) => { + bridge = next; + }, + }); attachDisconnectHandler(bridge); const shutdown = () => { diff --git a/packages/core/src/telemetry/event-loop-lag.test.ts b/packages/core/src/telemetry/event-loop-lag.test.ts index 5fc5123dcea..c57afbe05e4 100644 --- a/packages/core/src/telemetry/event-loop-lag.test.ts +++ b/packages/core/src/telemetry/event-loop-lag.test.ts @@ -204,7 +204,7 @@ describe('startEventLoopLagMonitor', () => { it('reports a stall just below the default suspend threshold', async () => { const onNewMaxStall = vi.fn(); - histogram.max = 599_000_000_000; + histogram.max = 299_000_000_000; const monitor = startEventLoopLagMonitor({ resolutionMs: 10, stallThresholdMs: 1_000, @@ -214,14 +214,14 @@ describe('startEventLoopLagMonitor', () => { await vi.advanceTimersByTimeAsync(10); expect(histogram.reset).not.toHaveBeenCalled(); - expect(onNewMaxStall).toHaveBeenCalledWith(599_000); + expect(onNewMaxStall).toHaveBeenCalledWith(299_000); monitor.dispose(); }); it('resets a low-CPU sample at the default suspend threshold', async () => { const onNewMaxStall = vi.fn(); - histogram.max = 600_000_000_000; + histogram.max = 300_000_000_000; const monitor = startEventLoopLagMonitor({ resolutionMs: 10, stallThresholdMs: 1_000, diff --git a/packages/core/src/telemetry/event-loop-lag.ts b/packages/core/src/telemetry/event-loop-lag.ts index 9a36a372407..331f13ee609 100644 --- a/packages/core/src/telemetry/event-loop-lag.ts +++ b/packages/core/src/telemetry/event-loop-lag.ts @@ -33,7 +33,13 @@ export interface EventLoopLagMonitorOptions { const DEFAULT_RESOLUTION_MS = 20; const DEFAULT_STALL_THRESHOLD_MS = 1_000; -const DEFAULT_SUSPEND_THRESHOLD_MS = 10 * 60 * 1_000; +/** + * Default minimum gap treated as host suspension. Kept at or below the ACP + * bridge stall-kill threshold (`ACP_EVENT_LOOP_STALL_RESTART_MS`) so a low-CPU + * sleep gap is always filtered before it can be reported as a kill-eligible + * stall. + */ +export const DEFAULT_EVENT_LOOP_SUSPEND_THRESHOLD_MS = 5 * 60 * 1_000; const DEFAULT_SUSPEND_CPU_RATIO = 0.01; const NS_PER_MS = 1_000_000; @@ -50,7 +56,7 @@ export function startEventLoopLagMonitor( ); const suspendThresholdMs = positiveFiniteOrDefault( options.suspendThresholdMs, - DEFAULT_SUSPEND_THRESHOLD_MS, + DEFAULT_EVENT_LOOP_SUSPEND_THRESHOLD_MS, ); const suspendCpuRatio = fractionOrDefault( options.suspendCpuRatio, diff --git a/packages/core/src/telemetry/index.ts b/packages/core/src/telemetry/index.ts index 0cea7f1d96f..11f2f140eb2 100644 --- a/packages/core/src/telemetry/index.ts +++ b/packages/core/src/telemetry/index.ts @@ -225,6 +225,7 @@ export type { DaemonPipeDirection, } from './daemon-metrics.js'; export { + DEFAULT_EVENT_LOOP_SUSPEND_THRESHOLD_MS, startEventLoopLagMonitor, type EventLoopLagMonitor, type EventLoopLagMonitorOptions, From 71e9272af97ce15e4a4c60b742f1216202dc695e Mon Sep 17 00:00:00 2001 From: "jinjing.zzj" Date: Fri, 31 Jul 2026 19:34:22 +0800 Subject: [PATCH 4/9] test(channel): use observe in observedContacts recovery-gate mock The mock provided a record method but ChannelBaseOptions.observedContacts expects observe, so the slow-contact-recording path was never exercised. --- packages/channels/base/src/ChannelBase.test.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/channels/base/src/ChannelBase.test.ts b/packages/channels/base/src/ChannelBase.test.ts index 6b62452a853..59b9b0a9162 100644 --- a/packages/channels/base/src/ChannelBase.test.ts +++ b/packages/channels/base/src/ChannelBase.test.ts @@ -16157,7 +16157,7 @@ describe('ChannelBase', () => { { bridgeRecovery: () => recoveryState.current, observedContacts: { - record: vi.fn().mockImplementation(() => contactRecord), + observe: vi.fn().mockImplementation(() => contactRecord), }, }, ); From 144d8d9d34b5939df564de208ff72088f3bfacd3 Mon Sep 17 00:00:00 2001 From: yiliang114 Date: Fri, 31 Jul 2026 19:35:14 +0800 Subject: [PATCH 5/9] fix(channel): harden wake recovery --- packages/channels/base/src/AcpBridge.test.ts | 31 ++++- packages/channels/base/src/AcpBridge.ts | 37 +++++- .../channels/base/src/ChannelBase.test.ts | 31 +++++ packages/channels/base/src/ChannelBase.ts | 13 +- .../cli/src/acp-integration/acpAgent.test.ts | 4 + packages/cli/src/acp-integration/acpAgent.ts | 4 + .../cli/src/commands/channel/start.test.ts | 123 +++++++++++++++++- packages/cli/src/commands/channel/start.ts | 67 ++++++---- .../core/src/telemetry/event-loop-lag.test.ts | 67 +++++++++- packages/core/src/telemetry/event-loop-lag.ts | 8 +- 10 files changed, 344 insertions(+), 41 deletions(-) diff --git a/packages/channels/base/src/AcpBridge.test.ts b/packages/channels/base/src/AcpBridge.test.ts index 2312d100187..1967b9cb3f3 100644 --- a/packages/channels/base/src/AcpBridge.test.ts +++ b/packages/channels/base/src/AcpBridge.test.ts @@ -3,6 +3,7 @@ import type { RequestPermissionResponse } from '@agentclientprotocol/sdk'; import { ACP_EVENT_LOOP_STALL_RESTART_MS, ACP_PERMISSION_RESPONSE_TIMEOUT_MS, + ACP_START_TIMEOUT_MS, AcpBridge, } from './AcpBridge.js'; import { CHANNEL_LOOP_MCP_SERVER_NAME } from './ChannelLoopTools.js'; @@ -47,6 +48,7 @@ const child = vi.hoisted(() => { }); } + let initializeImplementation: () => Promise = () => Promise.resolve(); return { instances: [] as MockChild[], clients: [] as Array<{ @@ -57,6 +59,13 @@ const child = vi.hoisted(() => { cancel: ReturnType; }>, MockChild, + initializeImplementation: () => initializeImplementation(), + resetInitializeImplementation: () => { + initializeImplementation = () => Promise.resolve(); + }, + setInitializeImplementation: (implementation: () => Promise) => { + initializeImplementation = implementation; + }, spawn: vi.fn(() => { const instance = new MockChild(); child.instances.push(instance); @@ -80,7 +89,7 @@ vi.mock('@agentclientprotocol/sdk', () => ({ ClientSideConnection: vi.fn().mockImplementation((createClient) => { const client = createClient(); const connection = { - initialize: vi.fn().mockResolvedValue(undefined), + initialize: vi.fn(() => child.initializeImplementation()), cancel: vi.fn().mockResolvedValue(undefined), }; child.clients.push(client); @@ -131,6 +140,26 @@ describe('AcpBridge', () => { child.clients.length = 0; child.connections.length = 0; child.spawn.mockClear(); + child.resetInitializeImplementation(); + }); + + it('times out bridge initialization and stops the child', async () => { + vi.useFakeTimers(); + child.setInitializeImplementation(() => new Promise(() => {})); + const bridge = new AcpBridge({ + cliEntryPath: '/tmp/qwen', + cwd: '/tmp', + }); + + const start = bridge.start(); + const rejection = expect(start).rejects.toThrow( + `ACP initialization timed out after ${ACP_START_TIMEOUT_MS}ms`, + ); + await vi.advanceTimersByTimeAsync(1000 + ACP_START_TIMEOUT_MS); + + await rejection; + expect(child.instances[0]!.kill).toHaveBeenCalledOnce(); + vi.useRealTimers(); }); it('registers the channel loop MCP server once across concurrent calls', async () => { diff --git a/packages/channels/base/src/AcpBridge.ts b/packages/channels/base/src/AcpBridge.ts index f2009bd05c1..881d0d905f5 100644 --- a/packages/channels/base/src/AcpBridge.ts +++ b/packages/channels/base/src/AcpBridge.ts @@ -43,6 +43,7 @@ export interface AcpBridgeOptions { } export const ACP_EVENT_LOOP_STALL_RESTART_MS = 5 * 60 * 1000; +export const ACP_START_TIMEOUT_MS = 30 * 1000; export const ACP_PERMISSION_RESPONSE_TIMEOUT_MS = 5 * 60 * 1000; const ACP_EVENT_LOOP_STALL_RE = /^\[perf\] acp agent event loop stall: max=(\d+(?:\.\d+)?)ms/m; @@ -178,11 +179,20 @@ export class AcpBridge extends EventEmitter implements ChannelAgentBridge { stream, ); - await this.connection.initialize({ - protocolVersion: PROTOCOL_VERSION, - clientCapabilities: {}, - }); - await this.registerChannelLoopMcpServer(); + try { + await withTimeout( + this.connection.initialize({ + protocolVersion: PROTOCOL_VERSION, + clientCapabilities: {}, + }), + ACP_START_TIMEOUT_MS, + `ACP initialization timed out after ${ACP_START_TIMEOUT_MS}ms`, + ); + await this.registerChannelLoopMcpServer(); + } catch (error) { + this.stop(); + throw error; + } } registerChannelLoopToolHandler(handler: ChannelLoopToolHandler): void { @@ -622,6 +632,23 @@ export class AcpBridge extends EventEmitter implements ChannelAgentBridge { } } +async function withTimeout( + operation: Promise, + timeoutMs: number, + message: string, +): Promise { + let timer: ReturnType | undefined; + const timeout = new Promise((_, reject) => { + timer = setTimeout(() => reject(new Error(message)), timeoutMs); + timer.unref?.(); + }); + try { + return await Promise.race([operation, timeout]); + } finally { + if (timer) clearTimeout(timer); + } +} + function isSkippedMcpRegistration(result: unknown): boolean { return ( typeof result === 'object' && diff --git a/packages/channels/base/src/ChannelBase.test.ts b/packages/channels/base/src/ChannelBase.test.ts index 59b9b0a9162..67c2563fd22 100644 --- a/packages/channels/base/src/ChannelBase.test.ts +++ b/packages/channels/base/src/ChannelBase.test.ts @@ -16145,6 +16145,37 @@ describe('ChannelBase', () => { expect(bridge.prompt).toHaveBeenCalledTimes(1); }); + it('waits for recovery barriers that replace one another', async () => { + let releaseFirst: (() => void) | undefined; + let releaseSecond: (() => void) | undefined; + const recoveryState: { current?: Promise } = { + current: new Promise((resolve) => { + releaseFirst = resolve; + }), + }; + const ch = createChannel( + {}, + { bridgeRecovery: () => recoveryState.current }, + ); + + const inbound = ch.handleInbound(envelope()); + await Promise.resolve(); + recoveryState.current = new Promise((resolve) => { + releaseSecond = resolve; + }); + releaseFirst!(); + await Promise.resolve(); + + expect(bridge.newSession).not.toHaveBeenCalled(); + expect(bridge.prompt).not.toHaveBeenCalled(); + + recoveryState.current = undefined; + releaseSecond!(); + await inbound; + + expect(bridge.prompt).toHaveBeenCalledOnce(); + }); + it('rechecks bridge recovery after inbound preprocessing has started', async () => { const recoveryState: { current?: Promise } = {}; let releaseRecovery: (() => void) | undefined; diff --git a/packages/channels/base/src/ChannelBase.ts b/packages/channels/base/src/ChannelBase.ts index 651c9d62ebd..7985826ccab 100644 --- a/packages/channels/base/src/ChannelBase.ts +++ b/packages/channels/base/src/ChannelBase.ts @@ -350,6 +350,10 @@ function isUnattendedWebhookApprovalMode(mode: string | undefined): boolean { export abstract class ChannelBase { protected config: ChannelConfig; + /** + * Recovery invariant: every path that resolves a session or calls the bridge + * must await waitForBridgeRecovery() immediately before that operation. + */ protected bridge: ChannelAgentBridge; protected groupGate: GroupGate; protected dmGate: DmGate; @@ -4914,8 +4918,13 @@ export abstract class ChannelBase { /** Wait until the currently active bridge recovery, if any, has completed. */ private async waitForBridgeRecovery(): Promise { - const bridgeRecovery = this.bridgeRecovery?.(); - if (bridgeRecovery) await bridgeRecovery; + let completedRecovery: Promise | undefined; + while (true) { + const bridgeRecovery = this.bridgeRecovery?.(); + if (!bridgeRecovery || bridgeRecovery === completedRecovery) return; + await bridgeRecovery; + completedRecovery = bridgeRecovery; + } } /** diff --git a/packages/cli/src/acp-integration/acpAgent.test.ts b/packages/cli/src/acp-integration/acpAgent.test.ts index fd967bde581..bd526930818 100644 --- a/packages/cli/src/acp-integration/acpAgent.test.ts +++ b/packages/cli/src/acp-integration/acpAgent.test.ts @@ -1398,6 +1398,10 @@ describe('runAcpAgent shutdown cleanup', () => { }); expect(snapshot).toHaveBeenCalledTimes(1); + expect(startEventLoopLagMonitor).toHaveBeenCalledWith( + expect.objectContaining({ suspendThresholdMs: 5 * 60 * 1000 }), + ); + mockConnectionState.resolve(); await agentPromise; diff --git a/packages/cli/src/acp-integration/acpAgent.ts b/packages/cli/src/acp-integration/acpAgent.ts index bdfc381826e..df551a7ab4a 100644 --- a/packages/cli/src/acp-integration/acpAgent.ts +++ b/packages/cli/src/acp-integration/acpAgent.ts @@ -345,6 +345,9 @@ const POSIX_TMP_LOCAL_READ_ROOT = '/tmp'; const BTW_CHILD_TIMEOUT_MS = 55_000; const MCP_OAUTH_START_TIMEOUT_MS = 30_000; const SESSION_DRAIN_TIMEOUT_MS = 30_000; +// Match the parent channel bridge's restart threshold. Low-CPU gaps at this +// boundary are host suspension; active stalls still reach the parent watchdog. +const ACP_EVENT_LOOP_SUSPEND_THRESHOLD_MS = 5 * 60 * 1000; // Must be less than WORKSPACE_MEMORY_REMEMBER_TIMEOUT_MS (300s) in bridge.ts. const WORKSPACE_MEMORY_REMEMBER_CHILD_TIMEOUT_MS = 295_000; @@ -2876,6 +2879,7 @@ export async function runAcpAgent( process.stderr.write(`${warning}\n`); } const eventLoopMonitor = startEventLoopLagMonitor({ + suspendThresholdMs: ACP_EVENT_LOOP_SUSPEND_THRESHOLD_MS, onNewMaxStall: (maxMs) => { console.error(`[perf] acp agent event loop stall: max=${maxMs}ms`); }, diff --git a/packages/cli/src/commands/channel/start.test.ts b/packages/cli/src/commands/channel/start.test.ts index 0929ba093a8..ed7b20c5b4d 100644 --- a/packages/cli/src/commands/channel/start.test.ts +++ b/packages/cli/src/commands/channel/start.test.ts @@ -160,6 +160,7 @@ vi.mock('@qwen-code/channel-base', () => ({ })); import { + BRIDGE_SESSION_RESTORE_TIMEOUT_MS, resolveExtensionChannelEntrySpecifier, resolveProxy, startCommand, @@ -851,10 +852,10 @@ describe('startCommand.handler', () => { expect(disconnectedListener).toBeDefined(); vi.useFakeTimers(); - const firstRestart = disconnectedListener!(); - const secondRestart = disconnectedListener!(); - await vi.advanceTimersByTimeAsync(3000); - await Promise.all([firstRestart, secondRestart]); + disconnectedListener!(); + disconnectedListener!(); + await vi.advanceTimersByTimeAsync(6000); + disconnectedListener!(); expect(mockAcpBridge).toHaveBeenCalledTimes(2); expect(mockRouterRestoreSessions).toHaveBeenCalledTimes(1); @@ -941,6 +942,120 @@ describe('startCommand.handler', () => { } }); + it('cleans up standalone service state when replacement startup fails', async () => { + mockChannelConnect.mockResolvedValue(undefined); + mockBridgeStart + .mockResolvedValueOnce(undefined) + .mockRejectedValueOnce(new Error('replacement failed')); + mockChannelDisconnect.mockImplementationOnce(() => { + throw new Error('disconnect failed'); + }); + const channels = { telegram: { type: 'telegram' } }; + mockLoadSettings.mockReturnValue({ merged: { channels } }); + const processOnSpy = vi + .spyOn(process, 'on') + .mockImplementation(() => process); + const exitSpy = vi + .spyOn(process, 'exit') + .mockImplementation(() => undefined as never); + + try { + void invokeStartHandler({ name: 'telegram' }); + await new Promise((resolve) => setImmediate(resolve)); + const disconnectedListener = mockBridgeOn.mock.calls.find( + ([eventName]) => eventName === 'disconnected', + )?.[1] as (() => void) | undefined; + + vi.useFakeTimers(); + disconnectedListener!(); + await vi.advanceTimersByTimeAsync(3000); + await vi.waitFor(() => expect(mockRemoveServiceInfo).toHaveBeenCalled()); + + expect(mockRouterClearAll).toHaveBeenCalled(); + expect(exitSpy).toHaveBeenCalledWith(1); + } finally { + processOnSpy.mockRestore(); + exitSpy.mockRestore(); + vi.useRealTimers(); + } + }); + + it('stops the replacement bridge when session restore fails', async () => { + mockChannelConnect.mockResolvedValue(undefined); + mockRouterRestoreSessions.mockRejectedValueOnce( + new Error('restore failed'), + ); + const channels = { telegram: { type: 'telegram' } }; + mockLoadSettings.mockReturnValue({ merged: { channels } }); + const processOnSpy = vi + .spyOn(process, 'on') + .mockImplementation(() => process); + const exitSpy = vi + .spyOn(process, 'exit') + .mockImplementation(() => undefined as never); + + try { + void invokeStartHandler({ name: 'telegram' }); + await new Promise((resolve) => setImmediate(resolve)); + const disconnectedListener = mockBridgeOn.mock.calls.find( + ([eventName]) => eventName === 'disconnected', + )?.[1] as (() => void) | undefined; + + vi.useFakeTimers(); + disconnectedListener!(); + await vi.advanceTimersByTimeAsync(3000); + await vi.waitFor(() => expect(mockRemoveServiceInfo).toHaveBeenCalled()); + + expect(mockBridgeStart).toHaveBeenCalledTimes(2); + expect(mockBridgeStop).toHaveBeenCalledOnce(); + expect(exitSpy).toHaveBeenCalledWith(1); + } finally { + processOnSpy.mockRestore(); + exitSpy.mockRestore(); + vi.useRealTimers(); + } + }); + + it('times out a wedged session restore and stops the replacement bridge', async () => { + mockChannelConnect.mockResolvedValue(undefined); + mockRouterRestoreSessions.mockImplementationOnce( + () => new Promise(() => {}), + ); + const channels = { telegram: { type: 'telegram' } }; + mockLoadSettings.mockReturnValue({ merged: { channels } }); + const processOnSpy = vi + .spyOn(process, 'on') + .mockImplementation(() => process); + const exitSpy = vi + .spyOn(process, 'exit') + .mockImplementation(() => undefined as never); + + try { + void invokeStartHandler({ name: 'telegram' }); + await new Promise((resolve) => setImmediate(resolve)); + const disconnectedListener = mockBridgeOn.mock.calls.find( + ([eventName]) => eventName === 'disconnected', + )?.[1] as (() => void) | undefined; + + vi.useFakeTimers(); + disconnectedListener!(); + await vi.advanceTimersByTimeAsync( + 3000 + BRIDGE_SESSION_RESTORE_TIMEOUT_MS, + ); + await vi.waitFor(() => expect(mockRemoveServiceInfo).toHaveBeenCalled()); + + expect(mockBridgeStop).toHaveBeenCalledOnce(); + expect(mockWriteStderrLine).toHaveBeenCalledWith( + expect.stringContaining('Session restore timed out'), + ); + expect(exitSpy).toHaveBeenCalledWith(1); + } finally { + processOnSpy.mockRestore(); + exitSpy.mockRestore(); + vi.useRealTimers(); + } + }); + it('starts all channels with one shared bridge and router', async () => { const channels = { first: { type: 'telegram' }, diff --git a/packages/cli/src/commands/channel/start.ts b/packages/cli/src/commands/channel/start.ts index 4631b900969..952d8566c4c 100644 --- a/packages/cli/src/commands/channel/start.ts +++ b/packages/cli/src/commands/channel/start.ts @@ -55,6 +55,7 @@ export { resolveProxy } from './proxy.js'; const MAX_CRASH_RESTARTS = 3; const CRASH_WINDOW_MS = 5 * 60 * 1000; // 5-minute window for counting crashes const RESTART_DELAY_MS = 3000; +export const BRIDGE_SESSION_RESTORE_TIMEOUT_MS = 60 * 1000; function isFileExistsError(err: unknown): boolean { return ( @@ -152,6 +153,29 @@ function createBridgeReadinessGate(): { }; } +async function restoreBridgeSessions( + router: SessionRouter, +): ReturnType { + let timeout: ReturnType | undefined; + const expired = new Promise((_, reject) => { + timeout = setTimeout( + () => + reject( + new Error( + `Session restore timed out after ${BRIDGE_SESSION_RESTORE_TIMEOUT_MS}ms`, + ), + ), + BRIDGE_SESSION_RESTORE_TIMEOUT_MS, + ); + timeout.unref?.(); + }); + try { + return await Promise.race([router.restoreSessions(), expired]); + } finally { + clearTimeout(timeout); + } +} + interface BridgeRecoveryOptions { bridgeOpts: AcpBridgeOptions; router: SessionRouter; @@ -184,53 +208,44 @@ function createBridgeRecovery(options: BridgeRecoveryOptions): { const crashTimestamps: number[] = []; let recoveryTask: Promise | undefined; let recoveryRequested = false; + let recoverySourceBridge: AcpBridge | undefined; const attachDisconnectHandler = (failedBridge: AcpBridge): void => { failedBridge.on('disconnected', () => { - if (isShuttingDown()) return; + if (isShuttingDown() || failedBridge !== getBridge()) return; if (recoveryTask) { - if (failedBridge === getBridge()) recoveryRequested = true; + if (failedBridge !== recoverySourceBridge) recoveryRequested = true; return; } recoverBridge(); }); }; - const disconnectAllChannels = (): void => { - for (const channel of channels.values()) { - try { - channel.disconnect(); - } catch { - // best-effort - } - } - }; - const recoverBridge = (): void => { bridgeReadiness.block(); const task = (async () => { do { recoveryRequested = false; + recoverySourceBridge = getBridge(); const now = Date.now(); crashTimestamps.push(now); - // Only count crashes within the recent window - const recentCrashes = crashTimestamps.filter( - (timestamp) => now - timestamp < CRASH_WINDOW_MS, - ); + while (now - crashTimestamps[0]! >= CRASH_WINDOW_MS) { + crashTimestamps.shift(); + } + const recentCrashCount = crashTimestamps.length; - if (recentCrashes.length > MAX_CRASH_RESTARTS) { + if (recentCrashCount > MAX_CRASH_RESTARTS) { writeStderrLine( - `[Channel] Bridge crashed ${recentCrashes.length} times in ${CRASH_WINDOW_MS / 1000}s. Giving up.`, + `[Channel] Bridge crashed ${recentCrashCount} times in ${CRASH_WINDOW_MS / 1000}s. Giving up.`, ); scheduler?.stop(); - disconnectAllChannels(); - router.clearAll(); + cleanupStartedChannels(channels.values(), getBridge(), router); removeServiceInfo(); process.exit(1); } writeStderrLine( - `[Channel] Bridge crashed (${recentCrashes.length}/${MAX_CRASH_RESTARTS} in window). Restarting in ${RESTART_DELAY_MS / 1000}s...`, + `[Channel] Bridge crashed (${recentCrashCount}/${MAX_CRASH_RESTARTS} in window). Restarting in ${RESTART_DELAY_MS / 1000}s...`, ); await new Promise((resolve) => setTimeout(resolve, RESTART_DELAY_MS)); @@ -247,7 +262,7 @@ function createBridgeRecovery(options: BridgeRecoveryOptions): { registerPermissionRelay(bridge, router, channels); registerSessionCleanup(bridge, router, channels); - const result = await router.restoreSessions(); + const result = await restoreBridgeSessions(router); writeStdoutLine( `[Channel] Bridge restarted. Sessions restored: ${result.restored}, failed: ${result.failed}`, ); @@ -258,14 +273,14 @@ function createBridgeRecovery(options: BridgeRecoveryOptions): { `[Channel] Failed to restart bridge: ${err instanceof Error ? err.message : String(err)}`, ); scheduler?.stop(); - disconnectAllChannels(); - router.clearAll(); + cleanupStartedChannels(channels.values(), getBridge(), router); removeServiceInfo(); process.exit(1); }) .finally(() => { if (recoveryTask === task) { recoveryTask = undefined; + recoverySourceBridge = undefined; bridgeReadiness.release(); } }); @@ -329,6 +344,7 @@ async function startSingle( const cliEntryPath = findCliEntryPath(); let shuttingDown = false; + const bridgeReadiness = createBridgeReadinessGate(); const bridgeOpts = { cliEntryPath, cwd: config.cwd, model: config.model }; @@ -381,6 +397,7 @@ async function startSingle( writeServiceInfoOrExit([name], () => cleanupStartedChannels([channel], bridge, router), ); + // Keep scheduled loops active; their prompt paths wait on bridgeReadiness. scheduler?.start(); writeStdoutLine(`[Channel] "${name}" is running. Press Ctrl+C to stop.`); @@ -446,6 +463,7 @@ async function startAll( const cliEntryPath = findCliEntryPath(); const defaultCwd = process.cwd(); let shuttingDown = false; + const bridgeReadiness = createBridgeReadinessGate(); const bridgeOpts = { @@ -522,6 +540,7 @@ async function startAll( parsed.map((p) => p.name), () => cleanupStartedChannels(channels.values(), bridge, router), ); + // Keep scheduled loops active; their prompt paths wait on bridgeReadiness. scheduler?.start(); writeStdoutLine( `[Channel] Running ${connectedCount} channel(s). Press Ctrl+C to stop.`, diff --git a/packages/core/src/telemetry/event-loop-lag.test.ts b/packages/core/src/telemetry/event-loop-lag.test.ts index c57afbe05e4..7cac2746395 100644 --- a/packages/core/src/telemetry/event-loop-lag.test.ts +++ b/packages/core/src/telemetry/event-loop-lag.test.ts @@ -74,6 +74,29 @@ describe('startEventLoopLagMonitor', () => { monitor.dispose(); }); + it('reads snapshots without advancing suspension detection state', async () => { + const onNewMaxStall = vi.fn(); + histogram.max = 300_000_000_000; + const monitor = startEventLoopLagMonitor({ + resolutionMs: 10, + stallThresholdMs: 1_000, + suspendThresholdMs: 300_000, + onNewMaxStall, + }); + + vi.setSystemTime(Date.now() + 300_000); + expect(monitor.snapshot().maxMs).toBe(300_000); + expect(cpuUsage).toHaveBeenCalledOnce(); + expect(histogram.reset).not.toHaveBeenCalled(); + + await vi.advanceTimersByTimeAsync(10); + + expect(histogram.reset).toHaveBeenCalledOnce(); + expect(onNewMaxStall).not.toHaveBeenCalled(); + + monitor.dispose(); + }); + it('actively reports only new max stalls above threshold', async () => { const onNewMaxStall = vi.fn(); histogram.max = 15_000_000; @@ -126,6 +149,7 @@ describe('startEventLoopLagMonitor', () => { onNewMaxStall, }); + vi.setSystemTime(Date.now() + 10_000); await vi.advanceTimersByTimeAsync(10); expect(histogram.reset).toHaveBeenCalledTimes(1); @@ -134,7 +158,7 @@ describe('startEventLoopLagMonitor', () => { monitor.dispose(); }); - it('resets suspended samples without a stall callback', () => { + it('keeps snapshots pure while suspension filtering runs on the interval', async () => { histogram.mean = 15_000_000_000; histogram.max = 15_000_000_000; histogram.percentile.mockReturnValue(15_000_000_000); @@ -149,13 +173,22 @@ describe('startEventLoopLagMonitor', () => { }); vi.setSystemTime(Date.now() + 10_000); + expect(monitor.snapshot()).toEqual({ + meanMs: 15_000, + p50Ms: 15_000, + p99Ms: 15_000, + maxMs: 15_000, + }); + expect(histogram.reset).not.toHaveBeenCalled(); + + await vi.advanceTimersByTimeAsync(10); + expect(histogram.reset).toHaveBeenCalledOnce(); expect(monitor.snapshot()).toEqual({ meanMs: 0, p50Ms: 0, p99Ms: 0, maxMs: 0, }); - expect(histogram.reset).toHaveBeenCalledTimes(1); monitor.dispose(); }); @@ -173,6 +206,7 @@ describe('startEventLoopLagMonitor', () => { onNewMaxStall, }); + vi.setSystemTime(Date.now() + 10_000); await vi.advanceTimersByTimeAsync(10); histogram.max = 5_000_000_000; await vi.advanceTimersByTimeAsync(10); @@ -202,15 +236,17 @@ describe('startEventLoopLagMonitor', () => { monitor.dispose(); }); - it('reports a stall just below the default suspend threshold', async () => { + it('resets a low-CPU gap at the configured suspend threshold', async () => { const onNewMaxStall = vi.fn(); histogram.max = 299_000_000_000; const monitor = startEventLoopLagMonitor({ resolutionMs: 10, stallThresholdMs: 1_000, + suspendThresholdMs: 300_000, onNewMaxStall, }); + vi.setSystemTime(Date.now() + 300_000); await vi.advanceTimersByTimeAsync(10); expect(histogram.reset).not.toHaveBeenCalled(); @@ -228,6 +264,7 @@ describe('startEventLoopLagMonitor', () => { onNewMaxStall, }); + vi.setSystemTime(Date.now() + 600_000); await vi.advanceTimersByTimeAsync(10); expect(histogram.reset).toHaveBeenCalledTimes(1); @@ -277,6 +314,30 @@ describe('startEventLoopLagMonitor', () => { monitor.dispose(); }); + it('does not suppress an old histogram max after a short idle check', async () => { + const onNewMaxStall = vi.fn(); + cpuUsage + .mockReset() + .mockReturnValueOnce({ user: 0, system: 0 }) + .mockReturnValueOnce({ user: 20_000_000, system: 0 }) + .mockReturnValue({ user: 20_000_000, system: 0 }); + histogram.max = 600_000_000_000; + const monitor = startEventLoopLagMonitor({ + resolutionMs: 10, + stallThresholdMs: 1_000, + suspendThresholdMs: 300_000, + onNewMaxStall, + }); + + vi.setSystemTime(Date.now() + 600_000); + await vi.advanceTimersByTimeAsync(10); + await vi.advanceTimersByTimeAsync(10); + + expect(onNewMaxStall).toHaveBeenCalledOnce(); + expect(histogram.reset).not.toHaveBeenCalled(); + + monitor.dispose(); + }); it('enables and disables the underlying histogram', () => { const monitor = startEventLoopLagMonitor(); diff --git a/packages/core/src/telemetry/event-loop-lag.ts b/packages/core/src/telemetry/event-loop-lag.ts index 331f13ee609..64863b2be22 100644 --- a/packages/core/src/telemetry/event-loop-lag.ts +++ b/packages/core/src/telemetry/event-loop-lag.ts @@ -67,6 +67,7 @@ export function startEventLoopLagMonitor( let disposed = false; let lastReportedMaxMs = 0; + let lastObservedMaxMs = 0; let lastCheckTimeMs = Date.now(); let lastCpuUsage = safeCpuUsage(); const readMaxMs = () => nsToMs(histogram.max); @@ -76,15 +77,19 @@ export function startEventLoopLagMonitor( const cpuUsage = safeCpuUsage(); const elapsedMs = Math.max(0, nowMs - lastCheckTimeMs); const maxMs = readMaxMs(); + const newMaxMs = maxMs > lastObservedMaxMs ? maxMs : 0; const cpuRatio = calculateCpuRatio(lastCpuUsage, cpuUsage, elapsedMs); lastCheckTimeMs = nowMs; if (cpuUsage) lastCpuUsage = cpuUsage; + lastObservedMaxMs = maxMs; if ( - maxMs >= suspendThresholdMs && + newMaxMs >= suspendThresholdMs && + elapsedMs >= suspendThresholdMs && cpuRatio !== undefined && cpuRatio <= suspendCpuRatio ) { histogram.reset(); + lastObservedMaxMs = 0; lastReportedMaxMs = 0; return; } @@ -106,7 +111,6 @@ export function startEventLoopLagMonitor( return { snapshot(): EventLoopLagSnapshot { - checkHistogram(); return { meanMs: nsToMs(histogram.mean), p50Ms: nsToMs(histogram.percentile(50)), From ec51dfc8ad0d67b0b222e6f388523f5e7e671599 Mon Sep 17 00:00:00 2001 From: yiliang114 Date: Fri, 31 Jul 2026 21:35:17 +0800 Subject: [PATCH 6/9] test(channel): align wake recovery threshold checks --- .../cli/src/acp-integration/acpAgent.test.ts | 5 ++++- .../acp-integration/stall-thresholds.test.ts | 22 ------------------- .../core/src/telemetry/event-loop-lag.test.ts | 2 +- 3 files changed, 5 insertions(+), 24 deletions(-) delete mode 100644 packages/cli/src/acp-integration/stall-thresholds.test.ts diff --git a/packages/cli/src/acp-integration/acpAgent.test.ts b/packages/cli/src/acp-integration/acpAgent.test.ts index bd526930818..2aa7bdd33f9 100644 --- a/packages/cli/src/acp-integration/acpAgent.test.ts +++ b/packages/cli/src/acp-integration/acpAgent.test.ts @@ -18,6 +18,7 @@ import * as fs from 'node:fs/promises'; import * as os from 'node:os'; import * as path from 'node:path'; import { NOT_CURRENTLY_GENERATING_CANCEL_MESSAGE } from '@qwen-code/acp-bridge/bridgeErrors'; +import { ACP_EVENT_LOOP_STALL_RESTART_MS } from '@qwen-code/channel-base'; // Mock cleanup module before importing anything else const { mockRunExitCleanup } = vi.hoisted(() => ({ @@ -1399,7 +1400,9 @@ describe('runAcpAgent shutdown cleanup', () => { expect(snapshot).toHaveBeenCalledTimes(1); expect(startEventLoopLagMonitor).toHaveBeenCalledWith( - expect.objectContaining({ suspendThresholdMs: 5 * 60 * 1000 }), + expect.objectContaining({ + suspendThresholdMs: ACP_EVENT_LOOP_STALL_RESTART_MS, + }), ); mockConnectionState.resolve(); diff --git a/packages/cli/src/acp-integration/stall-thresholds.test.ts b/packages/cli/src/acp-integration/stall-thresholds.test.ts deleted file mode 100644 index bad3917677b..00000000000 --- a/packages/cli/src/acp-integration/stall-thresholds.test.ts +++ /dev/null @@ -1,22 +0,0 @@ -/** - * @license - * Copyright 2026 Qwen Team - * SPDX-License-Identifier: Apache-2.0 - */ - -import { describe, expect, it } from 'vitest'; -import { ACP_EVENT_LOOP_STALL_RESTART_MS } from '@qwen-code/channel-base'; -import { DEFAULT_EVENT_LOOP_SUSPEND_THRESHOLD_MS } from '@qwen-code/qwen-code-core'; - -describe('acp stall thresholds', () => { - // The ACP agent starts its event-loop lag monitor with default options - // (acpAgent.ts), while AcpBridge kills the child once a reported stall - // reaches ACP_EVENT_LOOP_STALL_RESTART_MS. If the suspend threshold ever - // exceeded the kill threshold, a host sleep in between the two values would - // be reported as a stall and kill a healthy child on wake. - it('keeps host-suspension filtering at or below the bridge stall-kill threshold', () => { - expect(DEFAULT_EVENT_LOOP_SUSPEND_THRESHOLD_MS).toBeLessThanOrEqual( - ACP_EVENT_LOOP_STALL_RESTART_MS, - ); - }); -}); diff --git a/packages/core/src/telemetry/event-loop-lag.test.ts b/packages/core/src/telemetry/event-loop-lag.test.ts index 7cac2746395..f4c9e86009e 100644 --- a/packages/core/src/telemetry/event-loop-lag.test.ts +++ b/packages/core/src/telemetry/event-loop-lag.test.ts @@ -236,7 +236,7 @@ describe('startEventLoopLagMonitor', () => { monitor.dispose(); }); - it('resets a low-CPU gap at the configured suspend threshold', async () => { + it('reports a low-CPU gap just below the configured suspend threshold', async () => { const onNewMaxStall = vi.fn(); histogram.max = 299_000_000_000; const monitor = startEventLoopLagMonitor({ From b8baca2a4b77413bf173a52989a94b2f61b50441 Mon Sep 17 00:00:00 2001 From: yiliang114 Date: Fri, 31 Jul 2026 14:33:22 +0000 Subject: [PATCH 7/9] fix(channel): clear recovery-aborted loop runs and strengthen gate tests (#8211) --- .../channels/base/src/ChannelBase.test.ts | 44 ++++++++++++++----- .../base/src/ChannelLoopScheduler.test.ts | 33 ++++++++++++++ .../channels/base/src/ChannelLoopScheduler.ts | 17 ++++++- .../cli/src/commands/channel/start.test.ts | 2 + packages/cli/src/commands/channel/start.ts | 1 + 5 files changed, 86 insertions(+), 11 deletions(-) diff --git a/packages/channels/base/src/ChannelBase.test.ts b/packages/channels/base/src/ChannelBase.test.ts index 67c2563fd22..213d7e4c611 100644 --- a/packages/channels/base/src/ChannelBase.test.ts +++ b/packages/channels/base/src/ChannelBase.test.ts @@ -16065,9 +16065,15 @@ describe('ChannelBase', () => { ch.proactiveSupported = true; const run = ch.runWebhookTask(webhookTask); - await Promise.resolve(); + // Absent the entry gate the (fully mocked) flow reaches bridge.prompt + // quickly; wait long enough that a missing gate would be caught. + await expect( + vi.waitFor(() => expect(bridge.prompt).toHaveBeenCalled(), { + timeout: 500, + interval: 25, + }), + ).rejects.toThrow(); expect(bridge.newSession).not.toHaveBeenCalled(); - expect(bridge.prompt).not.toHaveBeenCalled(); releaseBridge!(); await expect(run).resolves.toBe('agent response'); @@ -16098,9 +16104,15 @@ describe('ChannelBase', () => { releaseRecovery = resolve; }); releaseMemoryRead!(); - await Promise.resolve(); - - expect(bridge.prompt).not.toHaveBeenCalled(); + // Absent the recheck gate the (fully mocked) flow reaches bridge.prompt + // quickly; wait long enough that a missing gate would be caught, then + // confirm the gate held the prompt back until recovery resolved. + await expect( + vi.waitFor(() => expect(bridge.prompt).toHaveBeenCalled(), { + timeout: 500, + interval: 25, + }), + ).rejects.toThrow(); releaseRecovery!(); await expect(run).resolves.toBe('agent response'); @@ -16135,9 +16147,15 @@ describe('ChannelBase', () => { const ch = createChannel({}, { bridgeRecovery: () => bridgeReady }); const inbound = ch.processAfterAdapterPreflight(envelope()); - await Promise.resolve(); + // Absent the gate the (fully mocked) flow reaches bridge.prompt quickly; + // wait long enough that a missing gate would be caught. + await expect( + vi.waitFor(() => expect(bridge.prompt).toHaveBeenCalled(), { + timeout: 500, + interval: 25, + }), + ).rejects.toThrow(); expect(bridge.newSession).not.toHaveBeenCalled(); - expect(bridge.prompt).not.toHaveBeenCalled(); releaseBridge!(); await inbound; @@ -16288,9 +16306,15 @@ describe('ChannelBase', () => { releaseRecovery = resolve; }); releaseMemoryRead!(); - await Promise.resolve(); - - expect(bridge.prompt).not.toHaveBeenCalled(); + // Absent the recheck gate the (fully mocked) flow reaches bridge.prompt + // quickly; wait long enough that a missing gate would be caught, then + // confirm the gate held the prompt back until recovery resolved. + await expect( + vi.waitFor(() => expect(bridge.prompt).toHaveBeenCalled(), { + timeout: 500, + interval: 25, + }), + ).rejects.toThrow(); releaseRecovery!(); await expect(loopRun).resolves.toBe('agent response'); diff --git a/packages/channels/base/src/ChannelLoopScheduler.test.ts b/packages/channels/base/src/ChannelLoopScheduler.test.ts index 51dc37aa34c..2bb111cf306 100644 --- a/packages/channels/base/src/ChannelLoopScheduler.test.ts +++ b/packages/channels/base/src/ChannelLoopScheduler.test.ts @@ -366,6 +366,39 @@ describe('ChannelLoopScheduler', () => { writeSpy.mockRestore(); }); + it('clears an in-flight loop aborted by bridge recovery instead of failing it', async () => { + jobs = [{ ...baseJob, consecutiveFailures: 4 }]; + const scheduler = new ChannelLoopScheduler({ + store, + channels: new Map([['feishu-main', { runLoopPrompt }]]), + now: () => new Date(nowMs), + nextFireTime: () => new Date(nowMs - 60_000), + maxConsecutiveFailures: 5, + }); + runLoopPrompt.mockImplementation(async () => { + // The bridge is replaced mid-prompt; recovery marks the scheduler before + // the in-flight prompt rejects. + scheduler.markBridgeRecovery(); + throw new Error('bridge replaced during prompt'); + }); + + await scheduler.tick(); + + await vi.waitFor(() => { + expect(store.update).toHaveBeenCalledWith('job-1', { + runningSince: undefined, + }); + }); + const failureWrites = ( + store.update as ReturnType + ).mock.calls.filter( + ([, patch]) => (patch as { lastStatus?: string }).lastStatus === 'error', + ); + expect(failureWrites).toHaveLength(0); + expect(jobs[0]?.consecutiveFailures).toBe(4); + expect(jobs[0]?.enabled).toBe(true); + }); + it('disables one-shot loops after a failed attempt', async () => { jobs = [{ ...baseJob, recurring: false }]; runLoopPrompt.mockRejectedValue(new Error('cannot cold send')); diff --git a/packages/channels/base/src/ChannelLoopScheduler.ts b/packages/channels/base/src/ChannelLoopScheduler.ts index 091c819e3bb..c6d589a58be 100644 --- a/packages/channels/base/src/ChannelLoopScheduler.ts +++ b/packages/channels/base/src/ChannelLoopScheduler.ts @@ -48,6 +48,7 @@ export class ChannelLoopScheduler { private runningTick: Promise | undefined; private readonly inFlightJobs = new Map(); private generation = 0; + private recoveryEpoch = 0; constructor(options: ChannelLoopSchedulerOptions) { this.store = options.store; @@ -89,6 +90,15 @@ export class ChannelLoopScheduler { this.inFlightJobs.clear(); } + /** + * Mark that a bridge recovery started. Loop prompts already in flight that + * the bridge replacement aborts are cleared instead of counted as agent + * failures, so a crash-restart cannot auto-disable a loop. + */ + markBridgeRecovery(): void { + this.recoveryEpoch++; + } + private async reconcileStartupState(): Promise { const jobs = await this.store.list(); const staleRunning = jobs.filter((job) => job.runningSince); @@ -177,6 +187,7 @@ export class ChannelLoopScheduler { const latestJob = await this.findJob(job.id); if (!latestJob?.enabled) return; if (this.generation !== generation) return; + const recoveryEpoch = this.recoveryEpoch; const runningSince = now.toISOString(); let resultPreview: string | undefined; @@ -214,7 +225,11 @@ export class ChannelLoopScheduler { await this.clearRunningSince(latestJob.id, runningSince); return; } - if (this.generation !== generation || !currentJob?.enabled) { + if ( + this.generation !== generation || + recoveryEpoch !== this.recoveryEpoch || + !currentJob?.enabled + ) { await this.clearRunningSince(latestJob.id, runningSince); return; } diff --git a/packages/cli/src/commands/channel/start.test.ts b/packages/cli/src/commands/channel/start.test.ts index ed7b20c5b4d..751648f0b22 100644 --- a/packages/cli/src/commands/channel/start.test.ts +++ b/packages/cli/src/commands/channel/start.test.ts @@ -83,10 +83,12 @@ const mockChannelLoopStore = vi.hoisted(() => ); const mockChannelLoopSchedulerStart = vi.hoisted(() => vi.fn()); const mockChannelLoopSchedulerStop = vi.hoisted(() => vi.fn()); +const mockChannelLoopSchedulerMarkRecovery = vi.hoisted(() => vi.fn()); const mockChannelLoopScheduler = vi.hoisted(() => vi.fn((_options?: unknown) => ({ start: mockChannelLoopSchedulerStart, stop: mockChannelLoopSchedulerStop, + markBridgeRecovery: mockChannelLoopSchedulerMarkRecovery, })), ); const mockSessionRouter = vi.hoisted(() => diff --git a/packages/cli/src/commands/channel/start.ts b/packages/cli/src/commands/channel/start.ts index 952d8566c4c..47c8629a648 100644 --- a/packages/cli/src/commands/channel/start.ts +++ b/packages/cli/src/commands/channel/start.ts @@ -223,6 +223,7 @@ function createBridgeRecovery(options: BridgeRecoveryOptions): { const recoverBridge = (): void => { bridgeReadiness.block(); + scheduler?.markBridgeRecovery(); const task = (async () => { do { recoveryRequested = false; From 3ed3a9f3f72d1bdee39e31c2050290d712780694 Mon Sep 17 00:00:00 2001 From: qwen-code-bot Date: Fri, 31 Jul 2026 16:54:05 +0000 Subject: [PATCH 8/9] fix(channel): scope recovery suppression to prompts actually aborted (#8211) --- packages/channels/base/src/AcpBridge.test.ts | 29 ++++++++------- .../base/src/ChannelLoopScheduler.test.ts | 36 +++++++++++++++++++ .../channels/base/src/ChannelLoopScheduler.ts | 3 +- 3 files changed, 54 insertions(+), 14 deletions(-) diff --git a/packages/channels/base/src/AcpBridge.test.ts b/packages/channels/base/src/AcpBridge.test.ts index 1967b9cb3f3..9e33aca4002 100644 --- a/packages/channels/base/src/AcpBridge.test.ts +++ b/packages/channels/base/src/AcpBridge.test.ts @@ -145,21 +145,24 @@ describe('AcpBridge', () => { it('times out bridge initialization and stops the child', async () => { vi.useFakeTimers(); - child.setInitializeImplementation(() => new Promise(() => {})); - const bridge = new AcpBridge({ - cliEntryPath: '/tmp/qwen', - cwd: '/tmp', - }); + try { + child.setInitializeImplementation(() => new Promise(() => {})); + const bridge = new AcpBridge({ + cliEntryPath: '/tmp/qwen', + cwd: '/tmp', + }); - const start = bridge.start(); - const rejection = expect(start).rejects.toThrow( - `ACP initialization timed out after ${ACP_START_TIMEOUT_MS}ms`, - ); - await vi.advanceTimersByTimeAsync(1000 + ACP_START_TIMEOUT_MS); + const start = bridge.start(); + const rejection = expect(start).rejects.toThrow( + `ACP initialization timed out after ${ACP_START_TIMEOUT_MS}ms`, + ); + await vi.advanceTimersByTimeAsync(1000 + ACP_START_TIMEOUT_MS); - await rejection; - expect(child.instances[0]!.kill).toHaveBeenCalledOnce(); - vi.useRealTimers(); + await rejection; + expect(child.instances[0]!.kill).toHaveBeenCalledOnce(); + } finally { + vi.useRealTimers(); + } }); it('registers the channel loop MCP server once across concurrent calls', async () => { diff --git a/packages/channels/base/src/ChannelLoopScheduler.test.ts b/packages/channels/base/src/ChannelLoopScheduler.test.ts index 2bb111cf306..82aa7417582 100644 --- a/packages/channels/base/src/ChannelLoopScheduler.test.ts +++ b/packages/channels/base/src/ChannelLoopScheduler.test.ts @@ -399,6 +399,42 @@ describe('ChannelLoopScheduler', () => { expect(jobs[0]?.enabled).toBe(true); }); + it('records a genuine failure when recovery starts before the prompt runs', async () => { + let releaseUpdate!: () => void; + const updateBlocked = new Promise((r) => (releaseUpdate = r)); + store.update = vi.fn(async (id, patch) => { + if ((patch as { runningSince?: string }).runningSince) { + await updateBlocked; + } + jobs = jobs.map((job) => (job.id === id ? { ...job, ...patch } : job)); + return true; + }); + runLoopPrompt.mockRejectedValue(new Error('genuine failure on new bridge')); + const scheduler = new ChannelLoopScheduler({ + store, + channels: new Map([['feishu-main', { runLoopPrompt }]]), + now: () => new Date(nowMs), + nextFireTime: () => new Date(nowMs - 60_000), + }); + + const tick = scheduler.tick(); + await vi.waitFor(() => expect(store.update).toHaveBeenCalledOnce()); + scheduler.markBridgeRecovery(); + releaseUpdate(); + await tick; + + await vi.waitFor(() => { + expect(store.update).toHaveBeenCalledWith( + 'job-1', + expect.objectContaining({ + lastStatus: 'error', + lastError: 'genuine failure on new bridge', + consecutiveFailures: 1, + }), + ); + }); + }); + it('disables one-shot loops after a failed attempt', async () => { jobs = [{ ...baseJob, recurring: false }]; runLoopPrompt.mockRejectedValue(new Error('cannot cold send')); diff --git a/packages/channels/base/src/ChannelLoopScheduler.ts b/packages/channels/base/src/ChannelLoopScheduler.ts index c6d589a58be..2faf71b4f5d 100644 --- a/packages/channels/base/src/ChannelLoopScheduler.ts +++ b/packages/channels/base/src/ChannelLoopScheduler.ts @@ -187,7 +187,7 @@ export class ChannelLoopScheduler { const latestJob = await this.findJob(job.id); if (!latestJob?.enabled) return; if (this.generation !== generation) return; - const recoveryEpoch = this.recoveryEpoch; + let recoveryEpoch = this.recoveryEpoch; const runningSince = now.toISOString(); let resultPreview: string | undefined; @@ -200,6 +200,7 @@ export class ChannelLoopScheduler { await this.clearRunningSince(latestJob.id, runningSince); return; } + recoveryEpoch = this.recoveryEpoch; resultPreview = await channel.runLoopPrompt(latestJob, { timeoutMs: this.loopTimeoutMs, shouldContinue: async () => { From e70d198d58defbf1f03be1e25d1a50f432aaa75a Mon Sep 17 00:00:00 2001 From: Qwen Code Bot Date: Fri, 31 Jul 2026 19:11:58 +0000 Subject: [PATCH 9/9] fix(channel): carry low-CPU gap across ticks and address review feedback (#8211) --- .../channels/base/src/ChannelBase.test.ts | 11 ++++---- packages/channels/base/src/ChannelBase.ts | 4 +-- .../base/src/ChannelLoopScheduler.test.ts | 19 +++++++++----- .../channels/base/src/ChannelLoopScheduler.ts | 21 +++++++++++---- packages/cli/src/acp-integration/acpAgent.ts | 7 +++-- .../cli/src/commands/channel/start.test.ts | 1 + .../core/src/telemetry/event-loop-lag.test.ts | 26 +++++++++++++++++++ packages/core/src/telemetry/event-loop-lag.ts | 15 ++++++++--- packages/core/src/telemetry/index.ts | 1 - 9 files changed, 79 insertions(+), 26 deletions(-) diff --git a/packages/channels/base/src/ChannelBase.test.ts b/packages/channels/base/src/ChannelBase.test.ts index 213d7e4c611..b2114f59763 100644 --- a/packages/channels/base/src/ChannelBase.test.ts +++ b/packages/channels/base/src/ChannelBase.test.ts @@ -16212,15 +16212,16 @@ describe('ChannelBase', () => { ); const inbound = ch.handleInbound(envelope()); - await vi.waitFor(() => expect(bridge.newSession).not.toHaveBeenCalled()); recoveryState.current = new Promise((resolve) => { releaseRecovery = resolve; }); releaseContactRecord!(); - await Promise.resolve(); - - expect(bridge.newSession).not.toHaveBeenCalled(); - expect(bridge.prompt).not.toHaveBeenCalled(); + await expect( + vi.waitFor(() => expect(bridge.prompt).toHaveBeenCalled(), { + timeout: 500, + interval: 25, + }), + ).rejects.toThrow(); releaseRecovery!(); await inbound; diff --git a/packages/channels/base/src/ChannelBase.ts b/packages/channels/base/src/ChannelBase.ts index 7985826ccab..d7a18f5b62c 100644 --- a/packages/channels/base/src/ChannelBase.ts +++ b/packages/channels/base/src/ChannelBase.ts @@ -351,8 +351,8 @@ function isUnattendedWebhookApprovalMode(mode: string | undefined): boolean { export abstract class ChannelBase { protected config: ChannelConfig; /** - * Recovery invariant: every path that resolves a session or calls the bridge - * must await waitForBridgeRecovery() immediately before that operation. + * Recovery invariant: session-resolution and prompt-capture paths must await + * waitForBridgeRecovery() immediately before that operation. */ protected bridge: ChannelAgentBridge; protected groupGate: GroupGate; diff --git a/packages/channels/base/src/ChannelLoopScheduler.test.ts b/packages/channels/base/src/ChannelLoopScheduler.test.ts index 82aa7417582..9d7a5bff1cc 100644 --- a/packages/channels/base/src/ChannelLoopScheduler.test.ts +++ b/packages/channels/base/src/ChannelLoopScheduler.test.ts @@ -385,16 +385,23 @@ describe('ChannelLoopScheduler', () => { await scheduler.tick(); await vi.waitFor(() => { - expect(store.update).toHaveBeenCalledWith('job-1', { - runningSince: undefined, - }); + expect(store.update).toHaveBeenCalledWith( + 'job-1', + expect.objectContaining({ + lastStatus: 'error', + lastError: 'bridge replaced during prompt', + runningSince: undefined, + }), + ); }); - const failureWrites = ( + const failureCountWrites = ( store.update as ReturnType ).mock.calls.filter( - ([, patch]) => (patch as { lastStatus?: string }).lastStatus === 'error', + ([, patch]) => + (patch as { consecutiveFailures?: number }).consecutiveFailures !== + undefined, ); - expect(failureWrites).toHaveLength(0); + expect(failureCountWrites).toHaveLength(0); expect(jobs[0]?.consecutiveFailures).toBe(4); expect(jobs[0]?.enabled).toBe(true); }); diff --git a/packages/channels/base/src/ChannelLoopScheduler.ts b/packages/channels/base/src/ChannelLoopScheduler.ts index 2faf71b4f5d..c8c0ab9a533 100644 --- a/packages/channels/base/src/ChannelLoopScheduler.ts +++ b/packages/channels/base/src/ChannelLoopScheduler.ts @@ -226,14 +226,25 @@ export class ChannelLoopScheduler { await this.clearRunningSince(latestJob.id, runningSince); return; } - if ( - this.generation !== generation || - recoveryEpoch !== this.recoveryEpoch || - !currentJob?.enabled - ) { + if (this.generation !== generation || !currentJob?.enabled) { await this.clearRunningSince(latestJob.id, runningSince); return; } + if (recoveryEpoch !== this.recoveryEpoch) { + try { + await this.store.update(latestJob.id, { + lastFinishedAt: this.now().toISOString(), + lastStatus: 'error', + lastError: truncateError( + err instanceof Error ? err.message : String(err), + ), + runningSince: undefined, + }); + } catch { + await this.clearRunningSince(latestJob.id, runningSince); + } + return; + } await this.recordFailure( currentJob, now, diff --git a/packages/cli/src/acp-integration/acpAgent.ts b/packages/cli/src/acp-integration/acpAgent.ts index df551a7ab4a..eed5acded8d 100644 --- a/packages/cli/src/acp-integration/acpAgent.ts +++ b/packages/cli/src/acp-integration/acpAgent.ts @@ -162,6 +162,7 @@ import { } from './authMethods.js'; import { AcpFileSystemService } from './service/filesystem.js'; import { ndJsonStream } from '@qwen-code/acp-bridge/ndJsonStream'; +import { ACP_EVENT_LOOP_STALL_RESTART_MS } from '@qwen-code/channel-base'; import { Readable, Writable } from 'node:stream'; import { normalizeDisabledToolList } from '../config/normalizeDisabledTools.js'; import { pipeline } from 'node:stream/promises'; @@ -345,9 +346,7 @@ const POSIX_TMP_LOCAL_READ_ROOT = '/tmp'; const BTW_CHILD_TIMEOUT_MS = 55_000; const MCP_OAUTH_START_TIMEOUT_MS = 30_000; const SESSION_DRAIN_TIMEOUT_MS = 30_000; -// Match the parent channel bridge's restart threshold. Low-CPU gaps at this -// boundary are host suspension; active stalls still reach the parent watchdog. -const ACP_EVENT_LOOP_SUSPEND_THRESHOLD_MS = 5 * 60 * 1000; + // Must be less than WORKSPACE_MEMORY_REMEMBER_TIMEOUT_MS (300s) in bridge.ts. const WORKSPACE_MEMORY_REMEMBER_CHILD_TIMEOUT_MS = 295_000; @@ -2879,7 +2878,7 @@ export async function runAcpAgent( process.stderr.write(`${warning}\n`); } const eventLoopMonitor = startEventLoopLagMonitor({ - suspendThresholdMs: ACP_EVENT_LOOP_SUSPEND_THRESHOLD_MS, + suspendThresholdMs: ACP_EVENT_LOOP_STALL_RESTART_MS, onNewMaxStall: (maxMs) => { console.error(`[perf] acp agent event loop stall: max=${maxMs}ms`); }, diff --git a/packages/cli/src/commands/channel/start.test.ts b/packages/cli/src/commands/channel/start.test.ts index 751648f0b22..d40988baecc 100644 --- a/packages/cli/src/commands/channel/start.test.ts +++ b/packages/cli/src/commands/channel/start.test.ts @@ -829,6 +829,7 @@ describe('startCommand.handler', () => { expect(mockChannelConnect).toHaveBeenCalledTimes(1); expect(mockChannelDisconnect).not.toHaveBeenCalled(); + expect(mockChannelLoopSchedulerMarkRecovery).toHaveBeenCalled(); expect(mockChannelLoopSchedulerStop).not.toHaveBeenCalled(); } finally { processOnSpy.mockRestore(); diff --git a/packages/core/src/telemetry/event-loop-lag.test.ts b/packages/core/src/telemetry/event-loop-lag.test.ts index f4c9e86009e..e457f42be37 100644 --- a/packages/core/src/telemetry/event-loop-lag.test.ts +++ b/packages/core/src/telemetry/event-loop-lag.test.ts @@ -338,6 +338,32 @@ describe('startEventLoopLagMonitor', () => { monitor.dispose(); }); + it('suppresses a suspension gap whose histogram max lands one tick late', async () => { + const onNewMaxStall = vi.fn(); + const monitor = startEventLoopLagMonitor({ + resolutionMs: 10, + stallThresholdMs: 1_000, + suspendThresholdMs: 10_000, + onNewMaxStall, + }); + + vi.setSystemTime(Date.now() + 10_000); + // Tick 1: our interval sees the 10 s gap but the histogram has not + // published the matching max yet. + await vi.advanceTimersByTimeAsync(10); + expect(histogram.reset).not.toHaveBeenCalled(); + + // The histogram publishes the gap between ticks. + histogram.max = 10_000_000_000; + // Tick 2: the carried gap qualifies the late histogram max. + await vi.advanceTimersByTimeAsync(10); + + expect(histogram.reset).toHaveBeenCalledTimes(1); + expect(onNewMaxStall).not.toHaveBeenCalled(); + + monitor.dispose(); + }); + it('enables and disables the underlying histogram', () => { const monitor = startEventLoopLagMonitor(); diff --git a/packages/core/src/telemetry/event-loop-lag.ts b/packages/core/src/telemetry/event-loop-lag.ts index 64863b2be22..7b79e19c75a 100644 --- a/packages/core/src/telemetry/event-loop-lag.ts +++ b/packages/core/src/telemetry/event-loop-lag.ts @@ -70,6 +70,7 @@ export function startEventLoopLagMonitor( let lastObservedMaxMs = 0; let lastCheckTimeMs = Date.now(); let lastCpuUsage = safeCpuUsage(); + let pendingSuspendGapMs = 0; const readMaxMs = () => nsToMs(histogram.max); const checkHistogram = () => { if (disposed) return; @@ -82,15 +83,23 @@ export function startEventLoopLagMonitor( lastCheckTimeMs = nowMs; if (cpuUsage) lastCpuUsage = cpuUsage; lastObservedMaxMs = maxMs; - if ( - newMaxMs >= suspendThresholdMs && + const isLowCpuGap = elapsedMs >= suspendThresholdMs && cpuRatio !== undefined && - cpuRatio <= suspendCpuRatio + cpuRatio <= suspendCpuRatio; + // The histogram's own libuv timer can land the gap one tick after ours, + // so a low-CPU gap stays eligible for one further check. + const suspendGapMs = isLowCpuGap ? elapsedMs : pendingSuspendGapMs; + pendingSuspendGapMs = isLowCpuGap ? elapsedMs : 0; + if ( + newMaxMs >= suspendThresholdMs && + suspendGapMs >= suspendThresholdMs && + newMaxMs <= suspendGapMs * 1.5 ) { histogram.reset(); lastObservedMaxMs = 0; lastReportedMaxMs = 0; + pendingSuspendGapMs = 0; return; } if ( diff --git a/packages/core/src/telemetry/index.ts b/packages/core/src/telemetry/index.ts index 11f2f140eb2..0cea7f1d96f 100644 --- a/packages/core/src/telemetry/index.ts +++ b/packages/core/src/telemetry/index.ts @@ -225,7 +225,6 @@ export type { DaemonPipeDirection, } from './daemon-metrics.js'; export { - DEFAULT_EVENT_LOOP_SUSPEND_THRESHOLD_MS, startEventLoopLagMonitor, type EventLoopLagMonitor, type EventLoopLagMonitorOptions,