diff --git a/packages/channels/base/src/AcpBridge.test.ts b/packages/channels/base/src/AcpBridge.test.ts index 8f0c0e5e68d..5133ebe2f1c 100644 --- a/packages/channels/base/src/AcpBridge.test.ts +++ b/packages/channels/base/src/AcpBridge.test.ts @@ -362,6 +362,62 @@ describe('AcpBridge', () => { ); }); + it('emits a completed background response separately from the active turn', () => { + const bridge = new AcpBridge({ + cliEntryPath: '/tmp/qwen', + cwd: '/tmp', + }) as unknown as TestableAcpBridge; + const backgroundResponses: Array<[string, string]> = []; + bridge.on('backgroundResponse', (sessionId, text) => { + backgroundResponses.push([sessionId, text]); + }); + const textChunks: Array<[string, string]> = []; + bridge.on('textChunk', (sessionId, text) => { + textChunks.push([sessionId, text]); + }); + + bridge.handleSessionUpdate({ + sessionId: 's-1', + update: { + sessionUpdate: 'agent_message_chunk', + content: { type: 'text', text: 'Background final answer.' }, + _meta: { + source: 'background_notification_response', + qwenDiscreteMessage: true, + }, + }, + }); + + expect(backgroundResponses).toEqual([['s-1', 'Background final answer.']]); + expect(textChunks).toEqual([]); + }); + + it('ignores a rewritten background response to avoid duplicate delivery', () => { + const bridge = new AcpBridge({ + cliEntryPath: '/tmp/qwen', + cwd: '/tmp', + }) as unknown as TestableAcpBridge; + const backgroundResponses: Array<[string, string]> = []; + bridge.on('backgroundResponse', (sessionId, text) => { + backgroundResponses.push([sessionId, text]); + }); + + bridge.handleSessionUpdate({ + sessionId: 's-1', + update: { + sessionUpdate: 'agent_message_chunk', + content: { type: 'text', text: 'Background final answer.' }, + _meta: { + source: 'background_notification_response', + qwenDiscreteMessage: true, + rewritten: true, + }, + }, + }); + + expect(backgroundResponses).toEqual([]); + }); + it('returns only the final slash-command output', async () => { const bridge = new AcpBridge({ cliEntryPath: '/tmp/qwen', diff --git a/packages/channels/base/src/AcpBridge.ts b/packages/channels/base/src/AcpBridge.ts index 5c1e60e411c..b3d4e8c214f 100644 --- a/packages/channels/base/src/AcpBridge.ts +++ b/packages/channels/base/src/AcpBridge.ts @@ -314,15 +314,23 @@ export class AcpBridge extends EventEmitter implements ChannelAgentBridge { switch (type) { case 'agent_message_chunk': { const meta = update['_meta'] as Record | undefined; - if ( - typeof meta?.['parentToolCallId'] === 'string' || - meta?.['qwenDiscreteMessage'] === true - ) { + if (typeof meta?.['parentToolCallId'] === 'string') { break; } const content = update['content'] as | { type?: string; text?: string } | undefined; + if (meta?.['qwenDiscreteMessage'] === true) { + if ( + meta['source'] === 'background_notification_response' && + meta['rewritten'] !== true && + content?.type === 'text' && + content.text + ) { + this.emit('backgroundResponse', sessionId, content.text); + } + break; + } if (content?.type === 'text' && content.text) { this.emit( meta?.['source'] === 'slash_command' diff --git a/packages/channels/base/src/ChannelAgentBridge.ts b/packages/channels/base/src/ChannelAgentBridge.ts index 3191057ad51..4f6d166f6c7 100644 --- a/packages/channels/base/src/ChannelAgentBridge.ts +++ b/packages/channels/base/src/ChannelAgentBridge.ts @@ -66,6 +66,7 @@ export interface PermissionResolvedEvent { interface ChannelAgentBridgeEventMap { sessionDied: [SessionDiedEvent]; textChunk: [sessionId: string, chunk: string]; + backgroundResponse: [sessionId: string, text: string]; responseBoundary: [sessionId: string]; toolCall: [ToolCallEvent]; permissionRequest: [PermissionRequestEvent]; diff --git a/packages/channels/base/src/ChannelBase.test.ts b/packages/channels/base/src/ChannelBase.test.ts index b20737aab39..245a0fd011d 100644 --- a/packages/channels/base/src/ChannelBase.test.ts +++ b/packages/channels/base/src/ChannelBase.test.ts @@ -7060,6 +7060,70 @@ describe('ChannelBase', () => { expect(router.handleSessionDied).toHaveBeenCalledWith('s-1'); }); + it('proactively delivers a completed background response to the session route', async () => { + const target: SessionTarget = { + channelName: 'test-chan', + senderId: 'user1', + chatId: 'chat1', + isGroup: true, + }; + const router = { + getTarget: vi.fn().mockReturnValue(target), + handleSessionDied: vi.fn(), + setBridge: vi.fn(), + }; + const ch = createChannel({}, { + router, + registerBridgeEvents: true, + } as unknown as ChannelBaseOptions); + ch.proactiveSupported = true; + + (bridge as unknown as EventEmitter).emit( + 'backgroundResponse', + 's-1', + 'Background final answer.', + ); + + await vi.waitFor(() => { + expect(ch.proactive).toEqual([ + { chatId: 'chat1', text: 'Background final answer.' }, + ]); + }); + expect(ch.proactiveTargets).toEqual([target]); + expect(ch.sent).toEqual([]); + }); + + it('falls back to sendResponseMessage when proactive send is unsupported', async () => { + const target: SessionTarget = { + channelName: 'test-chan', + senderId: 'user1', + chatId: 'chat1', + isGroup: true, + }; + const router = { + getTarget: vi.fn().mockReturnValue(target), + handleSessionDied: vi.fn(), + setBridge: vi.fn(), + }; + const ch = createChannel({}, { + router, + registerBridgeEvents: true, + } as unknown as ChannelBaseOptions); + + (bridge as unknown as EventEmitter).emit( + 'backgroundResponse', + 's-1', + 'Background final answer.', + ); + + await vi.waitFor(() => { + expect(ch.sent).toEqual([ + { chatId: 'chat1', text: 'Background final answer.' }, + ]); + }); + expect(ch.proactive).toEqual([]); + }); + it('leaves supplied router bridge events to the gateway by default', () => { const router = { getTarget: vi.fn(), diff --git a/packages/channels/base/src/ChannelBase.ts b/packages/channels/base/src/ChannelBase.ts index 21302f46bdd..f67bd31f3fe 100644 --- a/packages/channels/base/src/ChannelBase.ts +++ b/packages/channels/base/src/ChannelBase.ts @@ -346,6 +346,18 @@ export abstract class ChannelBase { private readonly bridgeToolCallListener = (event: ToolCallEvent): void => { this.dispatchToolCall(event); }; + private readonly bridgeBackgroundResponseListener = ( + sessionId: string, + text: string, + ): void => { + void this.dispatchBackgroundResponse(sessionId, text).catch( + (err: unknown) => { + process.stderr.write( + `[${this.name}] background response delivery failed for session ${sanitizeLogText(sessionId, 128)}: ${this.lifecycleError(err)}\n`, + ); + }, + ); + }; private readonly bridgeSessionDiedListener = ( event: SessionDiedEvent, ): void => { @@ -403,6 +415,25 @@ export abstract class ChannelBase { this.onToolCall(chatId, event); } + async dispatchBackgroundResponse( + sessionId: string, + text: string, + ): Promise { + const target = this.router.getTarget(sessionId); + if ( + !target || + target.channelName !== this.name || + text.trim().length === 0 + ) { + return; + } + if (this.supportsProactiveSend() && this.supportsProactiveTarget(target)) { + await this.pushProactive(target, text); + return; + } + await this.sendResponseMessage(target.chatId, text, sessionId); + } + async dispatchPermissionRequest( event: PermissionRequestEvent, ): Promise { @@ -1704,6 +1735,7 @@ export abstract class ChannelBase { private attachBridgeEvents(bridge: ChannelAgentBridge): void { bridge.on('toolCall', this.bridgeToolCallListener); + bridge.on('backgroundResponse', this.bridgeBackgroundResponseListener); bridge.on('sessionDied', this.bridgeSessionDiedListener); bridge.on('permissionRequest', this.bridgePermissionRequestListener); bridge.on('permissionResolved', this.bridgePermissionResolvedListener); @@ -1711,6 +1743,7 @@ export abstract class ChannelBase { private detachBridgeEvents(bridge: ChannelAgentBridge): void { bridge.off('toolCall', this.bridgeToolCallListener); + bridge.off('backgroundResponse', this.bridgeBackgroundResponseListener); bridge.off('sessionDied', this.bridgeSessionDiedListener); bridge.off('permissionRequest', this.bridgePermissionRequestListener); bridge.off('permissionResolved', this.bridgePermissionResolvedListener); diff --git a/packages/channels/base/src/DaemonChannelBridge.test.ts b/packages/channels/base/src/DaemonChannelBridge.test.ts index 649de326065..50c8a343200 100644 --- a/packages/channels/base/src/DaemonChannelBridge.test.ts +++ b/packages/channels/base/src/DaemonChannelBridge.test.ts @@ -566,6 +566,112 @@ describe('DaemonChannelBridge', () => { bridge.stop(); }); + it('emits the completed background response without appending it to the active turn', async () => { + const events = new EventQueue(); + const session = createFakeSession(events); + session.prompt.mockImplementation(async () => { + events.push({ + id: 1, + v: 1, + type: 'session_update', + data: { + sessionId: 'session-1', + update: { + sessionUpdate: 'agent_message_chunk', + content: { type: 'text', text: 'Initial answer.' }, + }, + }, + }); + events.push(turnCompleteEvent()); + return { stopReason: 'end_turn' }; + }); + const bridge = new DaemonChannelBridge({ + cwd: '/repo', + sessionFactory: vi.fn().mockResolvedValue(session), + }); + const backgroundResponses: Array<[string, string]> = []; + bridge.on('backgroundResponse', (sessionId, text) => { + backgroundResponses.push([sessionId, text]); + }); + + await bridge.start(); + await bridge.newSession('/repo'); + await expect(bridge.prompt('session-1', 'investigate')).resolves.toBe( + 'Initial answer.', + ); + + events.push({ + id: 2, + v: 1, + type: 'session_update', + data: { + sessionId: 'session-1', + update: { + sessionUpdate: 'agent_message_chunk', + content: { type: 'text', text: 'Background final answer.' }, + _meta: { + source: 'background_notification_response', + qwenDiscreteMessage: true, + }, + }, + }, + }); + + await vi.waitFor(() => { + expect(backgroundResponses).toEqual([ + ['session-1', 'Background final answer.'], + ]); + }); + + events.close(); + bridge.stop(); + }); + + it('ignores a rewritten background response to avoid duplicate delivery', async () => { + const events = new EventQueue(); + const session = createFakeSession(events); + session.prompt.mockImplementation(async () => { + events.push(turnCompleteEvent()); + return { stopReason: 'end_turn' }; + }); + const bridge = new DaemonChannelBridge({ + cwd: '/repo', + sessionFactory: vi.fn().mockResolvedValue(session), + }); + const backgroundResponses: Array<[string, string]> = []; + bridge.on('backgroundResponse', (sessionId, text) => { + backgroundResponses.push([sessionId, text]); + }); + + await bridge.start(); + await bridge.newSession('/repo'); + await bridge.prompt('session-1', 'investigate'); + + events.push({ + id: 2, + v: 1, + type: 'session_update', + data: { + sessionId: 'session-1', + update: { + sessionUpdate: 'agent_message_chunk', + content: { type: 'text', text: 'Background final answer.' }, + _meta: { + source: 'background_notification_response', + qwenDiscreteMessage: true, + rewritten: true, + }, + }, + }, + }); + + await new Promise((r) => setTimeout(r, 50)); + expect(backgroundResponses).toEqual([]); + + events.close(); + bridge.stop(); + }); + it('returns only the final slash-command output from the daemon', async () => { const events = new EventQueue(); const session = createFakeSession(events); diff --git a/packages/channels/base/src/DaemonChannelBridge.ts b/packages/channels/base/src/DaemonChannelBridge.ts index d5fea4c7e58..eaddd39c144 100644 --- a/packages/channels/base/src/DaemonChannelBridge.ts +++ b/packages/channels/base/src/DaemonChannelBridge.ts @@ -649,13 +649,20 @@ export class DaemonChannelBridge switch (type) { case 'agent_message_chunk': { const meta = isRecord(update['_meta']) ? update['_meta'] : undefined; - if ( - typeof meta?.['parentToolCallId'] === 'string' || - meta?.['qwenDiscreteMessage'] === true - ) { + if (typeof meta?.['parentToolCallId'] === 'string') { break; } const text = getTextContent(update['content']); + if (meta?.['qwenDiscreteMessage'] === true) { + if ( + meta['source'] === 'background_notification_response' && + meta['rewritten'] !== true && + text + ) { + this.emit('backgroundResponse', sessionId, text); + } + break; + } if (text) { this.emit( meta?.['source'] === 'slash_command' diff --git a/packages/cli/src/commands/channel/daemon-worker.test.ts b/packages/cli/src/commands/channel/daemon-worker.test.ts index 38395db5e7e..2fab222480c 100644 --- a/packages/cli/src/commands/channel/daemon-worker.test.ts +++ b/packages/cli/src/commands/channel/daemon-worker.test.ts @@ -13,6 +13,7 @@ const mockUpdateChannelMemoryEntry = vi.hoisted(() => vi.fn()); const mockRemoveChannelMemoryEntries = vi.hoisted(() => vi.fn()); const mockClearChannelMemory = vi.hoisted(() => vi.fn()); const mockRegisterToolCallDispatch = vi.hoisted(() => vi.fn()); +const mockRegisterBackgroundResponseRelay = vi.hoisted(() => vi.fn()); const mockRegisterPermissionRelay = vi.hoisted(() => vi.fn()); const mockRegisterSessionCleanup = vi.hoisted(() => vi.fn()); const mockSessionsPath = vi.hoisted(() => vi.fn(() => '/tmp/sessions.json')); @@ -174,6 +175,7 @@ vi.mock('./runtime.js', () => ({ loadChannelsConfig: mockLoadChannelsConfig, loadChannelsFromExtensions: mockLoadChannelsFromExtensions, parseConfiguredChannels: mockParseConfiguredChannels, + registerBackgroundResponseRelay: mockRegisterBackgroundResponseRelay, registerPermissionRelay: mockRegisterPermissionRelay, registerSessionCleanup: mockRegisterSessionCleanup, registerToolCallDispatch: mockRegisterToolCallDispatch, @@ -733,6 +735,11 @@ describe('runChannelDaemonWorker', () => { mockSessionRouter.mock.results[0]!.value, expect.any(Map), ); + expect(mockRegisterBackgroundResponseRelay).toHaveBeenCalledWith( + bridgeFacade, + mockSessionRouter.mock.results[0]!.value, + expect.any(Map), + ); expect(mockResolveProxyUrl).toHaveBeenCalledWith( undefined, 'http://settings-proxy:8080', diff --git a/packages/cli/src/commands/channel/daemon-worker.ts b/packages/cli/src/commands/channel/daemon-worker.ts index 994862ce2aa..b6606bcd218 100644 --- a/packages/cli/src/commands/channel/daemon-worker.ts +++ b/packages/cli/src/commands/channel/daemon-worker.ts @@ -57,6 +57,7 @@ import { loadChannelsConfig, loadChannelsFromExtensions, parseConfiguredChannels, + registerBackgroundResponseRelay, registerPermissionRelay, registerSessionCleanup, registerToolCallDispatch, @@ -516,6 +517,7 @@ export async function runChannelDaemonWorker( ); } registerToolCallDispatch(bridgeFacade, createdRouter, channels); + registerBackgroundResponseRelay(bridgeFacade, createdRouter, channels); registerPermissionRelay(bridgeFacade, createdRouter, channels); registerSessionCleanup(bridgeFacade, createdRouter, channels); diff --git a/packages/cli/src/commands/channel/runtime.test.ts b/packages/cli/src/commands/channel/runtime.test.ts index 4ac93f068bd..5eda92fbb65 100644 --- a/packages/cli/src/commands/channel/runtime.test.ts +++ b/packages/cli/src/commands/channel/runtime.test.ts @@ -4,6 +4,7 @@ import { daemonObservedContactsPath, daemonSessionRoutesPath, parseConfiguredChannels, + registerBackgroundResponseRelay, registerPermissionRelay, registerSessionCleanup, sessionsPath, @@ -308,6 +309,137 @@ describe('registerPermissionRelay', () => { }); }); +describe('registerBackgroundResponseRelay', () => { + it('routes the final background response without joining the active prompt', async () => { + const bridge = new EventEmitter(); + const router = { + getTarget: vi.fn(() => ({ + channelName: 'telegram', + chatId: 'chat1', + })), + }; + const channel = { + dispatchBackgroundResponse: vi.fn().mockResolvedValue(undefined), + }; + + registerBackgroundResponseRelay( + bridge as never, + router as never, + new Map([['telegram', channel as never]]), + ); + bridge.emit('backgroundResponse', 'session-1', 'Background final answer.'); + + await vi.waitFor(() => { + expect(channel.dispatchBackgroundResponse).toHaveBeenCalledWith( + 'session-1', + 'Background final answer.', + ); + }); + }); + + it('logs when no route exists for the background response', async () => { + const stderr = vi + .spyOn(process.stderr, 'write') + .mockImplementation(() => true); + const bridge = new EventEmitter(); + const router = { getTarget: vi.fn() }; + + try { + registerBackgroundResponseRelay( + bridge as never, + router as never, + new Map(), + ); + bridge.emit( + 'backgroundResponse', + 'session-1', + 'Background final answer.', + ); + + await vi.waitFor(() => { + expect(stderr.mock.calls.join('')).toContain( + 'No route for background response from session session-1', + ); + }); + } finally { + stderr.mockRestore(); + } + }); + + it('logs when the channel is not found for the background response', async () => { + const stderr = vi + .spyOn(process.stderr, 'write') + .mockImplementation(() => true); + const bridge = new EventEmitter(); + const router = { + getTarget: vi.fn(() => ({ + channelName: 'telegram', + chatId: 'chat1', + })), + }; + + try { + registerBackgroundResponseRelay( + bridge as never, + router as never, + new Map(), + ); + bridge.emit( + 'backgroundResponse', + 'session-1', + 'Background final answer.', + ); + + await vi.waitFor(() => { + expect(stderr.mock.calls.join('')).toContain( + 'No channel "telegram" for background response from session session-1', + ); + }); + } finally { + stderr.mockRestore(); + } + }); + + it('logs when dispatchBackgroundResponse rejects', async () => { + const stderr = vi + .spyOn(process.stderr, 'write') + .mockImplementation(() => true); + const bridge = new EventEmitter(); + const router = { + getTarget: vi.fn(() => ({ + channelName: 'telegram', + chatId: 'chat1', + })), + }; + const channel = { + dispatchBackgroundResponse: vi + .fn() + .mockRejectedValue(new Error('network down')), + }; + + try { + registerBackgroundResponseRelay( + bridge as never, + router as never, + new Map([['telegram', channel as never]]), + ); + bridge.emit( + 'backgroundResponse', + 'session-1', + 'Background final answer.', + ); + + await vi.waitFor(() => { + expect(stderr.mock.calls.join('')).toContain( + 'Background response relay failed for session session-1', + ); + }); + } finally { + stderr.mockRestore(); + } + }); +}); + describe('registerSessionCleanup', () => { it('updates routing state when no channel matches the dead session', () => { const bridge = new EventEmitter(); diff --git a/packages/cli/src/commands/channel/runtime.ts b/packages/cli/src/commands/channel/runtime.ts index 5908efef90a..4879b8eea35 100644 --- a/packages/cli/src/commands/channel/runtime.ts +++ b/packages/cli/src/commands/channel/runtime.ts @@ -186,6 +186,36 @@ export function registerToolCallDispatch( }); } +export function registerBackgroundResponseRelay( + bridge: ChannelAgentBridge, + router: SessionRouter, + channels: Map, +): void { + bridge.on('backgroundResponse', (sessionId: string, text: string) => { + const target = router.getTarget(sessionId); + if (!target) { + writeStderrLine( + `[Channel] No route for background response from session ${sanitizeLogText(sessionId, 128)}`, + ); + return; + } + const channel = channels.get(target.channelName); + if (!channel) { + writeStderrLine( + `[Channel] No channel "${sanitizeLogText(target.channelName, 64)}" for background response from session ${sanitizeLogText(sessionId, 128)}`, + ); + return; + } + void channel + .dispatchBackgroundResponse(sessionId, text) + .catch((err: unknown) => { + writeStderrLine( + `[Channel] Background response relay failed for session ${sanitizeLogText(sessionId, 128)}: ${err instanceof Error ? sanitizeLogText(err.message, 512) : sanitizeLogText(String(err), 512)}`, + ); + }); + }); +} + function cancelPermissionRequest( bridge: ChannelAgentBridge, requestId: string, diff --git a/packages/cli/src/commands/channel/start.ts b/packages/cli/src/commands/channel/start.ts index b9869b923c2..678a8012c42 100644 --- a/packages/cli/src/commands/channel/start.ts +++ b/packages/cli/src/commands/channel/start.ts @@ -36,6 +36,7 @@ import { loadChannelsConfig, loadChannelsFromExtensions, parseConfiguredChannels, + registerBackgroundResponseRelay, registerPermissionRelay, registerSessionCleanup, registerToolCallDispatch, @@ -232,6 +233,7 @@ async function startSingle( }) : undefined; registerToolCallDispatch(bridge, router, channels); + registerBackgroundResponseRelay(bridge, router, channels); registerPermissionRelay(bridge, router, channels); registerSessionCleanup(bridge, router, channels); @@ -286,6 +288,7 @@ async function startSingle( channel.disconnect(); await channel.connect(); registerToolCallDispatch(bridge, router, channels); + registerBackgroundResponseRelay(bridge, router, channels); registerPermissionRelay(bridge, router, channels); registerSessionCleanup(bridge, router, channels); attachDisconnectHandler(bridge); @@ -395,6 +398,7 @@ async function startAll( ); } registerToolCallDispatch(bridge, router, channels); + registerBackgroundResponseRelay(bridge, router, channels); registerPermissionRelay(bridge, router, channels); registerSessionCleanup(bridge, router, channels); @@ -494,6 +498,7 @@ async function startAll( process.exit(1); } registerToolCallDispatch(bridge, router, channels); + registerBackgroundResponseRelay(bridge, router, channels); registerPermissionRelay(bridge, router, channels); registerSessionCleanup(bridge, router, channels); attachDisconnectHandler(bridge);