diff --git a/packages/channels/base/src/AcpBridge.test.ts b/packages/channels/base/src/AcpBridge.test.ts index 2312d100187..9e33aca4002 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,29 @@ 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(); + 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); + + 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/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 ef9d4be4358..b2114f59763 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'); @@ -14961,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) => { @@ -15910,7 +15950,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 +16052,275 @@ describe('ChannelBase', () => { .calls[1][1] as string; 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); + // 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(); + + 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; + 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!(); + // 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'); + + 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()); + // 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(); + + releaseBridge!(); + await inbound; + + 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; + let releaseContactRecord: (() => void) | undefined; + const contactRecord = new Promise((resolve) => { + releaseContactRecord = resolve; + }); + const ch = createChannel( + {}, + { + bridgeRecovery: () => recoveryState.current, + observedContacts: { + observe: vi.fn().mockImplementation(() => contactRecord), + }, + }, + ); + + const inbound = ch.handleInbound(envelope()); + recoveryState.current = new Promise((resolve) => { + releaseRecovery = resolve; + }); + releaseContactRecord!(); + await expect( + vi.waitFor(() => expect(bridge.prompt).toHaveBeenCalled(), { + timeout: 500, + interval: 25, + }), + ).rejects.toThrow(); + + 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!(); + // 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'); + + expect(bridge.prompt).toHaveBeenCalledTimes(1); }); it('runs a loop prompt as a follow-up and pushes the result proactively', async () => { @@ -17253,7 +17562,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..d7a18f5b62c 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?: { @@ -348,6 +350,10 @@ function isUnattendedWebhookApprovalMode(mode: string | undefined): boolean { export abstract class ChannelBase { protected config: ChannelConfig; + /** + * Recovery invariant: session-resolution and prompt-capture paths must await + * waitForBridgeRecovery() immediately before that operation. + */ protected bridge: ChannelAgentBridge; protected groupGate: GroupGate; protected dmGate: DmGate; @@ -379,6 +385,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 +811,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 +1469,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 +1620,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 +1788,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 +1917,7 @@ export abstract class ChannelBase { releaseHeldChunks(); } }; + await this.waitForBridgeRecovery(); const promptBridge = this.bridge; promptBridge.on('textChunk', onChunk); @@ -4904,6 +4916,17 @@ export abstract class ChannelBase { this.preflightedEnvelopes.add(envelope); } + /** Wait until the currently active bridge recovery, if any, has completed. */ + private async waitForBridgeRecovery(): Promise { + let completedRecovery: Promise | undefined; + while (true) { + const bridgeRecovery = this.bridgeRecovery?.(); + if (!bridgeRecovery || bridgeRecovery === completedRecovery) return; + await bridgeRecovery; + completedRecovery = bridgeRecovery; + } + } + /** * Process an inbound message after preflight gates have passed. * @@ -4912,6 +4935,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 +5026,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 +5514,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/channels/base/src/ChannelLoopScheduler.test.ts b/packages/channels/base/src/ChannelLoopScheduler.test.ts index 51dc37aa34c..9d7a5bff1cc 100644 --- a/packages/channels/base/src/ChannelLoopScheduler.test.ts +++ b/packages/channels/base/src/ChannelLoopScheduler.test.ts @@ -366,6 +366,82 @@ 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', + expect.objectContaining({ + lastStatus: 'error', + lastError: 'bridge replaced during prompt', + runningSince: undefined, + }), + ); + }); + const failureCountWrites = ( + store.update as ReturnType + ).mock.calls.filter( + ([, patch]) => + (patch as { consecutiveFailures?: number }).consecutiveFailures !== + undefined, + ); + expect(failureCountWrites).toHaveLength(0); + expect(jobs[0]?.consecutiveFailures).toBe(4); + 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 091c819e3bb..c8c0ab9a533 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; + let recoveryEpoch = this.recoveryEpoch; const runningSince = now.toISOString(); let resultPreview: string | undefined; @@ -189,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 () => { @@ -218,6 +230,21 @@ export class ChannelLoopScheduler { 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/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/acpAgent.test.ts b/packages/cli/src/acp-integration/acpAgent.test.ts index fd967bde581..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(() => ({ @@ -1398,6 +1399,12 @@ describe('runAcpAgent shutdown cleanup', () => { }); expect(snapshot).toHaveBeenCalledTimes(1); + expect(startEventLoopLagMonitor).toHaveBeenCalledWith( + expect.objectContaining({ + suspendThresholdMs: ACP_EVENT_LOOP_STALL_RESTART_MS, + }), + ); + mockConnectionState.resolve(); await agentPromise; diff --git a/packages/cli/src/acp-integration/acpAgent.ts b/packages/cli/src/acp-integration/acpAgent.ts index bdfc381826e..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,6 +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; + // Must be less than WORKSPACE_MEMORY_REMEMBER_TIMEOUT_MS (300s) in bridge.ts. const WORKSPACE_MEMORY_REMEMBER_CHILD_TIMEOUT_MS = 295_000; @@ -2876,6 +2878,7 @@ export async function runAcpAgent( process.stderr.write(`${warning}\n`); } const eventLoopMonitor = startEventLoopLagMonitor({ + 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 3673184baa5..d40988baecc 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(() => @@ -160,6 +162,7 @@ vi.mock('@qwen-code/channel-base', () => ({ })); import { + BRIDGE_SESSION_RESTORE_TIMEOUT_MS, resolveExtensionChannelEntrySpecifier, resolveProxy, startCommand, @@ -775,7 +778,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 +790,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 +805,260 @@ 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(mockChannelLoopSchedulerMarkRecovery).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(); + disconnectedListener!(); + disconnectedListener!(); + await vi.advanceTimersByTimeAsync(6000); + disconnectedListener!(); + + 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('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' }, @@ -1064,4 +1321,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..47c8629a648 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 { @@ -51,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 ( @@ -124,6 +129,168 @@ 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?.(); + }, + }; +} + +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; + 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; + let recoverySourceBridge: AcpBridge | undefined; + + const attachDisconnectHandler = (failedBridge: AcpBridge): void => { + failedBridge.on('disconnected', () => { + if (isShuttingDown() || failedBridge !== getBridge()) return; + if (recoveryTask) { + if (failedBridge !== recoverySourceBridge) recoveryRequested = true; + return; + } + recoverBridge(); + }); + }; + + const recoverBridge = (): void => { + bridgeReadiness.block(); + scheduler?.markBridgeRecovery(); + const task = (async () => { + do { + recoveryRequested = false; + recoverySourceBridge = getBridge(); + const now = Date.now(); + crashTimestamps.push(now); + while (now - crashTimestamps[0]! >= CRASH_WINDOW_MS) { + crashTimestamps.shift(); + } + const recentCrashCount = crashTimestamps.length; + + if (recentCrashCount > MAX_CRASH_RESTARTS) { + writeStderrLine( + `[Channel] Bridge crashed ${recentCrashCount} times in ${CRASH_WINDOW_MS / 1000}s. Giving up.`, + ); + scheduler?.stop(); + cleanupStartedChannels(channels.values(), getBridge(), router); + removeServiceInfo(); + process.exit(1); + } + + writeStderrLine( + `[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)); + + 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 restoreBridgeSessions(router); + 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(); + cleanupStartedChannels(channels.values(), getBridge(), router); + removeServiceInfo(); + process.exit(1); + }) + .finally(() => { + if (recoveryTask === task) { + recoveryTask = undefined; + recoverySourceBridge = undefined; + bridgeReadiness.release(); + } + }); + recoveryTask = task; + }; + + return { attachDisconnectHandler }; +} + /** Check for duplicate instance and abort if one is already running. */ function checkDuplicateInstance(): void { const existing = readServiceInfo(); @@ -178,7 +345,8 @@ async function startSingle( const cliEntryPath = findCliEntryPath(); let shuttingDown = false; - const crashTimestamps: number[] = []; + + const bridgeReadiness = createBridgeReadinessGate(); const bridgeOpts = { cliEntryPath, cwd: config.cwd, model: config.model }; let bridge = new AcpBridge(bridgeOpts); @@ -203,6 +371,7 @@ async function startSingle( proxy, ...channelMemoryOptions(() => bridge, config.cwd), ...(loopController ? { loopController } : {}), + bridgeRecovery: bridgeReadiness.current, }); channels.set(name, channel); const scheduler = loopStore @@ -229,66 +398,22 @@ 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.`); - const attachDisconnectHandler = (b: AcpBridge): void => { - b.on('disconnected', async () => { - if (shuttingDown) return; - - const now = Date.now(); - crashTimestamps.push(now); - // Only count crashes within the recent window - 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.`, - ); - 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...`, - ); - scheduler?.stop(); - await new Promise((r) => setTimeout(r, RESTART_DELAY_MS)); - - try { - bridge = new AcpBridge(bridgeOpts); - 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) { - writeStderrLine( - `[Channel] Failed to restart bridge: ${err instanceof Error ? err.message : String(err)}`, - ); - channel.disconnect(); - router.clearAll(); - removeServiceInfo(); - process.exit(1); - } - }); - }; + const { attachDisconnectHandler } = createBridgeRecovery({ + bridgeOpts, + router, + channels, + scheduler, + bridgeReadiness, + isShuttingDown: () => shuttingDown, + getBridge: () => bridge, + setBridge: (next) => { + bridge = next; + }, + }); attachDisconnectHandler(bridge); const shutdown = () => { @@ -339,7 +464,8 @@ async function startAll( const cliEntryPath = findCliEntryPath(); const defaultCwd = process.cwd(); let shuttingDown = false; - const crashTimestamps: number[] = []; + + const bridgeReadiness = createBridgeReadinessGate(); const bridgeOpts = { cliEntryPath, @@ -374,6 +500,7 @@ async function startAll( proxy, ...channelMemoryOptions(() => bridge, config.cwd), ...(loopController ? { loopController } : {}), + bridgeRecovery: bridgeReadiness.current, }), ); } @@ -414,97 +541,24 @@ 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.`, ); - const attachDisconnectHandler = (b: AcpBridge): void => { - b.on('disconnected', async () => { - if (shuttingDown) return; - - 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.`, - ); - 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...`, - ); - scheduler?.stop(); - await new Promise((r) => setTimeout(r, RESTART_DELAY_MS)); - - try { - bridge = new AcpBridge(bridgeOpts); - 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) { - writeStderrLine( - `[Channel] Failed to restart bridge: ${err instanceof Error ? err.message : String(err)}`, - ); - for (const channel of channels.values()) { - try { - channel.disconnect(); - } catch { - // best-effort - } - } - router.clearAll(); - removeServiceInfo(); - process.exit(1); - } - }); - }; + 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 5e04f2d0bc1..e457f42be37 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); @@ -66,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; @@ -108,6 +139,231 @@ 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, + }); + + vi.setSystemTime(Date.now() + 10_000); + await vi.advanceTimersByTimeAsync(10); + + expect(histogram.reset).toHaveBeenCalledTimes(1); + expect(onNewMaxStall).not.toHaveBeenCalled(); + + monitor.dispose(); + }); + + 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); + 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: 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, + }); + + 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, + }); + + vi.setSystemTime(Date.now() + 10_000); + 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 low-CPU gap just below 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(); + 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 = 300_000_000_000; + const monitor = startEventLoopLagMonitor({ + resolutionMs: 10, + stallThresholdMs: 1_000, + onNewMaxStall, + }); + + vi.setSystemTime(Date.now() + 600_000); + 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('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('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 2aacdaf84cf..7b79e19c75a 100644 --- a/packages/core/src/telemetry/event-loop-lag.ts +++ b/packages/core/src/telemetry/event-loop-lag.ts @@ -21,11 +21,26 @@ 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; +/** + * 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; export function startEventLoopLagMonitor( @@ -39,16 +54,59 @@ export function startEventLoopLagMonitor( options.stallThresholdMs, DEFAULT_STALL_THRESHOLD_MS, ); + const suspendThresholdMs = positiveFiniteOrDefault( + options.suspendThresholdMs, + DEFAULT_EVENT_LOOP_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 lastObservedMaxMs = 0; + let lastCheckTimeMs = Date.now(); + let lastCpuUsage = safeCpuUsage(); + let pendingSuspendGapMs = 0; 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 newMaxMs = maxMs > lastObservedMaxMs ? maxMs : 0; + const cpuRatio = calculateCpuRatio(lastCpuUsage, cpuUsage, elapsedMs); + lastCheckTimeMs = nowMs; + if (cpuUsage) lastCpuUsage = cpuUsage; + lastObservedMaxMs = maxMs; + const isLowCpuGap = + elapsedMs >= suspendThresholdMs && + cpuRatio !== undefined && + 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 ( + options.onNewMaxStall && + maxMs >= stallThresholdMs && + maxMs > lastReportedMaxMs + ) { lastReportedMaxMs = maxMs; try { options.onNewMaxStall(maxMs); @@ -57,11 +115,8 @@ export function startEventLoopLagMonitor( } } }; - const interval = - options.onNewMaxStall !== undefined - ? setInterval(checkForNewMaxStall, resolutionMs) - : undefined; - interval?.unref(); + const interval = setInterval(checkHistogram, resolutionMs); + interval.unref(); return { snapshot(): EventLoopLagSnapshot { @@ -75,14 +130,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 +164,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; +}