diff --git a/docs/developers/daemon-client-adapters/channel-web.md b/docs/developers/daemon-client-adapters/channel-web.md new file mode 100644 index 00000000000..44caca3afd5 --- /dev/null +++ b/docs/developers/daemon-client-adapters/channel-web.md @@ -0,0 +1,120 @@ +# Channel And Web Backend Daemon Adapter Draft + +## Goal + +Let channel adapters and web chat backends consume `qwen serve` through +`DaemonSessionClient` while keeping existing channel ACP subprocess behavior as +the default. + +This draft covers server-side clients only: + +- Channel bot backend -> `qwen serve` +- Web browser -> web backend / BFF -> `qwen serve` + +It explicitly does not allow browser JavaScript to call the daemon directly. +The daemon currently rejects browser `Origin` requests by design. + +## Proposed Entry Points + +Channel backend: + +```bash +QWEN_CHANNEL_DAEMON_URL=http://127.0.0.1:4170 qwen channel start telegram +``` + +Web backend: + +```bash +QWEN_WEB_DAEMON_URL=http://127.0.0.1:4170 qwen web-chat-backend +``` + +Shared optional variables: + +```bash +QWEN_DAEMON_TOKEN=... +QWEN_DAEMON_WORKSPACE=/repo +``` + +## Minimal Channel Flow + +This PR adds `DaemonChannelBridge`, a locally verifiable server-side bridge for +channel and web-backend adapters. It keeps the existing ACP bridge as the +default and owns daemon session state inside the backend process. + +1. Resolve channel sender/thread to a channel session key. +2. Use `DaemonClient` + `DaemonSessionClient.createOrAttach()`. +3. Submit inbound user text with `session.prompt()`. +4. Subscribe to `session.events()` and collect assistant text chunks. +5. Send final text back through the platform adapter. +6. Cast permission votes through `session.respondToPermission()`. +7. Cancel active work through `session.cancel()`. + +## Minimal Web Backend Flow + +1. Browser opens a websocket or HTTP stream to the web backend. +2. Backend owns `DaemonSessionClient`. +3. Backend translates browser messages to daemon prompts. +4. Backend translates daemon SSE events to browser-safe app events. +5. Backend stores the daemon `sessionId` and last seen event id server-side. + +Browser clients must not receive daemon bearer tokens. + +## Session Isolation Constraint + +Current daemon Stage 1 behavior is effectively `sessionScope: single` at the +daemon setting level. Until per-request `sessionScope` lands, multi-user channel +or web deployments must choose one of these safe shapes: + +- one daemon per channel thread / web room +- one daemon per user workspace +- single-user demo only + +Do not silently multiplex unrelated channel threads into one daemon session. + +## Event Mapping Contract + +| Daemon event | Channel/web backend handling | +| ---------------------------------------- | -------------------------------------- | +| `session_update` / `agent_message_chunk` | Append assistant text | +| `session_update` / `agent_thought_chunk` | Optional hidden/debug stream | +| `session_update` / `tool_call` | Emit tool status card/message | +| `permission_request` | Platform-specific approval interaction | +| `permission_resolved` | Close/update approval interaction | +| `model_switched` | Update backend session metadata | +| `session_died` | Notify user and stop stream | + +Unknown daemon events must be ignored or forwarded as debug metadata, not fatal. + +The bridge is not wired into `qwen channel start` yet. Existing Telegram, +Weixin, Dingtalk, plugin channel, and browser behavior remains unchanged. + +## Explicit Non-Goals + +- No browser direct-to-daemon fetch or EventSource. +- No CORS relaxation in this adapter PR. +- No default migration of Telegram, Weixin, Dingtalk, or plugin channels. +- No file CRUD, memory CRUD, MCP restart, or provider mutation. +- No sessionScope emulation in the client when daemon-side support is absent. + +## Merge Safety + +- Default off. +- Existing ACP channel bridge remains the default. +- Web backend is an explicit BFF layer, not a daemon security change. +- No channel adapter should import daemon tokens into frontend/browser code. + +## Validation Plan + +- Unit-test channel session-key to daemon-session binding. +- Unit-test daemon event to channel/web message mapping. +- Unit-test prompt, cancel, model switch, and permission response forwarding. +- Smoke-test one single-user channel backend against local `qwen serve`. +- Smoke-test browser -> BFF -> daemon without exposing daemon token. + +## Blockers Before Default Migration + +- Per-request `sessionScope`. +- Session metadata + close/delete lifecycle. +- Daemon-stamped client identity. +- Session-scoped permission route. +- Read-only diagnostics for MCP, skills, providers, and environment. diff --git a/packages/channels/base/src/DaemonChannelBridge.test.ts b/packages/channels/base/src/DaemonChannelBridge.test.ts new file mode 100644 index 00000000000..8019613adfb --- /dev/null +++ b/packages/channels/base/src/DaemonChannelBridge.test.ts @@ -0,0 +1,1176 @@ +import { describe, expect, it, vi } from 'vitest'; +import type { + RequestPermissionRequest, + RequestPermissionResponse, +} from '@agentclientprotocol/sdk'; +import { + DaemonChannelBridge, + type DaemonChannelEvent, + type DaemonChannelSessionClient, +} from './DaemonChannelBridge.js'; + +class EventQueue implements AsyncGenerator { + private events: DaemonChannelEvent[] = []; + private waiters: Array<{ + resolve: (value: IteratorResult) => void; + reject: (error: unknown) => void; + }> = []; + private closed = false; + private failure: unknown; + + async next(): Promise> { + if (this.failure) { + throw this.failure; + } + const event = this.events.shift(); + if (event) { + return { done: false, value: event }; + } + if (this.closed) { + return { done: true, value: undefined }; + } + return await new Promise((resolve, reject) => { + this.waiters.push({ resolve, reject }); + }); + } + + async return(): Promise> { + this.close(); + return { done: true, value: undefined }; + } + + async throw(error?: unknown): Promise> { + this.close(); + throw error; + } + + [Symbol.asyncIterator](): AsyncGenerator { + return this; + } + + push(event: DaemonChannelEvent): void { + const waiter = this.waiters.shift(); + if (waiter) { + waiter.resolve({ done: false, value: event }); + return; + } + this.events.push(event); + } + + close(): void { + this.closed = true; + for (const waiter of this.waiters.splice(0)) { + waiter.resolve({ done: true, value: undefined }); + } + } + + fail(error: unknown): void { + this.failure = error; + for (const waiter of this.waiters.splice(0)) { + waiter.reject(error); + } + } +} + +interface FakeSession extends DaemonChannelSessionClient { + prompt: ReturnType; + events: ReturnType; + cancel: ReturnType; + setModel: ReturnType; + respondToPermission: ReturnType; +} + +function createFakeSession( + events: EventQueue, + sessionId = 'session-1', +): FakeSession { + return { + sessionId, + workspaceCwd: '/repo', + lastEventId: undefined, + prompt: vi.fn().mockImplementation(async () => undefined), + events: vi.fn((opts?: { signal?: AbortSignal }) => { + opts?.signal?.addEventListener('abort', () => events.close(), { + once: true, + }); + return events; + }), + cancel: vi.fn().mockResolvedValue(undefined), + setModel: vi.fn().mockResolvedValue({}), + respondToPermission: vi.fn().mockResolvedValue(true), + }; +} + +async function waitFor(assertion: () => void): Promise { + let lastError: unknown; + for (let i = 0; i < 20; i += 1) { + try { + assertion(); + return; + } catch (error) { + lastError = error; + await new Promise((resolve) => setTimeout(resolve, 0)); + } + } + throw lastError; +} + +describe('DaemonChannelBridge', () => { + it('binds a daemon session and collects assistant chunks during prompt', async () => { + const events = new EventQueue(); + const session = createFakeSession(events); + let resolvePrompt: () => void = () => {}; + session.prompt.mockImplementation( + () => + new Promise((resolve) => { + resolvePrompt = () => resolve({ stopReason: 'end_turn' }); + events.push({ + id: 1, + v: 1, + type: 'session_update', + data: { + sessionId: 'session-1', + update: { + sessionUpdate: 'agent_message_chunk', + content: { type: 'text', text: 'hello' }, + }, + }, + }); + }), + ); + const factory = vi.fn().mockResolvedValue(session); + const bridge = new DaemonChannelBridge({ + cwd: '/repo', + sessionFactory: factory, + }); + const promptComplete = vi.fn(); + bridge.on('promptComplete', promptComplete); + + await bridge.start(); + const sessionId = await bridge.newSession('/repo'); + const promptPromise = bridge.prompt(sessionId, 'summarize'); + await waitFor(() => expect(session.prompt).toHaveBeenCalledOnce()); + resolvePrompt(); + + await expect(promptPromise).resolves.toBe('hello'); + expect(promptComplete).toHaveBeenCalledWith({ + sessionId: 'session-1', + text: 'hello', + stopReason: 'end_turn', + }); + expect(factory).toHaveBeenCalledWith({ + workspaceCwd: '/repo', + modelServiceId: undefined, + sessionScope: 'thread', + }); + expect(session.prompt).toHaveBeenCalledWith( + { + prompt: [{ type: 'text', text: 'summarize' }], + }, + expect.any(AbortSignal), + ); + + events.close(); + bridge.stop(); + }); + + it('drains daemon chunks queued with prompt completion', async () => { + const events = new EventQueue(); + const session = createFakeSession(events); + session.prompt.mockImplementation(async () => { + setTimeout(() => { + events.push({ + id: 1, + v: 1, + type: 'session_update', + data: { + sessionId: 'session-1', + update: { + sessionUpdate: 'agent_message_chunk', + content: { type: 'text', text: 'late chunk' }, + }, + }, + }); + }, 0); + return { stopReason: 'end_turn' }; + }); + const bridge = new DaemonChannelBridge({ + cwd: '/repo', + sessionFactory: vi.fn().mockResolvedValue(session), + }); + + await bridge.start(); + await bridge.newSession('/repo'); + + await expect(bridge.prompt('session-1', 'summarize')).resolves.toBe( + 'late chunk', + ); + + events.close(); + bridge.stop(); + }); + + it('emits tool, thought, model, commands, and session lifecycle events', async () => { + const events = new EventQueue(); + const session = createFakeSession(events); + const bridge = new DaemonChannelBridge({ + cwd: '/repo', + sessionFactory: vi.fn().mockResolvedValue(session), + }); + const thoughtChunk = vi.fn(); + const toolCall = vi.fn(); + const modelSwitched = vi.fn(); + const modelSwitchFailed = vi.fn(); + const sessionDied = vi.fn(); + + bridge.on('thoughtChunk', thoughtChunk); + bridge.on('toolCall', toolCall); + bridge.on('modelSwitched', modelSwitched); + bridge.on('modelSwitchFailed', modelSwitchFailed); + bridge.on('sessionDied', sessionDied); + + await bridge.start(); + await bridge.newSession('/repo'); + + events.push({ + id: 2, + v: 1, + type: 'session_update', + data: { + sessionId: 'session-1', + update: { + sessionUpdate: 'agent_thought_chunk', + content: { type: 'text', text: 'thinking' }, + }, + }, + }); + events.push({ + id: 3, + v: 1, + type: 'session_update', + data: { + sessionId: 'session-1', + update: { + sessionUpdate: 'tool_call_update', + toolCallId: 'tool-1', + kind: 'read_file', + title: 'Read file', + status: 'completed', + rawInput: { path: 'README.md' }, + }, + }, + }); + events.push({ + id: 4, + v: 1, + type: 'session_update', + data: { + sessionId: 'session-1', + update: { + sessionUpdate: 'available_commands_update', + availableCommands: [ + { name: '/help', description: 'Show help', input: null }, + null, + { description: 'Missing name', input: null }, + ], + }, + }, + }); + await waitFor(() => + expect(bridge.getAvailableCommands('session-1')).toEqual([ + { name: '/help', description: 'Show help', input: null }, + ]), + ); + expect(bridge.availableCommands).toEqual([ + { name: '/help', description: 'Show help', input: null }, + ]); + + events.push({ + id: 5, + v: 1, + type: 'model_switched', + data: { sessionId: 'session-1', modelId: 'qwen3-coder-plus' }, + }); + events.push({ + id: 6, + v: 1, + type: 'model_switch_failed', + data: { + sessionId: 'session-1', + requestedModelId: 'missing-model', + error: 'not configured', + }, + }); + events.push({ + id: 7, + v: 1, + type: 'session_died', + data: { sessionId: 'session-1', reason: 'agent exited' }, + }); + + await waitFor(() => + expect(thoughtChunk).toHaveBeenCalledWith('session-1', 'thinking'), + ); + expect(toolCall).toHaveBeenCalledWith({ + sessionId: 'session-1', + toolCallId: 'tool-1', + kind: 'read_file', + title: 'Read file', + status: 'completed', + rawInput: { path: 'README.md' }, + }); + await waitFor(() => + expect(modelSwitched).toHaveBeenCalledWith({ + sessionId: 'session-1', + modelId: 'qwen3-coder-plus', + }), + ); + await waitFor(() => + expect(modelSwitchFailed).toHaveBeenCalledWith({ + sessionId: 'session-1', + requestedModelId: 'missing-model', + error: 'not configured', + }), + ); + await waitFor(() => + expect(sessionDied).toHaveBeenCalledWith({ + sessionId: 'session-1', + reason: 'agent exited', + }), + ); + expect(bridge.getAvailableCommands('session-1')).toEqual([]); + + events.close(); + }); + + it('keeps available commands scoped per daemon session', async () => { + const firstEvents = new EventQueue(); + const secondEvents = new EventQueue(); + const firstSession = createFakeSession(firstEvents, 'session-1'); + const secondSession = createFakeSession(secondEvents, 'session-2'); + const bridge = new DaemonChannelBridge({ + cwd: '/repo', + sessionFactory: vi + .fn() + .mockResolvedValueOnce(firstSession) + .mockResolvedValueOnce(secondSession), + }); + + await bridge.start(); + await bridge.newSession('/repo'); + await bridge.newSession('/repo'); + + firstEvents.push({ + id: 1, + v: 1, + type: 'session_update', + data: { + sessionId: 'session-1', + update: { + sessionUpdate: 'available_commands_update', + availableCommands: [ + { name: '/one', description: 'First', input: null }, + ], + }, + }, + }); + secondEvents.push({ + id: 2, + v: 1, + type: 'session_update', + data: { + sessionId: 'session-2', + update: { + sessionUpdate: 'available_commands_update', + availableCommands: [ + { name: '/two', description: 'Second', input: null }, + ], + }, + }, + }); + + await waitFor(() => + expect(bridge.getAvailableCommands('session-2')).toEqual([ + { name: '/two', description: 'Second', input: null }, + ]), + ); + expect(bridge.getAvailableCommands('session-1')).toEqual([ + { name: '/one', description: 'First', input: null }, + ]); + expect(bridge.availableCommands).toEqual([ + { name: '/two', description: 'Second', input: null }, + ]); + + secondEvents.push({ + id: 3, + v: 1, + type: 'session_died', + data: { reason: 'gone' }, + }); + await waitFor(() => + expect(bridge.getAvailableCommands('session-2')).toEqual([]), + ); + expect(bridge.availableCommands).toEqual([ + { name: '/one', description: 'First', input: null }, + ]); + + firstEvents.close(); + secondEvents.close(); + bridge.stop(); + }); + + it('routes permission responses back through the owning daemon session', async () => { + const events = new EventQueue(); + const session = createFakeSession(events); + const bridge = new DaemonChannelBridge({ + cwd: '/repo', + sessionFactory: vi.fn().mockResolvedValue(session), + }); + const permissionRequest = vi.fn(); + bridge.on('permissionRequest', permissionRequest); + + await bridge.start(); + await bridge.newSession('/repo'); + + const request: RequestPermissionRequest & { requestId: string } = { + requestId: 'req-1', + sessionId: 'session-1', + toolCall: { + toolCallId: 'tool-1', + kind: 'edit', + title: 'Edit file', + rawInput: {}, + }, + options: [ + { optionId: 'proceed_once', kind: 'allow_once', name: 'Allow' }, + ], + } as RequestPermissionRequest & { requestId: string }; + events.push({ + id: 6, + v: 1, + type: 'permission_request', + data: request, + }); + + await waitFor(() => + expect(permissionRequest).toHaveBeenCalledWith({ + requestId: 'req-1', + sessionId: 'session-1', + request, + }), + ); + + const response: RequestPermissionResponse = { + outcome: { outcome: 'selected', optionId: 'proceed_once' }, + }; + await expect(bridge.respondToPermission('req-1', response)).resolves.toBe( + true, + ); + expect(session.respondToPermission).toHaveBeenCalledWith('req-1', response); + await expect(bridge.respondToPermission('req-1', response)).resolves.toBe( + false, + ); + + const resolved = vi.fn(); + bridge.on('permissionResolved', resolved); + events.push({ + id: 7, + v: 1, + type: 'permission_resolved', + data: { requestId: 'req-1', outcome: response.outcome }, + }); + await waitFor(() => + expect(resolved).toHaveBeenCalledWith({ + requestId: 'req-1', + outcome: response.outcome, + }), + ); + await expect(bridge.respondToPermission('req-1', response)).resolves.toBe( + false, + ); + + events.push({ + id: 8, + v: 1, + type: 'permission_request', + data: request, + }); + await waitFor(() => expect(permissionRequest).toHaveBeenCalledTimes(2)); + let staleResponse: Promise | undefined; + bridge.once('sessionDied', () => { + staleResponse = bridge.respondToPermission('req-1', response); + }); + events.push({ + id: 9, + v: 1, + type: 'session_died', + data: { reason: 'gone' }, + }); + await waitFor(() => expect(staleResponse).toBeDefined()); + await expect(staleResponse).resolves.toBe(false); + + events.close(); + bridge.stop(); + }); + + it('rejects malformed permission resolution outcomes', async () => { + const events = new EventQueue(); + const session = createFakeSession(events); + const bridge = new DaemonChannelBridge({ + cwd: '/repo', + sessionFactory: vi.fn().mockResolvedValue(session), + }); + const permissionRequest = vi.fn(); + const permissionResolved = vi.fn(); + const errors = vi.fn(); + bridge.on('permissionRequest', permissionRequest); + bridge.on('permissionResolved', permissionResolved); + bridge.on('error', errors); + + await bridge.start(); + await bridge.newSession('/repo'); + + events.push({ + id: 10, + v: 1, + type: 'permission_request', + data: { + requestId: 'req-bad-outcome', + toolCall: { + toolCallId: 'tool-1', + kind: 'edit', + title: 'Edit file', + rawInput: {}, + }, + options: [ + { optionId: 'proceed_once', kind: 'allow_once', name: 'Allow' }, + ], + }, + }); + await waitFor(() => expect(permissionRequest).toHaveBeenCalledOnce()); + + events.push({ + id: 11, + v: 1, + type: 'permission_resolved', + data: { + requestId: 'req-bad-outcome', + outcome: { outcome: 'selected' }, + }, + }); + + await waitFor(() => + expect(errors).toHaveBeenCalledWith( + expect.objectContaining({ + message: 'Malformed daemon permission_resolved outcome', + }), + ), + ); + expect(permissionResolved).not.toHaveBeenCalled(); + + events.close(); + bridge.stop(); + }); + + it('ignores permission resolution events from non-owning sessions', async () => { + const firstEvents = new EventQueue(); + const secondEvents = new EventQueue(); + const firstSession = createFakeSession(firstEvents, 'session-1'); + const secondSession = createFakeSession(secondEvents, 'session-2'); + const bridge = new DaemonChannelBridge({ + cwd: '/repo', + sessionFactory: vi + .fn() + .mockResolvedValueOnce(firstSession) + .mockResolvedValueOnce(secondSession), + }); + const permissionResolved = vi.fn(); + const permissionRequest = vi.fn(); + const errors = vi.fn(); + bridge.on('permissionRequest', permissionRequest); + bridge.on('permissionResolved', permissionResolved); + bridge.on('error', errors); + + await bridge.start(); + await bridge.newSession('/repo'); + await bridge.newSession('/repo'); + + firstEvents.push({ + id: 1, + v: 1, + type: 'permission_request', + data: { + requestId: 'req-1', + toolCall: { + toolCallId: 'tool-1', + kind: 'edit', + title: 'Edit file', + rawInput: {}, + }, + options: [ + { optionId: 'proceed_once', kind: 'allow_once', name: 'Allow' }, + ], + }, + }); + await waitFor(() => expect(permissionRequest).toHaveBeenCalledOnce()); + await expect( + bridge.respondToPermission('req-1', { + outcome: { outcome: 'selected', optionId: 'proceed_once' }, + }), + ).resolves.toBe(true); + + secondEvents.push({ + id: 2, + v: 1, + type: 'permission_resolved', + data: { requestId: 'req-1', outcome: { outcome: 'selected' } }, + }); + + await waitFor(() => + expect(errors).toHaveBeenCalledWith( + expect.objectContaining({ + message: expect.stringContaining('non-owning session session-2'), + }), + ), + ); + expect(permissionResolved).not.toHaveBeenCalled(); + expect(firstSession.respondToPermission).toHaveBeenCalledWith('req-1', { + outcome: { outcome: 'selected', optionId: 'proceed_once' }, + }); + expect(secondSession.respondToPermission).not.toHaveBeenCalled(); + await expect( + bridge.respondToPermission('req-1', { + outcome: { outcome: 'selected', optionId: 'proceed_once' }, + }), + ).resolves.toBe(false); + + firstEvents.close(); + secondEvents.close(); + bridge.stop(); + }); + + it('replaces duplicate daemon sessions and clears stale ownership state', async () => { + const firstEvents = new EventQueue(); + const secondEvents = new EventQueue(); + const firstSession = createFakeSession(firstEvents, 'session-1'); + firstSession.events.mockImplementation(() => firstEvents); + const secondSession = createFakeSession(secondEvents, 'session-1'); + secondSession.prompt.mockResolvedValue({ stopReason: 'end_turn' }); + const bridge = new DaemonChannelBridge({ + cwd: '/repo', + sessionFactory: vi + .fn() + .mockResolvedValueOnce(firstSession) + .mockResolvedValueOnce(secondSession), + }); + const sessionDied = vi.fn(); + const permissionRequest = vi.fn(); + bridge.on('sessionDied', sessionDied); + bridge.on('permissionRequest', permissionRequest); + + await bridge.start(); + await bridge.newSession('/repo'); + + firstEvents.push({ + id: 1, + v: 1, + type: 'permission_request', + data: { + requestId: 'req-1', + toolCall: { + toolCallId: 'tool-1', + kind: 'edit', + title: 'Edit file', + rawInput: {}, + }, + options: [ + { optionId: 'proceed_once', kind: 'allow_once', name: 'Allow' }, + ], + }, + }); + await waitFor(() => expect(permissionRequest).toHaveBeenCalledOnce()); + await expect( + bridge.respondToPermission('req-1', { + outcome: { outcome: 'selected', optionId: 'proceed_once' }, + }), + ).resolves.toBe(true); + + await expect(bridge.newSession('/repo')).resolves.toBe('session-1'); + + await waitFor(() => + expect(sessionDied).toHaveBeenCalledWith({ + sessionId: 'session-1', + reason: 'session_replaced', + }), + ); + await expect( + bridge.respondToPermission('req-1', { + outcome: { outcome: 'selected', optionId: 'proceed_once' }, + }), + ).resolves.toBe(false); + + firstEvents.push({ + id: 2, + v: 1, + type: 'session_died', + data: { reason: 'old pump finished late' }, + }); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(sessionDied).toHaveBeenCalledTimes(1); + await expect(bridge.prompt('session-1', 'still alive')).resolves.toBe(''); + expect(secondSession.prompt).toHaveBeenCalledOnce(); + + firstEvents.close(); + secondEvents.close(); + bridge.stop(); + }); + + it('rejects unknown sessions and concurrent prompts for one session', async () => { + const events = new EventQueue(); + const session = createFakeSession(events); + let resolvePrompt: () => void = () => {}; + session.prompt.mockImplementation( + () => + new Promise((resolve) => { + resolvePrompt = () => resolve({ stopReason: 'end_turn' }); + }), + ); + const bridge = new DaemonChannelBridge({ + cwd: '/repo', + sessionFactory: vi.fn().mockResolvedValue(session), + }); + const promptComplete = vi.fn(); + bridge.on('promptComplete', promptComplete); + + await bridge.start(); + await bridge.newSession('/repo'); + + await expect(bridge.cancelSession('missing')).rejects.toThrow( + 'No daemon session bound for missing', + ); + await expect(bridge.prompt('missing', 'hello')).rejects.toThrow( + 'No daemon session bound for missing', + ); + await expect( + bridge.setSessionModel('missing', 'qwen3-coder-plus'), + ).rejects.toThrow('No daemon session bound for missing'); + + const firstPrompt = bridge.prompt('session-1', 'first'); + await waitFor(() => expect(session.prompt).toHaveBeenCalledOnce()); + await expect(bridge.prompt('session-1', 'second')).rejects.toThrow( + 'Prompt already in flight for daemon session session-1', + ); + resolvePrompt(); + await expect(firstPrompt).resolves.toBe(''); + expect(promptComplete).toHaveBeenCalledWith({ + sessionId: 'session-1', + text: '', + stopReason: 'end_turn', + }); + + events.close(); + bridge.stop(); + }); + + it('passes image prompt blocks and aborts prompts when a session dies', async () => { + const events = new EventQueue(); + const session = createFakeSession(events); + session.prompt.mockImplementation( + (_req: unknown, signal?: AbortSignal) => + new Promise((_resolve, reject) => { + signal?.addEventListener( + 'abort', + () => reject(new DOMException('aborted', 'AbortError')), + { once: true }, + ); + }), + ); + const bridge = new DaemonChannelBridge({ + cwd: '/repo', + sessionFactory: vi.fn().mockResolvedValue(session), + }); + + await bridge.start(); + await bridge.newSession('/repo'); + + const promptPromise = bridge.prompt('session-1', 'describe', { + imageBase64: 'base64-image', + imageMimeType: 'image/png', + }); + await waitFor(() => expect(session.prompt).toHaveBeenCalledOnce()); + expect(session.prompt).toHaveBeenCalledWith( + { + prompt: [ + { type: 'image', data: 'base64-image', mimeType: 'image/png' }, + { type: 'text', text: 'describe' }, + ], + }, + expect.any(AbortSignal), + ); + + events.push({ + id: 10, + v: 1, + type: 'session_died', + data: { reason: 'agent exited' }, + }); + await expect(promptPromise).rejects.toThrow('aborted'); + + events.close(); + bridge.stop(); + }); + + it('aborts in-flight prompts when the bridge stops', async () => { + const events = new EventQueue(); + const session = createFakeSession(events); + session.prompt.mockImplementation( + (_req: unknown, signal?: AbortSignal) => + new Promise((_resolve, reject) => { + signal?.addEventListener( + 'abort', + () => reject(new DOMException('aborted', 'AbortError')), + { once: true }, + ); + }), + ); + const bridge = new DaemonChannelBridge({ + cwd: '/repo', + sessionFactory: vi.fn().mockResolvedValue(session), + }); + const sessionDied = vi.fn(); + bridge.on('sessionDied', sessionDied); + + await bridge.start(); + await bridge.newSession('/repo'); + const promptPromise = bridge.prompt('session-1', 'hello'); + await waitFor(() => expect(session.prompt).toHaveBeenCalledOnce()); + + bridge.stop(); + await expect(promptPromise).rejects.toThrow('aborted'); + expect(session.cancel).toHaveBeenCalledOnce(); + expect(sessionDied).toHaveBeenCalledWith({ + sessionId: 'session-1', + reason: 'bridge_stopped', + }); + }); + + it('aborts in-flight prompts when cancelling a session', async () => { + const events = new EventQueue(); + const session = createFakeSession(events); + const order: string[] = []; + session.cancel.mockImplementation(async () => { + order.push('cancel'); + }); + session.prompt.mockImplementation( + (_req: unknown, signal?: AbortSignal) => + new Promise((_resolve, reject) => { + signal?.addEventListener( + 'abort', + () => { + order.push('abort'); + reject(new DOMException('aborted', 'AbortError')); + }, + { once: true }, + ); + }), + ); + const bridge = new DaemonChannelBridge({ + cwd: '/repo', + sessionFactory: vi.fn().mockResolvedValue(session), + }); + + await bridge.start(); + await bridge.newSession('/repo'); + const promptPromise = bridge.prompt('session-1', 'hello'); + await waitFor(() => expect(session.prompt).toHaveBeenCalledOnce()); + + await bridge.cancelSession('session-1'); + + await expect(promptPromise).rejects.toThrow('aborted'); + expect(session.cancel).toHaveBeenCalledOnce(); + expect(order).toEqual(['cancel', 'abort']); + + events.close(); + bridge.stop(); + }); + + it('clears permission ownership when daemon permission responses fail', async () => { + const events = new EventQueue(); + const session = createFakeSession(events); + const bridge = new DaemonChannelBridge({ + cwd: '/repo', + sessionFactory: vi.fn().mockResolvedValue(session), + }); + const permissionRequest = vi.fn(); + bridge.on('permissionRequest', permissionRequest); + + await bridge.start(); + await bridge.newSession('/repo'); + + events.push({ + id: 1, + v: 1, + type: 'permission_request', + data: { + requestId: 'req-fail', + toolCall: { + toolCallId: 'tool-1', + kind: 'edit', + title: 'Edit file', + rawInput: {}, + }, + options: [ + { optionId: 'proceed_once', kind: 'allow_once', name: 'Allow' }, + ], + }, + }); + await waitFor(() => expect(permissionRequest).toHaveBeenCalledOnce()); + + session.respondToPermission.mockRejectedValueOnce( + new Error('permission failed'), + ); + const response: RequestPermissionResponse = { + outcome: { outcome: 'selected', optionId: 'proceed_once' }, + }; + await expect( + bridge.respondToPermission('req-fail', response), + ).rejects.toThrow('permission failed'); + await expect( + bridge.respondToPermission('req-fail', response), + ).resolves.toBe(false); + + events.close(); + bridge.stop(); + }); + + it('treats terminal stream frames and completion as session death', async () => { + const failedEvents = new EventQueue(); + failedEvents.fail(new Error('network down')); + const failedSession = createFakeSession(failedEvents); + const failedBridge = new DaemonChannelBridge({ + cwd: '/repo', + sessionFactory: vi.fn().mockResolvedValue(failedSession), + }); + const failedDied = vi.fn(); + failedBridge.on('sessionDied', failedDied); + + await failedBridge.start(); + await failedBridge.newSession('/repo'); + await waitFor(() => + expect(failedDied).toHaveBeenCalledWith({ + sessionId: 'session-1', + reason: 'network down', + }), + ); + expect(failedBridge.lastDaemonError).toMatchObject({ + message: 'network down', + }); + await expect(failedBridge.prompt('session-1', 'hello')).rejects.toThrow( + 'No daemon session bound for session-1', + ); + + const endedEvents = new EventQueue(); + const endedSession = createFakeSession(endedEvents); + const endedBridge = new DaemonChannelBridge({ + cwd: '/repo', + sessionFactory: vi.fn().mockResolvedValue(endedSession), + }); + const endedDied = vi.fn(); + endedBridge.on('sessionDied', endedDied); + + await endedBridge.start(); + await endedBridge.newSession('/repo'); + endedEvents.close(); + await waitFor(() => + expect(endedDied).toHaveBeenCalledWith({ + sessionId: 'session-1', + reason: 'stream_ended', + }), + ); + + const terminalEvents = new EventQueue(); + const terminalSession = createFakeSession(terminalEvents); + const terminalBridge = new DaemonChannelBridge({ + cwd: '/repo', + sessionFactory: vi.fn().mockResolvedValue(terminalSession), + }); + const terminalDied = vi.fn(); + terminalBridge.on('sessionDied', terminalDied); + + await terminalBridge.start(); + await terminalBridge.newSession('/repo'); + terminalEvents.push({ + id: 20, + v: 1, + type: 'stream_error', + data: { error: 'subscriber limit reached' }, + }); + await waitFor(() => + expect(terminalDied).toHaveBeenCalledWith({ + sessionId: 'session-1', + reason: 'subscriber limit reached', + }), + ); + await expect(terminalBridge.prompt('session-1', 'hello')).rejects.toThrow( + 'No daemon session bound for session-1', + ); + + const evictedEvents = new EventQueue(); + const evictedSession = createFakeSession(evictedEvents); + const evictedBridge = new DaemonChannelBridge({ + cwd: '/repo', + sessionFactory: vi.fn().mockResolvedValue(evictedSession), + }); + const evictedDied = vi.fn(); + evictedBridge.on('sessionDied', evictedDied); + + await evictedBridge.start(); + await evictedBridge.newSession('/repo'); + evictedEvents.push({ + id: 21, + v: 1, + type: 'client_evicted', + data: { reason: 'queue_overflow' }, + }); + await waitFor(() => + expect(evictedDied).toHaveBeenCalledWith({ + sessionId: 'session-1', + reason: 'queue_overflow', + }), + ); + }); + + it('loads an existing daemon session and forwards cancel/model changes', async () => { + const events = new EventQueue(); + const session = createFakeSession(events, 'existing-session'); + const factory = vi.fn().mockResolvedValue(session); + const bridge = new DaemonChannelBridge({ + cwd: '/repo', + modelServiceId: 'default', + sessionScope: 'user', + sessionFactory: factory, + }); + + await bridge.start(); + await expect(bridge.loadSession('existing-session', '/repo')).resolves.toBe( + 'existing-session', + ); + await bridge.cancelSession('existing-session'); + await bridge.setSessionModel('existing-session', 'qwen3-coder-plus'); + + expect(factory).toHaveBeenCalledWith({ + workspaceCwd: '/repo', + modelServiceId: 'default', + sessionId: 'existing-session', + sessionScope: 'user', + }); + expect(session.cancel).toHaveBeenCalledOnce(); + expect(session.setModel).toHaveBeenCalledWith('qwen3-coder-plus'); + + events.close(); + bridge.stop(); + }); + + it('rejects mismatched daemon session ids while loading', async () => { + const events = new EventQueue(); + const session = createFakeSession(events, 'different-session'); + const bridge = new DaemonChannelBridge({ + cwd: '/repo', + sessionFactory: vi.fn().mockResolvedValue(session), + }); + + await bridge.start(); + await expect( + bridge.loadSession('existing-session', '/repo'), + ).rejects.toThrow( + 'Daemon returned session different-session while loading existing-session', + ); + await expect(bridge.prompt('different-session', 'hello')).rejects.toThrow( + 'No daemon session bound for different-session', + ); + + events.close(); + bridge.stop(); + }); + + it('surfaces malformed daemon events through the error channel', async () => { + const events = new EventQueue(); + const session = createFakeSession(events); + const bridge = new DaemonChannelBridge({ + cwd: '/repo', + sessionFactory: vi.fn().mockResolvedValue(session), + }); + const errors = vi.fn(); + bridge.on('error', errors); + + await bridge.start(); + await bridge.newSession('/repo'); + + events.push({ + id: 1, + v: 1, + type: 'permission_request', + data: { requestId: 'req-1' }, + }); + events.push({ + id: 2, + v: 1, + type: 'session_update', + data: { + update: { + sessionUpdate: 'available_commands_update', + availableCommands: 'not-an-array', + }, + }, + }); + events.push({ + id: 3, + v: 1, + type: 'model_switched', + data: {}, + }); + events.push({ + id: 4, + v: 1, + type: 'session_update', + data: { + update: { + sessionUpdate: 'tool_call_update', + status: 'running', + }, + }, + }); + + await waitFor(() => expect(errors).toHaveBeenCalledTimes(4)); + expect(errors).toHaveBeenNthCalledWith( + 1, + expect.objectContaining({ + message: 'Malformed daemon permission_request event', + }), + ); + expect(errors).toHaveBeenNthCalledWith( + 2, + expect.objectContaining({ + message: 'Malformed daemon available_commands_update event', + }), + ); + expect(errors).toHaveBeenNthCalledWith( + 3, + expect.objectContaining({ + message: 'Malformed daemon model_switched event', + }), + ); + expect(errors).toHaveBeenNthCalledWith( + 4, + expect.objectContaining({ + message: 'Malformed daemon tool_call_update event', + }), + ); + expect(bridge.lastDaemonError).toMatchObject({ + message: 'Malformed daemon tool_call_update event', + }); + + events.close(); + bridge.stop(); + }); +}); diff --git a/packages/channels/base/src/DaemonChannelBridge.ts b/packages/channels/base/src/DaemonChannelBridge.ts new file mode 100644 index 00000000000..04c04fe1e91 --- /dev/null +++ b/packages/channels/base/src/DaemonChannelBridge.ts @@ -0,0 +1,705 @@ +import { EventEmitter } from 'node:events'; +import type { + RequestPermissionRequest, + RequestPermissionResponse, +} from '@agentclientprotocol/sdk'; +import type { AvailableCommand, ToolCallEvent } from './AcpBridge.js'; +import type { SessionScope } from './types.js'; + +const MAX_RESPONDED_PERMISSION_REQUESTS = 256; + +export interface DaemonChannelEvent { + id?: number; + v: 1; + type: string; + data: unknown; + originatorClientId?: string; +} + +export interface DaemonChannelSessionClient { + readonly sessionId: string; + readonly workspaceCwd: string; + readonly lastEventId?: number; + prompt( + req: { + prompt: Array>; + }, + signal?: AbortSignal, + ): Promise<{ stopReason?: string; [key: string]: unknown }>; + events(opts?: { + signal?: AbortSignal; + lastEventId?: number; + resume?: boolean; + }): AsyncGenerator; + cancel(): Promise; + setModel(modelId: string): Promise>; + respondToPermission( + requestId: string, + response: RequestPermissionResponse, + ): Promise; +} + +export interface DaemonChannelSessionFactoryRequest { + workspaceCwd: string; + modelServiceId?: string; + sessionId?: string; + sessionScope?: SessionScope; +} + +export type DaemonChannelSessionFactory = ( + req: DaemonChannelSessionFactoryRequest, +) => Promise; + +export interface DaemonChannelBridgeOptions { + cwd: string; + sessionFactory: DaemonChannelSessionFactory; + modelServiceId?: string; + sessionScope?: SessionScope; +} + +export interface DaemonPermissionRequestEvent { + requestId: string; + sessionId: string; + request: RequestPermissionRequest; +} + +export interface DaemonPermissionResolvedEvent { + requestId: string; + outcome?: DaemonPermissionOutcome; +} + +export interface DaemonPromptCompleteEvent { + sessionId: string; + text: string; + stopReason?: string; +} + +function isRecord(value: unknown): value is Record { + return typeof value === 'object' && value !== null; +} + +function getString(value: unknown): string | undefined { + return typeof value === 'string' ? value : undefined; +} + +function getTextContent(content: unknown): string | undefined { + if (!isRecord(content)) { + return undefined; + } + return getString(content['text']); +} + +function getSessionUpdate(data: unknown): Record | undefined { + if (!isRecord(data) || !isRecord(data['update'])) { + return undefined; + } + return data['update']; +} + +function isAvailableCommand(value: unknown): value is AvailableCommand { + return isRecord(value) && typeof value['name'] === 'string'; +} + +function isPermissionRequestData( + value: unknown, +): value is RequestPermissionRequest & { requestId: string } { + if ( + !isRecord(value) || + typeof value['requestId'] !== 'string' || + !isRecord(value['toolCall']) || + typeof value['toolCall']['toolCallId'] !== 'string' || + typeof value['toolCall']['kind'] !== 'string' || + !Array.isArray(value['options']) + ) { + return false; + } + return value['options'].every( + (option) => isRecord(option) && typeof option['optionId'] === 'string', + ); +} + +type DaemonPermissionOutcome = + | { outcome: 'cancelled' } + | { outcome: 'selected'; optionId: string }; + +function parsePermissionOutcome( + value: unknown, +): DaemonPermissionOutcome | undefined { + if (!isRecord(value)) { + return undefined; + } + if (value['outcome'] === 'cancelled') { + return { outcome: 'cancelled' }; + } + if ( + value['outcome'] === 'selected' && + typeof value['optionId'] === 'string' + ) { + return { outcome: 'selected', optionId: value['optionId'] }; + } + return undefined; +} + +function summarizeProtocolDetails(details: unknown): unknown { + if (!isRecord(details)) { + return { type: typeof details }; + } + const summary: Record = {}; + for (const key of [ + 'requestId', + 'sessionId', + 'sessionUpdate', + 'modelId', + 'requestedModelId', + 'toolCallId', + 'kind', + ]) { + const value = details[key]; + if (typeof value === 'string') { + summary[key] = value; + } + } + return summary; +} + +async function drainDaemonEventLoop(): Promise { + // TODO(daemon-roadmap): replace this bounded client-side drain with a daemon + // terminal turn event / SSE waterline once the typed event schema defines it. + await new Promise((resolve) => setTimeout(resolve, 0)); +} + +export class DaemonChannelBridge extends EventEmitter { + private readonly options: DaemonChannelBridgeOptions; + private readonly sessions = new Map(); + private readonly eventControllers = new Map(); + private readonly requestToSession = new Map(); + private readonly respondedRequestToSession = new Map(); + private readonly activePrompts = new Set(); + private readonly activePromptControllers = new Map< + string, + Set + >(); + private readonly availableCommandsBySession = new Map< + string, + AvailableCommand[] + >(); + private connected = false; + private latestAvailableCommandsSessionId: string | undefined; + private lastError: unknown; + + constructor(options: DaemonChannelBridgeOptions) { + super(); + this.options = options; + this.on('error', (error) => { + this.lastError = error; + }); + } + + get availableCommands(): AvailableCommand[] { + if (this.latestAvailableCommandsSessionId) { + return ( + this.availableCommandsBySession.get( + this.latestAvailableCommandsSessionId, + ) ?? [] + ); + } + return Array.from(this.availableCommandsBySession.values()).at(-1) ?? []; + } + + get lastDaemonError(): unknown { + return this.lastError; + } + + getAvailableCommands(sessionId: string): AvailableCommand[] { + return this.availableCommandsBySession.get(sessionId) ?? []; + } + + async start(): Promise { + this.connected = true; + } + + async newSession(cwd: string): Promise { + const session = await this.options.sessionFactory({ + workspaceCwd: cwd || this.options.cwd, + modelServiceId: this.options.modelServiceId, + sessionScope: this.options.sessionScope ?? 'thread', + }); + this.attachSession(session); + return session.sessionId; + } + + async loadSession(sessionId: string, cwd: string): Promise { + const session = await this.options.sessionFactory({ + workspaceCwd: cwd || this.options.cwd, + modelServiceId: this.options.modelServiceId, + sessionId, + sessionScope: this.options.sessionScope ?? 'thread', + }); + if (session.sessionId !== sessionId) { + throw new Error( + `Daemon returned session ${session.sessionId} while loading ${sessionId}`, + ); + } + this.attachSession(session); + return session.sessionId; + } + + async prompt( + sessionId: string, + text: string, + options?: { imageBase64?: string; imageMimeType?: string }, + ): Promise { + const session = this.ensureSession(sessionId); + if (this.activePrompts.has(sessionId)) { + throw new Error( + `Prompt already in flight for daemon session ${sessionId}`, + ); + } + this.activePrompts.add(sessionId); + + const controller = new AbortController(); + let controllers = this.activePromptControllers.get(sessionId); + if (!controllers) { + controllers = new Set(); + this.activePromptControllers.set(sessionId, controllers); + } + controllers.add(controller); + + const chunks: string[] = []; + const onChunk = (sid: string, chunk: string) => { + if (sid === sessionId) { + chunks.push(chunk); + } + }; + const onSessionDied = (info: { sessionId: string }) => { + if (info.sessionId === sessionId) { + controller.abort(); + } + }; + this.on('textChunk', onChunk); + this.on('sessionDied', onSessionDied); + + const prompt: Array> = []; + if (options?.imageBase64 && options.imageMimeType) { + prompt.push({ + type: 'image', + data: options.imageBase64, + mimeType: options.imageMimeType, + }); + } + prompt.push({ type: 'text', text }); + + try { + const result = await session.prompt({ prompt }, controller.signal); + await drainDaemonEventLoop(); + const textResult = chunks.join(''); + this.emit('promptComplete', { + sessionId, + text: textResult, + stopReason: result.stopReason, + } satisfies DaemonPromptCompleteEvent); + return textResult; + } finally { + this.off('textChunk', onChunk); + this.off('sessionDied', onSessionDied); + this.activePrompts.delete(sessionId); + controllers.delete(controller); + if ( + controllers.size === 0 && + this.activePromptControllers.get(sessionId) === controllers + ) { + this.activePromptControllers.delete(sessionId); + } + } + } + + async cancelSession(sessionId: string): Promise { + const session = this.ensureSession(sessionId); + await session.cancel(); + this.abortActivePrompts(sessionId); + this.activePrompts.delete(sessionId); + } + + async setSessionModel( + sessionId: string, + modelId: string, + ): Promise> { + return await this.ensureSession(sessionId).setModel(modelId); + } + + async respondToPermission( + requestId: string, + response: RequestPermissionResponse, + ): Promise { + const sessionId = this.requestToSession.get(requestId); + if (!sessionId) { + return false; + } + const session = this.sessions.get(sessionId); + if (!session) { + this.requestToSession.delete(requestId); + this.respondedRequestToSession.delete(requestId); + return false; + } + try { + const accepted = await session.respondToPermission(requestId, response); + this.requestToSession.delete(requestId); + if (accepted) { + this.rememberRespondedPermissionRequest(requestId, sessionId); + } else { + this.respondedRequestToSession.delete(requestId); + } + return accepted; + } catch (error) { + this.requestToSession.delete(requestId); + this.respondedRequestToSession.delete(requestId); + throw error; + } + } + + stop(): void { + for (const sessionId of Array.from(this.sessions.keys())) { + const session = this.sessions.get(sessionId); + if (session) { + void session.cancel().catch((error: unknown) => { + this.lastError = error; + }); + } + this.dropSession(sessionId, 'bridge_stopped'); + } + this.latestAvailableCommandsSessionId = undefined; + this.connected = false; + } + + get isConnected(): boolean { + return this.connected; + } + + private attachSession(session: DaemonChannelSessionClient): void { + if (this.sessions.has(session.sessionId)) { + this.dropSession(session.sessionId, 'session_replaced'); + } + + this.sessions.set(session.sessionId, session); + const controller = new AbortController(); + this.eventControllers.set(session.sessionId, controller); + void this.pumpEvents(session, controller.signal); + } + + private ensureSession(sessionId: string): DaemonChannelSessionClient { + const session = this.sessions.get(sessionId); + if (!session) { + throw new Error(`No daemon session bound for ${sessionId}`); + } + return session; + } + + private async pumpEvents( + session: DaemonChannelSessionClient, + signal: AbortSignal, + ): Promise { + try { + for await (const event of session.events({ + signal, + lastEventId: session.lastEventId, + resume: true, + })) { + if (!this.isCurrentPump(session, signal)) { + return; + } + this.handleEvent(session, event); + } + if (!signal.aborted && this.isCurrentPump(session, signal)) { + this.dropSession(session.sessionId, 'stream_ended'); + } + } catch (error) { + if (!signal.aborted && this.isCurrentPump(session, signal)) { + this.emit('error', error); + this.dropSession( + session.sessionId, + error instanceof Error ? error.message : String(error), + ); + } + } + } + + private isCurrentPump( + session: DaemonChannelSessionClient, + signal: AbortSignal, + ): boolean { + return ( + this.sessions.get(session.sessionId) === session && + this.eventControllers.get(session.sessionId)?.signal === signal + ); + } + + private handleEvent( + session: DaemonChannelSessionClient, + event: DaemonChannelEvent, + ): void { + switch (event.type) { + case 'session_update': + this.handleSessionUpdate(session.sessionId, event.data); + break; + case 'permission_request': + this.handlePermissionRequest(session.sessionId, event.data); + break; + case 'permission_resolved': + this.handlePermissionResolved(session.sessionId, event.data); + break; + case 'model_switched': + this.handleModelSwitched(session.sessionId, event.data); + break; + case 'model_switch_failed': + this.handleModelSwitchFailed(session.sessionId, event.data); + break; + case 'session_died': + this.handleSessionDied(session.sessionId, event.data); + break; + case 'client_evicted': + this.dropSession( + session.sessionId, + this.getReason(event.data, 'client_evicted'), + ); + break; + case 'stream_error': + this.dropSession( + session.sessionId, + this.getError(event.data, 'stream_error'), + ); + break; + default: + break; + } + } + + private handleSessionUpdate(sessionId: string, data: unknown): void { + const update = getSessionUpdate(data); + if (!update) { + this.emitProtocolError('Malformed daemon session_update event', data); + return; + } + + const type = getString(update['sessionUpdate']); + switch (type) { + case 'agent_message_chunk': { + const text = getTextContent(update['content']); + if (text) { + this.emit('textChunk', sessionId, text); + } + break; + } + case 'agent_thought_chunk': { + const text = getTextContent(update['content']); + if (text) { + this.emit('thoughtChunk', sessionId, text); + } + break; + } + case 'tool_call': + case 'tool_call_update': { + const toolCallId = getString(update['toolCallId']); + const kind = getString(update['kind']); + if (!toolCallId || !kind) { + this.emitProtocolError(`Malformed daemon ${type} event`, update); + break; + } + const event: ToolCallEvent = { + sessionId, + toolCallId, + kind, + title: getString(update['title']) ?? '', + status: getString(update['status']) ?? 'pending', + rawInput: isRecord(update['rawInput']) + ? update['rawInput'] + : undefined, + }; + this.emit('toolCall', event); + break; + } + case 'available_commands_update': { + if (Array.isArray(update['availableCommands'])) { + const commands = + update['availableCommands'].filter(isAvailableCommand); + this.availableCommandsBySession.set(sessionId, commands); + this.latestAvailableCommandsSessionId = sessionId; + } else { + this.emitProtocolError( + 'Malformed daemon available_commands_update event', + data, + ); + } + break; + } + default: + break; + } + + this.emit('sessionUpdate', data); + } + + private handlePermissionRequest(sessionId: string, data: unknown): void { + if (!isPermissionRequestData(data)) { + this.emitProtocolError('Malformed daemon permission_request event', data); + return; + } + const requestId = data['requestId']; + this.requestToSession.set(requestId, sessionId); + this.emit('permissionRequest', { + requestId, + sessionId, + request: data as unknown as RequestPermissionRequest, + } satisfies DaemonPermissionRequestEvent); + } + + private rememberRespondedPermissionRequest( + requestId: string, + sessionId: string, + ): void { + this.respondedRequestToSession.set(requestId, sessionId); + while ( + this.respondedRequestToSession.size > MAX_RESPONDED_PERMISSION_REQUESTS + ) { + const oldestRequestId = this.respondedRequestToSession + .keys() + .next().value; + if (oldestRequestId === undefined) { + return; + } + this.respondedRequestToSession.delete(oldestRequestId); + } + } + + private handlePermissionResolved(sessionId: string, data: unknown): void { + if (!isRecord(data) || typeof data['requestId'] !== 'string') { + this.emitProtocolError( + 'Malformed daemon permission_resolved event', + data, + ); + return; + } + const requestId = data['requestId']; + const mappedSessionId = + this.requestToSession.get(requestId) ?? + this.respondedRequestToSession.get(requestId); + if (!mappedSessionId) { + this.emitProtocolError( + `Ignoring daemon permission_resolved for unknown request ${requestId}`, + data, + ); + return; + } + if (mappedSessionId !== sessionId) { + this.requestToSession.delete(requestId); + this.respondedRequestToSession.delete(requestId); + this.emitProtocolError( + `Ignoring daemon permission_resolved for request ${requestId} from non-owning session ${sessionId}`, + data, + ); + return; + } + const outcome = parsePermissionOutcome(data['outcome']); + if (!outcome) { + this.requestToSession.delete(requestId); + this.respondedRequestToSession.delete(requestId); + this.emitProtocolError( + 'Malformed daemon permission_resolved outcome', + data, + ); + return; + } + this.requestToSession.delete(requestId); + this.respondedRequestToSession.delete(requestId); + this.emit('permissionResolved', { + requestId, + outcome, + } satisfies DaemonPermissionResolvedEvent); + } + + private handleModelSwitched(sessionId: string, data: unknown): void { + if (!isRecord(data) || typeof data['modelId'] !== 'string') { + this.emitProtocolError('Malformed daemon model_switched event', data); + return; + } + this.emit('modelSwitched', { + sessionId, + modelId: data['modelId'], + }); + } + + private handleModelSwitchFailed(sessionId: string, data: unknown): void { + if (!isRecord(data)) { + this.emitProtocolError( + 'Malformed daemon model_switch_failed event', + data, + ); + return; + } + this.emit('modelSwitchFailed', { + sessionId, + requestedModelId: getString(data['requestedModelId']), + error: getString(data['error']) ?? 'model_switch_failed', + }); + } + + private handleSessionDied(sessionId: string, data: unknown): void { + this.dropSession(sessionId, this.getReason(data, 'session_died')); + } + + private dropSession(sessionId: string, reason: string): void { + if (!this.sessions.has(sessionId)) { + return; + } + this.eventControllers.get(sessionId)?.abort(); + this.eventControllers.delete(sessionId); + this.sessions.delete(sessionId); + this.abortActivePrompts(sessionId); + this.activePrompts.delete(sessionId); + this.availableCommandsBySession.delete(sessionId); + if (this.latestAvailableCommandsSessionId === sessionId) { + this.latestAvailableCommandsSessionId = Array.from( + this.availableCommandsBySession.keys(), + ).at(-1); + } + for (const [requestId, mappedSessionId] of this.requestToSession) { + if (mappedSessionId === sessionId) { + this.requestToSession.delete(requestId); + } + } + for (const [requestId, mappedSessionId] of this.respondedRequestToSession) { + if (mappedSessionId === sessionId) { + this.respondedRequestToSession.delete(requestId); + } + } + this.emit('sessionDied', { sessionId, reason }); + } + + private getReason(data: unknown, fallback: string): string { + return isRecord(data) && typeof data['reason'] === 'string' + ? data['reason'] + : fallback; + } + + private getError(data: unknown, fallback: string): string { + return isRecord(data) && typeof data['error'] === 'string' + ? data['error'] + : fallback; + } + + private abortActivePrompts(sessionId: string): void { + const promptControllers = this.activePromptControllers.get(sessionId); + if (!promptControllers) { + return; + } + for (const controller of promptControllers) { + controller.abort(); + } + this.activePromptControllers.delete(sessionId); + } + + private emitProtocolError(message: string, details: unknown): void { + const error = new Error(message) as Error & { details?: unknown }; + error.details = summarizeProtocolDetails(details); + this.emit('error', error); + } +} diff --git a/packages/channels/base/src/index.ts b/packages/channels/base/src/index.ts index b93c8e3e9bf..bf5a39f7a0b 100644 --- a/packages/channels/base/src/index.ts +++ b/packages/channels/base/src/index.ts @@ -5,6 +5,17 @@ export type { AvailableCommand, ToolCallEvent, } from './AcpBridge.js'; +export { DaemonChannelBridge } from './DaemonChannelBridge.js'; +export type { + DaemonChannelBridgeOptions, + DaemonChannelEvent, + DaemonChannelSessionClient, + DaemonChannelSessionFactory, + DaemonChannelSessionFactoryRequest, + DaemonPromptCompleteEvent, + DaemonPermissionRequestEvent, + DaemonPermissionResolvedEvent, +} from './DaemonChannelBridge.js'; export { BlockStreamer } from './BlockStreamer.js'; export type { BlockStreamerOptions } from './BlockStreamer.js'; export { ChannelBase } from './ChannelBase.js';