diff --git a/docs/design/daemon-multi-workspace-phase4b-voice.md b/docs/design/daemon-multi-workspace-phase4b-voice.md new file mode 100644 index 00000000000..9488b8131e3 --- /dev/null +++ b/docs/design/daemon-multi-workspace-phase4b-voice.md @@ -0,0 +1,41 @@ +# Workspace-qualified Voice + +## Goal + +Expose the existing daemon Voice settings, batch transcription, and streaming +transcription surfaces for every trusted workspace runtime without changing +legacy primary-only routes. + +## Design + +`GET`/`POST /workspaces/:workspace/voice`, +`POST /workspaces/:workspace/voice/transcribe`, and +`WS /workspaces/:workspace/voice/stream` resolve a registered trusted runtime +by id or encoded cwd. They use that runtime's cwd, effective environment, +bridge, and workspace settings. Voice setting writes through plural REST always +use workspace scope; secondary ACP voice writes use the same scope so they +cannot mutate shared user settings. + +One process-scoped `WorkspaceVoiceCoordinator` owns the existing limit of +eight active Voice operations. It accounts for both WebSocket and REST batch +work across legacy and workspace-qualified paths. A removal drain rejects new +admission but leaves existing Voice work visible to the non-force removal +activity snapshot. Runtime disposal aborts only the selected runtime's Voice +leases before its bridge is shut down. + +## Compatibility + +Legacy `/workspace/voice`, `/workspace/voice/transcribe`, and `/voice/stream` +remain bound to the primary workspace. ACP method names and Voice settings +schema are unchanged. `workspace_qualified_voice` advertises all qualified +Voice modalities when the shared ACP/Voice WebSocket listener is enabled. The +existing Voice modality capability tags remain +primary-workspace signals and are not prerequisites for a secondary runtime, +whose configuration is validated by the selected route. + +Unknown workspace selectors return `400 workspace_mismatch`; registered but +untrusted runtimes return `403 untrusted_workspace` before Voice settings or +audio are read. The shared eight-operation admission cap covers batch and +streaming work for both legacy and plural routes. Batch capacity failures return +`503 voice_capacity_exceeded` with `Retry-After: 5`; streaming capacity failures +send an error frame and close with code `1013`. diff --git a/docs/developers/qwen-serve-protocol.md b/docs/developers/qwen-serve-protocol.md index 7dcb2cfa50a..9fd2cb64569 100644 --- a/docs/developers/qwen-serve-protocol.md +++ b/docs/developers/qwen-serve-protocol.md @@ -189,8 +189,8 @@ registry. Clients **must** gate UI off `features`, not off `mode` (per design 'session_branch', 'rate_limit', 'workspace_reload', 'multi_workspace_sessions', 'multi_workspace_session_rewind', 'multi_workspace_session_shell', 'persistent_workspace_registration', - 'workspace_qualified_rest_core', 'extension_management_v2', - 'workspace_persisted_transcript', + 'workspace_qualified_rest_core', 'workspace_qualified_voice', + 'extension_management_v2', 'workspace_persisted_transcript', 'client_mcp_over_ws', 'cdp_tunnel_over_ws', 'browser_automation_mcp'] ``` @@ -222,6 +222,8 @@ registry. Clients **must** gate UI off `features`, not off `mode` (per design `workspace_qualified_rest_core` advertises plural core REST routes under `/workspaces/:workspace/...`. The selector resolves as exact workspace id first, then as a URL-encoded absolute cwd after canonicalization. Newer single-workspace daemons include the primary runtime in `workspaces[]` even when `multi_workspace_sessions` is absent, allowing clients to discover the id required by workspace-qualified routes; clients should fall back to `capabilities.workspaceCwd` for older daemons that omit the array. Trust status and trust request routes are available for registered untrusted workspaces; file read routes follow the existing filesystem read policy. Registered untrusted secondary workspaces also expose persisted-only session and session-group catalogs: these reads do not attach to a session, start ACP, or merge live bridge state. File writes, catalog mutations, and other plural core routes require a trusted workspace unless a separate capability explicitly defines a narrower read-only policy, such as `workspace_persisted_transcript`. An untrusted primary continues to receive `403 { code: "untrusted_workspace" }` from the plural catalog and transcript routes; legacy singular primary routes keep their existing compatibility behavior. This tag covers the core file, status, settings, permissions, trust, lifecycle, MCP control, tool and skill toggles, memory, workspace agent CRUD, and session storage surfaces. It does not cover auth, voice, extensions, ACP/WebSocket transport, or channel-worker routing. Workspace trust is not an ACL: a client holding the daemon token can read every registered workspace surface allowed by this policy. +`workspace_qualified_voice` advertises Voice routes selected by a trusted workspace runtime: `GET` and `POST /workspaces/:workspace/voice`, `POST /workspaces/:workspace/voice/transcribe`, and `WS /workspaces/:workspace/voice/stream`. It is advertised only when multi-workspace runtimes and the shared ACP/Voice WebSocket listener are both enabled. The selector follows the same id-or-encoded-absolute-cwd rules as other plural routes. For REST, an unknown selector returns `400 { code: "workspace_mismatch" }` and an untrusted selector returns `403 { code: "untrusted_workspace" }`; WebSocket upgrade rejection exposes the corresponding HTTP 400/403 status without a structured JSON envelope. Neither transport falls back to primary. Legacy `/workspace/voice`, `/workspace/voice/transcribe`, and `/voice/stream` remain primary-only. Clients use `workspace_qualified_voice` for all qualified Voice modalities and let the selected runtime report configuration-specific errors. The legacy `workspace_voice`, `workspace_voice_transcription`, and `voice_transcribe` tags describe only the primary-bound routes and must not hide a qualified secondary configuration. + `session_lsp` advertises `GET /session/:id/lsp`, the read-only structured LSP status snapshot for daemon clients. Older daemons return `404`; pre-flight this tag before exposing remote LSP status. `session_status` advertises `GET /session/:id/status`, the live bridge summary for a single session by id. In addition to `clientCount` and `hasActivePrompt`, live sessions expose `isWaitingForPermission`, `isWaitingForUserQuestion`, `pendingInteractionCount`, and a retained `turnError` after a failed turn. The error clears when the next prompt actually starts. Both the single-session status response and workspace session lists include `turnError` and `pendingInteractions`: render-ready permission actions or `ask_user_question` questions plus the `requestId` and selectable options required by the existing permission vote routes. Each user question has an `answerKey`; vote with `answers`, for example `{ "0": "Polling" }`, keyed by that value. Persisted-only sessions omit runtime state because no runtime exists. Older daemons return `404`; pre-flight this tag before polling a single session's status instead of scanning the full session list. @@ -412,6 +414,7 @@ operator diagnostic snapshot documented below. | `prompt_absolute_deadline` | `--prompt-deadline-ms` / `QWEN_SERVE_PROMPT_DEADLINE_MS` / `ServeOptions.promptDeadlineMs` is set to a positive integer. | | `writer_idle_timeout` | `--writer-idle-timeout-ms` / `QWEN_SERVE_WRITER_IDLE_TIMEOUT_MS` / `ServeOptions.writerIdleTimeoutMs` is set to a positive integer. | | `workspace_settings` | the daemon was created with settings persistence available. | +| `workspace_qualified_voice` | multi-workspace runtimes and the shared ACP/Voice WebSocket listener are active, so every workspace-qualified Voice modality is reachable for a secondary runtime. | | `session_shell_command` | session shell execution is explicitly enabled. | | `multi_workspace_session_rewind` | more than one workspace runtime is registered; singular live-session rewind routes resolve the owning runtime. | | `multi_workspace_session_shell` | more than one workspace runtime is registered and session shell execution is explicitly enabled; singular REST shell resolves the owning runtime. | @@ -785,7 +788,8 @@ Non-force removal returns `409 workspace_busy` with an `activity` snapshot when "pendingSessionStarts": 0, "acpConnections": 1, "memoryTasks": 0, - "channelWorkers": 0 + "channelWorkers": 0, + "voiceSessions": 0 } } ``` diff --git a/packages/cli/src/serve/acp-http/index.ts b/packages/cli/src/serve/acp-http/index.ts index 9b1654ef81d..e688abd541f 100644 --- a/packages/cli/src/serve/acp-http/index.ts +++ b/packages/cli/src/serve/acp-http/index.ts @@ -24,6 +24,7 @@ import type { WorkspaceRuntime, } from '../workspace-registry.js'; import { + isPortableAbsolutePath, resolveManagedWorkspaceRuntimeFromParam, resolveManagedWorkspaceRuntimeByPathSelector, } from '../workspace-route-runtime.js'; @@ -100,9 +101,10 @@ function isActiveDrainCorrelation( ); } -/** Prefix/suffix of the Phase 4 workspace-qualified ACP WS path. */ -const PLURAL_ACP_WS_PREFIX = '/workspaces/'; +/** Prefix of workspace-qualified WebSocket routes. */ +const PLURAL_WS_PREFIX = '/workspaces/'; const PLURAL_ACP_WS_SUFFIX = '/acp'; +const PLURAL_VOICE_WS_SUFFIX = '/voice/stream'; /** * Extract the raw (undecoded, un-normalized) pathname from a request-target. @@ -121,29 +123,29 @@ function rawRequestPathname(reqUrl: string | undefined): string { } /** - * Match `/workspaces//acp` (with an optional single trailing slash) - * against a RAW request-target pathname and return the still-encoded selector, - * or null when the shape does not match. Rejects empty selectors, extra path - * segments (slash/backslash), and dot-segment traversal shapes -- including - * percent-encoded variants -- so decoding afterwards can never reintroduce a - * `/` or `..` that bypassed classification. + * Match `/workspaces/` (with an optional single trailing + * slash) against a RAW request-target pathname and return the still-encoded + * selector, or null when the shape does not match. Rejects empty selectors, + * extra path segments (slash/backslash), and dot-segment traversal shapes -- + * including percent-encoded variants -- so decoding afterwards can never + * reintroduce a `/` or `..` that bypassed classification. */ -function pluralAcpRawSelector(rawPath: string): string | null { +function pluralWorkspaceRawSelector( + rawPath: string, + suffix: string, +): string | null { let p = rawPath; - if (p.endsWith(`${PLURAL_ACP_WS_SUFFIX}/`)) { + if (p.endsWith(`${suffix}/`)) { p = p.slice(0, -1); } if ( - !p.startsWith(PLURAL_ACP_WS_PREFIX) || - !p.endsWith(PLURAL_ACP_WS_SUFFIX) || - p.length <= PLURAL_ACP_WS_PREFIX.length + PLURAL_ACP_WS_SUFFIX.length + !p.startsWith(PLURAL_WS_PREFIX) || + !p.endsWith(suffix) || + p.length <= PLURAL_WS_PREFIX.length + suffix.length ) { return null; } - const selector = p.slice( - PLURAL_ACP_WS_PREFIX.length, - p.length - PLURAL_ACP_WS_SUFFIX.length, - ); + const selector = p.slice(PLURAL_WS_PREFIX.length, p.length - suffix.length); if ( selector.length === 0 || selector.includes('/') || @@ -436,6 +438,11 @@ export interface MountAcpHttpOptions { * upgrade listener's security checks. Matched paths skip the ACP init flow. */ extraWsRoutes?: readonly ExtraWsRoute[]; + workspaceVoiceConnection?: ( + runtime: WorkspaceRuntime, + ws: WebSocket, + req: IncomingMessage, + ) => void; } /** @@ -1383,7 +1390,8 @@ export function mountAcpHttp( // rather than `url.pathname`. WHATWG URL normalizes dot-segments, so // `/workspaces/%2e%2e/acp` would collapse to `/acp` and silently bind to // the primary mount. `rawRequestPathname` keeps it un-normalized and - // `pluralAcpRawSelector` rejects traversal / backslash / empty selectors. + // `pluralWorkspaceRawSelector` rejects traversal / backslash / empty + // selectors for both ACP and Voice workspace-qualified routes. const rawPath = rawRequestPathname(req.url); const isCdpPath = opts.cdpTunnelOverWs === true && @@ -1393,10 +1401,20 @@ export function mountAcpHttp( (route) => route.path === rawPath, ); const pluralRawSelector = workspaceQualifiedAcpEnabled - ? pluralAcpRawSelector(rawPath) + ? pluralWorkspaceRawSelector(rawPath, PLURAL_ACP_WS_SUFFIX) : null; const isPluralAcpShape = pluralRawSelector !== null; - if (rawPath !== path && !isCdpPath && !extraRoute && !isPluralAcpShape) { + const pluralVoiceRawSelector = opts.workspaceVoiceConnection + ? pluralWorkspaceRawSelector(rawPath, PLURAL_VOICE_WS_SUFFIX) + : null; + const isPluralVoiceShape = pluralVoiceRawSelector !== null; + if ( + rawPath !== path && + !isCdpPath && + !extraRoute && + !isPluralAcpShape && + !isPluralVoiceShape + ) { logReject(`unknown-path ${logSafe(rawPath)}`); socket.destroy(); return; @@ -1504,6 +1522,48 @@ export function mountAcpHttp( return; } + if (isPluralVoiceShape) { + let selector: string; + try { + selector = decodeURIComponent(pluralVoiceRawSelector!); + } catch { + logReject('workspace-selector-decode-error'); + socket.write('HTTP/1.1 400 Bad Request\r\n\r\n'); + socket.destroy(); + return; + } + const wsRegistry = opts.workspaceRegistry; + const runtime = wsRegistry + ? (wsRegistry.getManagedByWorkspaceId(selector) ?? + (isPortableAbsolutePath(selector) + ? resolveManagedWorkspaceRuntimeByPathSelector( + wsRegistry, + selector, + ) + : undefined)) + : undefined; + if (!runtime) { + logReject(`workspace-mismatch ${logSafe(selector)}`); + socket.write('HTTP/1.1 400 Bad Request\r\n\r\n'); + socket.destroy(); + return; + } + if (!runtime.trusted) { + logReject(`untrusted-workspace ${runtime.workspaceId}`); + socket.write('HTTP/1.1 403 Forbidden\r\n\r\n'); + socket.destroy(); + return; + } + wss!.handleUpgrade(req, socket, head, (ws: WebSocket) => { + if (disposed) { + ws.close(1012, 'Server shutting down'); + return; + } + opts.workspaceVoiceConnection!(runtime, ws, req); + }); + return; + } + // ── Phase 4: resolve the target ACP mount for this upgrade ── // Legacy `/acp` binds to the primary mount; `/workspaces/:workspace/acp` // resolves the registered runtime's mount. The shared security checks @@ -1525,7 +1585,12 @@ export function mountAcpHttp( const wsRegistry = opts.workspaceRegistry; const rt = wsRegistry ? (wsRegistry.getManagedByWorkspaceId(selector) ?? - resolveManagedWorkspaceRuntimeByPathSelector(wsRegistry, selector)) + (isPortableAbsolutePath(selector) + ? resolveManagedWorkspaceRuntimeByPathSelector( + wsRegistry, + selector, + ) + : undefined)) : undefined; if (!rt) { logReject(`workspace-mismatch ${logSafe(selector)}`); diff --git a/packages/cli/src/serve/acp-http/workspace-qualified-acp.test.ts b/packages/cli/src/serve/acp-http/workspace-qualified-acp.test.ts index 6d0ea08f214..c71094b8137 100644 --- a/packages/cli/src/serve/acp-http/workspace-qualified-acp.test.ts +++ b/packages/cli/src/serve/acp-http/workspace-qualified-acp.test.ts @@ -137,6 +137,7 @@ describe('workspace-qualified ACP (/workspaces/:workspace/acp)', () => { let secondaryBridge: HttpAcpBridge; let workspaceRegistry: ReturnType; let secondaryRuntime: WorkspaceRuntime; + let workspaceVoiceConnection: ReturnType; beforeEach(async () => { primaryBridge = makeBridge(); @@ -176,6 +177,12 @@ describe('workspace-qualified ACP (/workspaces/:workspace/acp)', () => { }); cdpRegistry = new CdpTunnelRegistry(); checkRate = vi.fn().mockReturnValue(true); + workspaceVoiceConnection = vi.fn( + (runtime: WorkspaceRuntime, ws: WebSocket) => { + ws.send(JSON.stringify({ workspaceCwd: runtime.workspaceCwd })); + ws.close(1000, 'done'); + }, + ); const app = express(); app.use(express.json()); @@ -190,6 +197,7 @@ describe('workspace-qualified ACP (/workspaces/:workspace/acp)', () => { checkRate, sessionShellCommandEnabled: true, workspaceRememberLane: new WorkspaceRememberTaskLane(primaryBridge), + workspaceVoiceConnection, }); await new Promise((resolve) => { @@ -674,6 +682,81 @@ describe('workspace-qualified ACP (/workspaces/:workspace/acp)', () => { ); }); + it('routes workspace-qualified Voice WS by id and encoded cwd', async () => { + const connect = (selector: string) => + new Promise<{ workspaceCwd?: string }>((resolve, reject) => { + const ws = new WebSocket( + `ws://127.0.0.1:${port}/workspaces/${selector}/voice/stream`, + { handshakeTimeout: 2000 }, + ); + ws.on('message', (data: WebSocket.RawData) => { + resolve(JSON.parse(data.toString()) as { workspaceCwd?: string }); + }); + ws.on('error', reject); + }); + + await expect(connect('secondary-id')).resolves.toEqual({ + workspaceCwd: '/ws-b', + }); + await expect(connect(encodeURIComponent('/ws-b'))).resolves.toEqual({ + workspaceCwd: '/ws-b', + }); + expect(workspaceVoiceConnection).toHaveBeenCalledTimes(2); + expect(workspaceVoiceConnection.mock.calls[0]?.[0]).toBe(secondaryRuntime); + expect(workspaceVoiceConnection.mock.calls[1]?.[0]).toBe(secondaryRuntime); + }); + + it('rejects unknown and untrusted workspace-qualified Voice WS upgrades', async () => { + const status = (selector: string) => + new Promise((resolve, reject) => { + const ws = new WebSocket( + `ws://127.0.0.1:${port}/workspaces/${selector}/voice/stream`, + { handshakeTimeout: 2000 }, + ); + ws.on('unexpected-response', (_req, response) => { + resolve(response.statusCode ?? 0); + ws.terminate(); + }); + ws.on('open', () => { + ws.close(); + reject(new Error('rejected Voice WS upgrade should not open')); + }); + ws.on('error', (err) => + reject( + new Error( + `unexpected Voice WS error for ${selector}: ${err.message}`, + ), + ), + ); + }); + + await expect(status('missing')).resolves.toBe(400); + await expect(status('untrusted-id')).resolves.toBe(403); + expect(workspaceVoiceConnection).not.toHaveBeenCalled(); + }); + + it('rejects an encoded relative workspace-qualified Voice selector', async () => { + const relativeSelector = path.relative(process.cwd(), '/ws-b'); + const status = await new Promise((resolve, reject) => { + const ws = new WebSocket( + `ws://127.0.0.1:${port}/workspaces/${encodeURIComponent(relativeSelector)}/voice/stream`, + { handshakeTimeout: 2000 }, + ); + ws.on('unexpected-response', (_req, response) => { + resolve(response.statusCode ?? 0); + ws.terminate(); + }); + ws.on('open', () => { + ws.close(); + reject(new Error('relative Voice WS selector should not open')); + }); + ws.on('error', reject); + }); + + expect(status).toBe(400); + expect(workspaceVoiceConnection).not.toHaveBeenCalled(); + }); + it('rejects a WS upgrade to an untrusted workspace', async () => { const status = await new Promise((resolve, reject) => { const ws = new WebSocket( @@ -688,9 +771,9 @@ describe('workspace-qualified ACP (/workspaces/:workspace/acp)', () => { ws.close(); reject(new Error('untrusted WS upgrade should not open')); }); - // Some ws versions surface a rejected upgrade as an error rather than - // `unexpected-response`; treat that as the expected 403. - ws.on('error', () => resolve(403)); + ws.on('error', (err) => + reject(new Error(`unexpected ACP WS error: ${err.message}`)), + ); }); expect(status).toBe(403); }); diff --git a/packages/cli/src/serve/capabilities.ts b/packages/cli/src/serve/capabilities.ts index c68d60c08e2..a1474ae843f 100644 --- a/packages/cli/src/serve/capabilities.ts +++ b/packages/cli/src/serve/capabilities.ts @@ -285,10 +285,15 @@ export const SERVE_CAPABILITY_REGISTRY = { // workspace agent CRUD, and persisted session organization surfaces. // Workspace-qualified settings also require the existing // `workspace_settings` tag because that surface depends on settings - // persistence. ACP/WebSocket, auth, and voice stay on their existing - // primary-workspace routes in this phase; V2 extension management is - // advertised separately via `extension_management_v2`. + // persistence. ACP/WebSocket and auth stay outside this core tag; + // workspace-qualified Voice REST/WebSocket routes use their separate + // `workspace_qualified_voice` capability below. V2 extension management + // is advertised separately via `extension_management_v2`. workspace_qualified_rest_core: { since: 'v1' }, + // Workspace-qualified Voice REST and WebSocket routes. This tag is enough + // to discover plural modalities because legacy Voice tags describe only + // the primary runtime and may be absent for a secondary-only setup. + workspace_qualified_voice: { since: 'v1' }, // Global extension catalog/mutations plus workspace-qualified activation // projections. This is additive to the legacy primary-workspace // `workspace_extensions` contract. @@ -490,6 +495,14 @@ export const CONDITIONAL_SERVE_FEATURES: ReadonlyMap< toggles.acpHttpEnabled === true && toggles.multiWorkspaceSessionsEnabled === true, ], + [ + 'workspace_qualified_voice', + // Like qualified ACP, the plural Voice surface is mounted ahead of time + // but only becomes useful once the daemon has a secondary runtime. + (toggles) => + toggles.acpHttpEnabled === true && + toggles.multiWorkspaceSessionsEnabled === true, + ], ['client_mcp_over_ws', (toggles) => toggles.clientMcpOverWsEnabled === true], ['cdp_tunnel_over_ws', (toggles) => toggles.cdpTunnelOverWsEnabled === true], [ diff --git a/packages/cli/src/serve/routes/workspace-management.test.ts b/packages/cli/src/serve/routes/workspace-management.test.ts index f6324847e79..e808fe29a67 100644 --- a/packages/cli/src/serve/routes/workspace-management.test.ts +++ b/packages/cli/src/serve/routes/workspace-management.test.ts @@ -127,7 +127,11 @@ function createRemovalController( beginDrain: vi.fn(), cancelDrain: vi.fn(), completeDrain: vi.fn(), - getActivity: vi.fn(() => ({ pendingSessionStarts, channelWorkers: 0 })), + getActivity: vi.fn(() => ({ + pendingSessionStarts, + channelWorkers: 0, + voiceSessions: 0, + })), disposeRuntime: vi.fn().mockResolvedValue(undefined), }; } @@ -600,12 +604,46 @@ describe('DELETE /workspaces/:workspace', () => { expect(deps.workspaceRegistry.beginDrain).not.toHaveBeenCalled(); }); + it('blocks non-force removal while a Voice operation is active', async () => { + const runtime = makeRuntime(REAL_DIR); + const runtimeRemoval = createRemovalController(); + vi.mocked(runtimeRemoval.getActivity).mockReturnValue({ + pendingSessionStarts: 0, + channelWorkers: 0, + voiceSessions: 1, + }); + const { app, deps } = createApp({ + workspaceRegistry: createMockRegistry([runtime]), + runtimeRemoval, + }); + + const res = await request(app).delete( + `/workspaces/${encodeURIComponent(runtime.workspaceId)}`, + ); + + expect(res.status).toBe(409); + expect(res.body).toMatchObject({ + code: 'workspace_busy', + activity: { voiceSessions: 1 }, + }); + expect(runtimeRemoval.beginDrain).not.toHaveBeenCalled(); + expect(deps.workspaceRegistry.beginDrain).not.toHaveBeenCalled(); + }); + it('rolls every gate back when the final frozen snapshot becomes busy', async () => { const runtime = makeRuntime(REAL_DIR); const runtimeRemoval = createRemovalController(); vi.mocked(runtimeRemoval.getActivity) - .mockReturnValueOnce({ pendingSessionStarts: 0, channelWorkers: 0 }) - .mockReturnValueOnce({ pendingSessionStarts: 1, channelWorkers: 0 }); + .mockReturnValueOnce({ + pendingSessionStarts: 0, + channelWorkers: 0, + voiceSessions: 0, + }) + .mockReturnValueOnce({ + pendingSessionStarts: 1, + channelWorkers: 0, + voiceSessions: 0, + }); const acpHandle = { beginWorkspaceDrain: vi.fn(), cancelWorkspaceDrain: vi.fn(), @@ -643,8 +681,16 @@ describe('DELETE /workspaces/:workspace', () => { const runtime = makeRuntime(REAL_DIR); const runtimeRemoval = createRemovalController(); vi.mocked(runtimeRemoval.getActivity) - .mockReturnValueOnce({ pendingSessionStarts: 0, channelWorkers: 0 }) - .mockReturnValueOnce({ pendingSessionStarts: 1, channelWorkers: 0 }); + .mockReturnValueOnce({ + pendingSessionStarts: 0, + channelWorkers: 0, + voiceSessions: 0, + }) + .mockReturnValueOnce({ + pendingSessionStarts: 1, + channelWorkers: 0, + voiceSessions: 0, + }); vi.mocked(runtimeRemoval.cancelDrain).mockImplementation(() => { throw new Error('controller rollback failed'); }); diff --git a/packages/cli/src/serve/routes/workspace-management.ts b/packages/cli/src/serve/routes/workspace-management.ts index 813d5516585..a54114e25c9 100644 --- a/packages/cli/src/serve/routes/workspace-management.ts +++ b/packages/cli/src/serve/routes/workspace-management.ts @@ -50,6 +50,7 @@ export interface WorkspaceRemovalActivity { acpConnections: number; memoryTasks: number; channelWorkers: number; + voiceSessions: number; } export interface WorkspaceRuntimeRemovalController { @@ -60,6 +61,7 @@ export interface WorkspaceRuntimeRemovalController { getActivity(runtime: WorkspaceRuntime): { pendingSessionStarts: number; channelWorkers: number; + voiceSessions: number; }; disposeRuntime( runtime: WorkspaceRuntime, @@ -492,6 +494,7 @@ export function registerWorkspaceManagementRoutes( const controllerActivity = runtimeRemoval?.getActivity(runtime) ?? { pendingSessionStarts: 0, channelWorkers: 0, + voiceSessions: 0, }; const acpActivity = getAcpHandle?.()?.getWorkspaceActivity( runtime.workspaceId, @@ -503,6 +506,7 @@ export function registerWorkspaceManagementRoutes( acpConnections: acpActivity.acpConnections, memoryTasks: acpActivity.memoryTasks, channelWorkers: controllerActivity.channelWorkers, + voiceSessions: controllerActivity.voiceSessions, }; }; const isBusy = (activity: WorkspaceRemovalActivity): boolean => diff --git a/packages/cli/src/serve/routes/workspace-qualified-voice.test.ts b/packages/cli/src/serve/routes/workspace-qualified-voice.test.ts new file mode 100644 index 00000000000..f2e08e2b2bd --- /dev/null +++ b/packages/cli/src/serve/routes/workspace-qualified-voice.test.ts @@ -0,0 +1,302 @@ +/** + * @license + * Copyright 2026 Qwen Team + * SPDX-License-Identifier: Apache-2.0 + */ + +import { promises as fsp } from 'node:fs'; +import * as os from 'node:os'; +import * as path from 'node:path'; +import express from 'express'; +import { afterEach, describe, expect, it, vi } from 'vitest'; +import request from 'supertest'; +import { + SettingScope, + resetHomeEnvBootstrapForTesting, +} from '../../config/settings.js'; +import { registerWorkspaceQualifiedVoiceRoutes } from './workspace-voice.js'; +import { + createWorkspaceRegistry, + type WorkspaceRuntime, +} from '../workspace-registry.js'; +import { + WorkspaceVoiceCoordinator, + type VoiceAdmissionResult, +} from '../voice/workspace-voice-coordinator.js'; + +const homes: string[] = []; + +function runtime( + workspaceId: string, + workspaceCwd: string, + opts: { + primary?: boolean; + trusted?: boolean; + envMode?: 'parent-process' | 'runtime-overlay'; + effectiveEnv?: Readonly>; + } = {}, +): WorkspaceRuntime { + return { + workspaceId, + workspaceCwd, + primary: opts.primary === true, + trusted: opts.trusted !== false, + env: + opts.envMode === 'parent-process' + ? { mode: 'parent-process', overlayKeys: [] } + : { + mode: 'runtime-overlay', + overlayKeys: [], + effectiveEnv: opts.effectiveEnv ?? {}, + }, + bridge: { publishWorkspaceEvent: vi.fn() }, + } as unknown as WorkspaceRuntime; +} + +async function createApp( + opts: { + acquireVoiceLease?: (runtime: WorkspaceRuntime) => VoiceAdmissionResult; + transcribe?: ReturnType; + secondaryEnvMode?: 'parent-process' | 'runtime-overlay'; + } = {}, +): Promise<{ + app: express.Application; + secondary: WorkspaceRuntime; + untrusted: WorkspaceRuntime; + registry: ReturnType; + persistSetting: ReturnType; + acquireVoiceLease: ReturnType; + transcribe: ReturnType; + invalidateServeFeaturesCache: ReturnType; +}> { + const home = await fsp.mkdtemp(path.join(os.tmpdir(), 'qwen-voice-home-')); + homes.push(home); + const primaryCwd = path.join(home, 'primary'); + const secondaryCwd = path.join(home, 'secondary'); + const untrustedCwd = path.join(home, 'untrusted'); + await Promise.all( + [primaryCwd, secondaryCwd, untrustedCwd].map((cwd) => + fsp.mkdir(cwd, { recursive: true }), + ), + ); + process.env['QWEN_HOME'] = home; + resetHomeEnvBootstrapForTesting(); + + const primary = runtime('primary-id', primaryCwd, { primary: true }); + const secondary = runtime('secondary-id', secondaryCwd, { + envMode: opts.secondaryEnvMode, + effectiveEnv: { SECONDARY_VOICE_KEY: 'secondary-secret' }, + }); + const untrusted = runtime('untrusted-id', untrustedCwd, { trusted: false }); + const registry = createWorkspaceRegistry([primary, secondary, untrusted]); + const persistSetting = vi.fn(async () => undefined); + const acquireVoiceLease = vi.fn( + opts.acquireVoiceLease ?? + (() => ({ + kind: 'admitted' as const, + lease: { signal: new AbortController().signal, release: () => {} }, + })), + ); + const transcribe = + opts.transcribe ?? + vi.fn(async () => ({ + text: 'secondary transcript', + model: 'secondary-asr', + transport: 'qwen-asr-chat' as const, + })); + const invalidateServeFeaturesCache = vi.fn(); + const app = express(); + app.use(express.json()); + registerWorkspaceQualifiedVoiceRoutes(app, { + workspaceRegistry: registry, + mutate: () => (_req, _res, next) => next(), + safeBody: (req) => req.body as Record, + persistSetting, + acquireVoiceLease: (target) => acquireVoiceLease(target), + transcribe, + parseAndValidateClientId: () => undefined, + invalidateServeFeaturesCache, + }); + return { + app, + secondary, + untrusted, + registry, + persistSetting, + acquireVoiceLease, + transcribe, + invalidateServeFeaturesCache, + }; +} + +async function enableSecondaryVoice(runtime: WorkspaceRuntime): Promise { + await fsp.mkdir(path.join(runtime.workspaceCwd, '.qwen'), { + recursive: true, + }); + await fsp.writeFile( + path.join(runtime.workspaceCwd, '.qwen', 'settings.json'), + JSON.stringify({ + voiceModel: 'secondary-asr', + general: { voice: { enabled: true } }, + }), + 'utf8', + ); +} + +describe('workspace-qualified Voice routes', () => { + const originalQwenHome = process.env['QWEN_HOME']; + + afterEach(async () => { + await Promise.all( + homes + .splice(0) + .map((home) => fsp.rm(home, { recursive: true, force: true })), + ); + if (originalQwenHome === undefined) delete process.env['QWEN_HOME']; + else process.env['QWEN_HOME'] = originalQwenHome; + resetHomeEnvBootstrapForTesting(); + }); + + it('selects only the trusted target runtime and writes Voice settings in workspace scope', async () => { + const { app, secondary, persistSetting, invalidateServeFeaturesCache } = + await createApp(); + + await expect( + request(app).get('/workspaces/secondary-id/voice'), + ).resolves.toMatchObject({ + status: 200, + body: { workspaceCwd: secondary.workspaceCwd }, + }); + await expect( + request(app) + .post('/workspaces/secondary-id/voice') + .send({ enabled: false }), + ).resolves.toMatchObject({ status: 200 }); + await expect( + request(app).get( + `/workspaces/${encodeURIComponent(secondary.workspaceCwd)}/voice`, + ), + ).resolves.toMatchObject({ + status: 200, + body: { workspaceCwd: secondary.workspaceCwd }, + }); + + expect(persistSetting).toHaveBeenCalledWith( + secondary.workspaceCwd, + SettingScope.Workspace, + 'general.voice.enabled', + false, + ); + expect( + secondary.bridge.publishWorkspaceEvent as ReturnType, + ).toHaveBeenCalled(); + expect(invalidateServeFeaturesCache).not.toHaveBeenCalled(); + }); + + it('invalidates the primary-derived feature cache for primary Voice updates', async () => { + const { app, invalidateServeFeaturesCache } = await createApp(); + + await expect( + request(app) + .post('/workspaces/primary-id/voice') + .send({ enabled: false }), + ).resolves.toMatchObject({ status: 200 }); + + expect(invalidateServeFeaturesCache).toHaveBeenCalledOnce(); + }); + + it('rejects an unknown selector before reading settings and untrusted targets without fallback', async () => { + const { app, acquireVoiceLease } = await createApp(); + + await expect( + request(app).get('/workspaces/missing/voice'), + ).resolves.toMatchObject({ + status: 400, + body: { code: 'workspace_mismatch' }, + }); + await expect( + request(app).get('/workspaces/untrusted-id/voice'), + ).resolves.toMatchObject({ + status: 403, + body: { code: 'untrusted_workspace' }, + }); + expect(acquireVoiceLease).not.toHaveBeenCalled(); + }); + + it('transcribes with the selected runtime cwd and effective environment', async () => { + const { app, secondary, transcribe } = await createApp(); + await enableSecondaryVoice(secondary); + + const response = await request(app) + .post('/workspaces/secondary-id/voice/transcribe') + .set('Content-Type', 'audio/wav') + .send(Buffer.from([1, 2, 3])); + + expect(response.status).toBe(200); + expect(transcribe).toHaveBeenCalledWith( + expect.objectContaining({ + workspaceCwd: secondary.workspaceCwd, + env: secondary.env.effectiveEnv, + voiceModel: 'secondary-asr', + abortSignal: expect.any(AbortSignal), + }), + ); + }); + + it('inherits the parent process environment for parent-process runtimes', async () => { + const { app, secondary, transcribe } = await createApp({ + secondaryEnvMode: 'parent-process', + }); + await enableSecondaryVoice(secondary); + + const response = await request(app) + .post('/workspaces/secondary-id/voice/transcribe') + .set('Content-Type', 'audio/wav') + .send(Buffer.from([1, 2, 3])); + + expect(response.status).toBe(200); + expect(transcribe).toHaveBeenCalledWith( + expect.objectContaining({ env: undefined }), + ); + }); + + it('rejects capacity before reading the audio body', async () => { + const { app, secondary, transcribe, acquireVoiceLease } = await createApp({ + acquireVoiceLease: () => ({ kind: 'rejected', reason: 'capacity' }), + }); + await enableSecondaryVoice(secondary); + + const response = await request(app) + .post('/workspaces/secondary-id/voice/transcribe') + .set('Content-Type', 'audio/wav') + .send(Buffer.from([1, 2, 3])); + + expect(response.status).toBe(503); + expect(response.headers['retry-after']).toBe('5'); + expect(response.body.code).toBe('voice_capacity_exceeded'); + expect(acquireVoiceLease).toHaveBeenCalledOnce(); + expect(transcribe).not.toHaveBeenCalled(); + }); + + it('reports workspace draining after the registry drain gate closes', async () => { + const coordinator = new WorkspaceVoiceCoordinator(); + const { app, secondary, registry, transcribe, acquireVoiceLease } = + await createApp({ + acquireVoiceLease: (runtime) => coordinator.acquire(runtime), + }); + await enableSecondaryVoice(secondary); + expect(registry.beginDrain(secondary)).toBe(true); + coordinator.beginWorkspaceDrain(secondary); + + const response = await request(app) + .post('/workspaces/secondary-id/voice/transcribe') + .set('Content-Type', 'audio/wav') + .send(Buffer.from([1, 2, 3])); + + expect(response.status).toBe(503); + expect(response.headers['retry-after']).toBe('5'); + expect(response.body.code).toBe('workspace_draining'); + expect(acquireVoiceLease).toHaveBeenCalledOnce(); + expect(transcribe).not.toHaveBeenCalled(); + }); +}); diff --git a/packages/cli/src/serve/routes/workspace-voice.test.ts b/packages/cli/src/serve/routes/workspace-voice.test.ts index 170bce8dce4..3fb80c9d98c 100644 --- a/packages/cli/src/serve/routes/workspace-voice.test.ts +++ b/packages/cli/src/serve/routes/workspace-voice.test.ts @@ -6,6 +6,8 @@ import { randomBytes } from 'node:crypto'; import { promises as fsp } from 'node:fs'; +import { request as httpRequest, type Server } from 'node:http'; +import type { AddressInfo } from 'node:net'; import * as os from 'node:os'; import * as path from 'node:path'; import express, { @@ -226,6 +228,21 @@ describe('workspace voice routes', () => { expect(serialized).not.toContain('envKey'); }); + it('GET returns a structured error when Voice settings are unreadable', async () => { + await fsp.mkdir(path.join(h.home, 'settings.json')); + + const res = await request(h.app) + .get('/workspace/voice') + .set('Host', hostHeader) + .set('Authorization', 'Bearer secret'); + + expect(res.status).toBe(500); + expect(res.body).toEqual({ + error: 'Failed to load voice settings', + code: 'internal_error', + }); + }); + it('POST updates voice settings only after validating the resulting state', async () => { await writeVoiceModelSettings(h); @@ -504,7 +521,7 @@ describe('workspace voice routes', () => { ); }); - it('POST does not broadcast fallback voice writes when a later write fails', async () => { + it('POST broadcasts committed fallback voice writes when a later write fails', async () => { await writeVoiceModelSettings(h); const broadcastSettingsChanged = vi.fn(); const persistSetting = vi.fn( @@ -538,7 +555,13 @@ describe('workspace voice routes', () => { }); expect(res.status).toBe(500); - expect(broadcastSettingsChanged).not.toHaveBeenCalled(); + expect(broadcastSettingsChanged).toHaveBeenCalledOnce(); + expect(broadcastSettingsChanged).toHaveBeenCalledWith( + 'voiceModel', + 'qwen3-asr-flash', + 'user', + 'client-1', + ); const error = mockWriteStderrLine.mock.calls[0]?.[0]; expect(mockWriteStderrLine).toHaveBeenCalledWith( expect.stringContaining('partial persist error'), @@ -546,6 +569,36 @@ describe('workspace voice routes', () => { expect(error).toContain('committed=1/2'); }); + it('returns a structured error when persisted settings cannot be reloaded', async () => { + const broadcastSettingsChanged = vi.fn(); + const persistSetting = vi.fn(async () => { + await fsp.mkdir(path.join(h.home, 'settings.json')); + }); + const app = express(); + app.use(express.json({ limit: '10mb' })); + registerWorkspaceVoiceRoutes(app, { + boundWorkspace: h.workspace, + mutate: () => (_req: Request, _res: Response, next: NextFunction) => + next(), + safeBody: (req) => req.body as Record, + persistSetting, + broadcastSettingsChanged, + parseAndValidateClientId: vi.fn(() => 'client-1'), + transcribe: h.transcribe, + }); + + const res = await request(app) + .post('/workspace/voice') + .send({ enabled: false }); + + expect(res.status).toBe(500); + expect(res.body).toEqual({ + error: 'Voice settings persisted but failed to reload', + code: 'internal_error', + }); + expect(broadcastSettingsChanged).toHaveBeenCalledOnce(); + }); + it('POST rejects enabling voice when no valid voice model is selected', async () => { const res = await request(h.app) .post('/workspace/voice') @@ -935,6 +988,210 @@ describe('workspace voice routes', () => { expect(stderrOutput).not.toContain('top-secret'); }); + it('POST /workspace/voice/transcribe reports draining when an active lease is aborted', async () => { + await writeVoiceModelSettings(h); + const leaseController = new AbortController(); + const release = vi.fn(); + const transcribe = vi.fn( + async (input: { abortSignal?: AbortSignal }): Promise => + await new Promise((_resolve, reject) => { + input.abortSignal?.addEventListener( + 'abort', + () => reject(new Error('aborted')), + { once: true }, + ); + }), + ); + const app = express(); + registerWorkspaceVoiceRoutes(app, { + boundWorkspace: h.workspace, + mutate: () => (_req: Request, _res: Response, next: NextFunction) => + next(), + safeBody: (req) => req.body as Record, + persistSetting: h.persistSetting, + broadcastSettingsChanged: vi.fn(), + parseAndValidateClientId: () => undefined, + acquireVoiceLease: () => ({ + kind: 'admitted', + lease: { signal: leaseController.signal, release }, + }), + transcribe, + }); + + const response = request(app) + .post('/workspace/voice/transcribe') + .set('Content-Type', 'audio/wav') + .send(Buffer.from([1])) + .then((result) => result); + await vi.waitFor(() => expect(transcribe).toHaveBeenCalledOnce()); + leaseController.abort(); + + await expect(response).resolves.toMatchObject({ + status: 503, + body: { code: 'workspace_draining' }, + }); + expect(release).toHaveBeenCalledOnce(); + }); + + it('releases an admitted lease once after successful transcription', async () => { + await writeVoiceModelSettings(h); + const release = vi.fn(); + const app = express(); + registerWorkspaceVoiceRoutes(app, { + boundWorkspace: h.workspace, + mutate: () => (_req: Request, _res: Response, next: NextFunction) => + next(), + safeBody: (req) => req.body as Record, + persistSetting: h.persistSetting, + broadcastSettingsChanged: vi.fn(), + parseAndValidateClientId: () => undefined, + acquireVoiceLease: () => ({ + kind: 'admitted', + lease: { signal: new AbortController().signal, release }, + }), + transcribe: h.transcribe, + }); + + await expect( + request(app) + .post('/workspace/voice/transcribe') + .set('Content-Type', 'audio/wav') + .send(Buffer.from([1])), + ).resolves.toMatchObject({ status: 200 }); + + expect(release).toHaveBeenCalledOnce(); + }); + + it('aborts a slow upload and releases its lease before transcription starts', async () => { + const leaseController = new AbortController(); + const release = vi.fn(); + const transcribe = vi.fn(); + const acquireVoiceLease = vi.fn(() => ({ + kind: 'admitted' as const, + lease: { signal: leaseController.signal, release }, + })); + const app = express(); + registerWorkspaceVoiceRoutes(app, { + boundWorkspace: h.workspace, + mutate: () => (_req: Request, _res: Response, next: NextFunction) => + next(), + safeBody: (req) => req.body as Record, + persistSetting: h.persistSetting, + broadcastSettingsChanged: vi.fn(), + parseAndValidateClientId: () => undefined, + acquireVoiceLease, + transcribe, + }); + const server = await new Promise((resolve) => { + const listening = app.listen(0, '127.0.0.1', () => resolve(listening)); + }); + try { + const port = (server.address() as AddressInfo).port; + const response = new Promise<{ status: number; body: string }>( + (resolve, reject) => { + const req = httpRequest( + { + host: '127.0.0.1', + port, + path: '/workspace/voice/transcribe', + method: 'POST', + headers: { + 'Content-Type': 'audio/wav', + 'Content-Length': '100', + }, + }, + (res) => { + let body = ''; + res.setEncoding('utf8'); + res.on('data', (chunk) => { + body += chunk; + }); + res.on('end', () => + resolve({ status: res.statusCode ?? 0, body }), + ); + }, + ); + req.on('error', reject); + req.write(Buffer.from([1])); + }, + ); + + await vi.waitFor(() => expect(acquireVoiceLease).toHaveBeenCalledOnce()); + leaseController.abort(new Error('Workspace runtime was removed')); + + await expect(response).resolves.toMatchObject({ + status: 503, + body: expect.stringContaining('workspace_draining'), + }); + expect(release).toHaveBeenCalledOnce(); + expect(transcribe).not.toHaveBeenCalled(); + } finally { + await new Promise((resolve) => server.close(() => resolve())); + } + }); + + it('retains the lease after disconnect until transcription cleanup settles', async () => { + await writeVoiceModelSettings(h); + const release = vi.fn(); + let operationSignal: AbortSignal | undefined; + let settle: + | ((value: Awaited>) => void) + | undefined; + const transcribe = vi.fn( + async (input: { abortSignal?: AbortSignal }) => + await new Promise>>( + (resolve) => { + operationSignal = input.abortSignal; + settle = resolve; + }, + ), + ); + const app = express(); + registerWorkspaceVoiceRoutes(app, { + boundWorkspace: h.workspace, + mutate: () => (_req: Request, _res: Response, next: NextFunction) => + next(), + safeBody: (req) => req.body as Record, + persistSetting: h.persistSetting, + broadcastSettingsChanged: vi.fn(), + parseAndValidateClientId: () => undefined, + acquireVoiceLease: () => ({ + kind: 'admitted', + lease: { signal: new AbortController().signal, release }, + }), + transcribe, + }); + const server = await new Promise((resolve) => { + const listening = app.listen(0, '127.0.0.1', () => resolve(listening)); + }); + try { + const port = (server.address() as AddressInfo).port; + const req = httpRequest({ + host: '127.0.0.1', + port, + path: '/workspace/voice/transcribe', + method: 'POST', + headers: { 'Content-Type': 'audio/wav', 'Content-Length': '1' }, + }); + req.on('error', () => {}); + req.end(Buffer.from([1])); + await vi.waitFor(() => expect(transcribe).toHaveBeenCalledOnce()); + + req.destroy(); + await vi.waitFor(() => expect(operationSignal?.aborted).toBe(true)); + expect(release).not.toHaveBeenCalled(); + + settle?.({ + text: 'discarded', + model: 'qwen3-asr-flash', + transport: 'qwen-asr-chat', + }); + await vi.waitFor(() => expect(release).toHaveBeenCalledOnce()); + } finally { + await new Promise((resolve) => server.close(() => resolve())); + } + }); + it('POST /workspace/voice/transcribe rejects unsupported content types', async () => { const res = await request(h.app) .post('/workspace/voice/transcribe') diff --git a/packages/cli/src/serve/routes/workspace-voice.ts b/packages/cli/src/serve/routes/workspace-voice.ts index cea16cb8b72..2112beff55b 100644 --- a/packages/cli/src/serve/routes/workspace-voice.ts +++ b/packages/cli/src/serve/routes/workspace-voice.ts @@ -11,7 +11,7 @@ import express, { } from 'express'; import { loadSettings, - type SettingScope, + SettingScope, type LoadedSettings, } from '../../config/settings.js'; import { getWorkspaceTrustStatus } from '../../config/trustedFolders.js'; @@ -43,6 +43,19 @@ import { WorkspaceSettingsPartialPersistError, type WorkspaceSettingsWrite, } from '../workspace-service/types.js'; +import type { + VoiceAdmissionLease, + VoiceAdmissionResult, +} from '../voice/workspace-voice-coordinator.js'; +import type { + WorkspaceRegistry, + WorkspaceRuntime, +} from '../workspace-registry.js'; +import { + requireTrustedWorkspaceRuntime, + resolveManagedWorkspaceRuntimeFromParam, + resolveWorkspaceRuntimeFromParam, +} from '../workspace-route-runtime.js'; type WorkspaceVoiceTranscriber = ( input: WorkspaceVoiceTranscriptionInput, @@ -76,7 +89,26 @@ export interface WorkspaceVoiceRouteDeps { req: Request, res: Response, ) => string | undefined | null; + env?: Readonly>; + scopeOverride?: SettingScope; + acquireVoiceLease?: () => VoiceAdmissionResult; + transcribe?: WorkspaceVoiceTranscriber; +} + +export interface WorkspaceQualifiedVoiceRouteDeps { + workspaceRegistry: WorkspaceRegistry; + mutate: WorkspaceVoiceRouteDeps['mutate']; + safeBody: WorkspaceVoiceRouteDeps['safeBody']; + persistSetting?: PersistSetting; + persistSettings?: PersistSettings; transcribe?: WorkspaceVoiceTranscriber; + acquireVoiceLease: (runtime: WorkspaceRuntime) => VoiceAdmissionResult; + parseAndValidateClientId: ( + req: Request, + res: Response, + runtime: WorkspaceRuntime, + ) => string | undefined | null; + invalidateServeFeaturesCache: () => void; } function sendVoiceError(res: Response, err: unknown): boolean { @@ -186,6 +218,7 @@ async function persistVoiceUpdate( } const writes = buildWorkspaceVoiceSettingsWrites(settings, update, { workspaceTrusted, + ...(deps.scopeOverride ? { scopeOverride: deps.scopeOverride } : {}), }); if (deps.persistSettings) { @@ -218,6 +251,9 @@ async function persistVoiceUpdate( err instanceof Error ? err.message : String(err) }`, ); + for (const committedWrite of committed) { + broadcastVoiceWrite(deps, committedWrite, clientId); + } throw new WorkspaceSettingsPartialPersistError( `Voice settings partial persist failed: committed=${committed.length}/${writes.length}`, committed, @@ -246,6 +282,89 @@ function requestAbortSignal(req: Request, res: Response): AbortSignal { return controller.signal; } +function loadVoiceSettings(deps: WorkspaceVoiceRouteDeps): LoadedSettings { + return loadSettings( + deps.boundWorkspace, + deps.env ? { skipLoadEnvironment: true } : true, + ); +} + +interface VoiceAdmissionState { + lease: VoiceAdmissionLease; + operationStarted: boolean; + release: () => void; +} + +function admissionState(req: Request): VoiceAdmissionState | undefined { + return (req as { voiceAdmissionState?: VoiceAdmissionState }) + .voiceAdmissionState; +} + +function admissionLease(req: Request): VoiceAdmissionLease | undefined { + return admissionState(req)?.lease; +} + +function installAdmissionLease( + req: Request, + res: Response, + deps: WorkspaceVoiceRouteDeps, +): boolean { + const result = deps.acquireVoiceLease?.(); + if (!result) return true; + if (result.kind === 'rejected') { + if (result.reason === 'draining') { + sendWorkspaceDraining(res); + } else { + res.set('Retry-After', '5').status(503).json({ + error: 'Too many voice sessions in progress; try again shortly.', + code: 'voice_capacity_exceeded', + }); + } + return false; + } + let released = false; + const release = () => { + if (released) return; + released = true; + result.lease.release(); + }; + const state: VoiceAdmissionState = { + lease: result.lease, + operationStarted: false, + release, + }; + (req as { voiceAdmissionState?: VoiceAdmissionState }).voiceAdmissionState = + state; + const releaseBeforeOperation = () => { + if (!state.operationStarted) release(); + }; + res.once('finish', releaseBeforeOperation); + res.once('close', releaseBeforeOperation); + return true; +} + +function beginVoiceOperation(req: Request): (() => void) | undefined { + const state = admissionState(req); + if (!state) return; + state.operationStarted = true; + return state.release; +} + +function combinedAbortSignal(req: Request, res: Response): AbortSignal { + const requestSignal = requestAbortSignal(req, res); + const leaseSignal = admissionLease(req)?.signal; + return leaseSignal + ? AbortSignal.any([requestSignal, leaseSignal]) + : requestSignal; +} + +function sendWorkspaceDraining(res: Response): void { + res.set('Retry-After', '5').status(503).json({ + error: 'Workspace runtime is being removed', + code: 'workspace_draining', + }); +} + function isSupportedAudioContentType( contentType: string | undefined, ): contentType is string { @@ -267,199 +386,389 @@ function readBinaryBody(req: Request): Uint8Array | undefined { return undefined; } -export function registerWorkspaceVoiceRoutes( - app: Application, +function sendUnsupportedVoiceContentType(res: Response): void { + res.status(415).json({ + error: + 'Content-Type must be audio/* or application/octet-stream for voice transcription', + code: 'unsupported_voice_content_type', + }); +} + +function handleVoiceStatus( + res: Response, deps: WorkspaceVoiceRouteDeps, + route: string, ): void { - const transcribe = deps.transcribe ?? transcribeWorkspaceVoiceAudio; + try { + res + .status(200) + .json( + buildWorkspaceVoiceStatus(deps.boundWorkspace, loadVoiceSettings(deps)), + ); + } catch (err) { + writeStderrLine( + `qwen serve: ${route} error: ${ + err instanceof Error ? err.message : String(err) + }`, + ); + res.status(500).json({ + error: 'Failed to load voice settings', + code: 'internal_error', + }); + } +} - app.get('/workspace/voice', (_req: Request, res: Response) => { - try { - res - .status(200) - .json( - buildWorkspaceVoiceStatus( - deps.boundWorkspace, - loadSettings(deps.boundWorkspace), - ), - ); - } catch (err) { - writeStderrLine( - `qwen serve: GET /workspace/voice error: ${ - err instanceof Error ? err.message : String(err) - }`, +async function handleVoiceUpdate( + req: Request, + res: Response, + deps: WorkspaceVoiceRouteDeps, + route: string, +): Promise { + if (!deps.persistSettings && !deps.persistSetting) { + res.status(501).json({ + error: 'Workspace voice settings persistence is not available', + code: 'not_implemented', + }); + return; + } + const parsed = parseWorkspaceVoiceUpdateParams(deps.safeBody(req)); + if ('error' in parsed) { + res.status(400).json(parsed); + return; + } + const settings = loadVoiceSettings(deps); + try { + validateWorkspaceVoiceState(settings, parsed, { env: deps.env }); + } catch (err) { + if (sendVoiceError(res, err)) return; + throw err; + } + const clientId = deps.parseAndValidateClientId(req, res); + if (clientId === null) return; + try { + const workspaceTrusted = + getWorkspaceTrustStatus(settings.merged, deps.boundWorkspace).effective + .state === 'trusted'; + await persistVoiceUpdate( + deps, + settings, + parsed, + clientId, + workspaceTrusted, + ); + } catch (err) { + writeStderrLine( + `qwen serve: ${route} persist error (workspace=${deps.boundWorkspace}): ${ + err instanceof Error ? err.message : String(err) + }`, + ); + res.status(500).json({ + error: 'Failed to persist voice settings', + code: 'persist_error', + }); + return; + } + try { + res + .status(200) + .json( + buildWorkspaceVoiceStatus(deps.boundWorkspace, loadVoiceSettings(deps)), ); - res.status(500).json({ - error: 'Failed to load voice settings', - code: 'internal_error', + } catch (err) { + writeStderrLine( + `qwen serve: ${route} reload error after persist (workspace=${deps.boundWorkspace}): ${ + err instanceof Error ? err.message : String(err) + }`, + ); + res.status(500).json({ + error: 'Voice settings persisted but failed to reload', + code: 'internal_error', + }); + } +} + +async function handleVoiceTranscription( + req: Request, + res: Response, + deps: WorkspaceVoiceRouteDeps, + route: string, +): Promise { + const contentType = normalizeContentType(req); + if (!isSupportedAudioContentType(contentType)) { + sendUnsupportedVoiceContentType(res); + return; + } + if (admissionLease(req)?.signal.aborted) { + sendWorkspaceDraining(res); + return; + } + const clientId = deps.parseAndValidateClientId(req, res); + if (clientId === null) return; + const data = readBinaryBody(req); + if (!data || data.byteLength === 0) { + res.status(400).json({ + error: 'Voice audio body must be non-empty binary data', + code: 'invalid_voice_audio', + }); + return; + } + const settings = loadVoiceSettings(deps); + if (!isVoiceEnabled(settings)) { + res.status(403).json({ + error: 'Voice transcription is disabled for this workspace', + code: 'voice_disabled', + }); + return; + } + const queryVoiceModel = req.query['voiceModel']; + if (queryVoiceModel !== undefined && typeof queryVoiceModel !== 'string') { + res.status(400).json({ + error: '`voiceModel` query parameter must be a string', + code: 'invalid_voice_model', + }); + return; + } + const requestedVoiceModel = + typeof queryVoiceModel === 'string' ? queryVoiceModel.trim() : ''; + if ( + requestedVoiceModel && + requestedVoiceModel.length > MAX_VOICE_MODEL_LENGTH + ) { + res.status(400).json({ + error: `\`voiceModel\` exceeds the ${MAX_VOICE_MODEL_LENGTH}-character limit`, + code: 'invalid_voice_model', + }); + return; + } + const voiceModel = requestedVoiceModel || readVoiceModel(settings); + if (!voiceModel) { + res.status(400).json({ + error: 'A valid voiceModel is required before transcription.', + code: 'voice_model_required', + }); + return; + } + try { + const releaseOperation = beginVoiceOperation(req); + let result: WorkspaceVoiceTranscriptionResult; + try { + result = await (deps.transcribe ?? transcribeWorkspaceVoiceAudio)({ + data, + mimeType: contentType, + voiceModel, + settings, + workspaceCwd: deps.boundWorkspace, + env: deps.env, + abortSignal: combinedAbortSignal(req, res), }); + } finally { + releaseOperation?.(); } - }); - - app.post( - '/workspace/voice', - deps.mutate({ strict: true }), - async (req: Request, res: Response) => { - if (!deps.persistSettings && !deps.persistSetting) { - res.status(501).json({ - error: 'Workspace voice settings persistence is not available', - code: 'not_implemented', - }); - return; - } + if (admissionLease(req)?.signal.aborted) { + sendWorkspaceDraining(res); + return; + } + res.status(200).json({ v: 1, ...result }); + } catch (err) { + if (admissionLease(req)?.signal.aborted) { + sendWorkspaceDraining(res); + return; + } + if (sendVoiceError(res, err)) return; + const message = sanitizeVoiceErrorMessage( + err instanceof Error ? err.message : String(err), + ); + writeStderrLine( + `qwen serve: ${route} error (workspace=${deps.boundWorkspace}): ${message}`, + ); + res.status(502).json({ + error: 'Voice transcription failed', + code: 'voice_transcription_failed', + }); + } +} - const parsed = parseWorkspaceVoiceUpdateParams(deps.safeBody(req)); - if ('error' in parsed) { - res.status(400).json(parsed); - return; - } +function admitVoiceTranscription( + req: Request, + res: Response, + deps: WorkspaceVoiceRouteDeps, +): boolean { + if (!isSupportedAudioContentType(normalizeContentType(req))) { + sendUnsupportedVoiceContentType(res); + return false; + } + return installAdmissionLease(req, res, deps); +} - const settings = loadSettings(deps.boundWorkspace); - try { - validateWorkspaceVoiceState(settings, parsed); - } catch (err) { - if (sendVoiceError(res, err)) return; - throw err; +function voiceAudioBodyParser(): import('express').RequestHandler { + const parse = express.raw({ + type: (req) => + isSupportedAudioContentType( + req.headers['content-type']?.split(';', 1)[0]?.trim().toLowerCase(), + ), + limit: '10mb', + }); + return (req, res, next) => { + const signal = admissionLease(req)?.signal; + let aborted = signal?.aborted === true; + const onAbort = () => { + aborted = true; + if (!res.headersSent) sendWorkspaceDraining(res); + const destroyRequest = () => { + if (!req.destroyed) req.destroy(); + }; + if (res.writableFinished) destroyRequest(); + else { + res.once('finish', destroyRequest); + res.once('close', destroyRequest); } - - const clientId = deps.parseAndValidateClientId(req, res); - if (clientId === null) return; - - try { - const workspaceTrusted = - getWorkspaceTrustStatus(settings.merged, deps.boundWorkspace) - .effective.state === 'trusted'; - await persistVoiceUpdate( - deps, - settings, - parsed, - clientId, - workspaceTrusted, - ); - } catch (err) { - writeStderrLine( - `qwen serve: POST /workspace/voice persist error (workspace=${deps.boundWorkspace}): ${ - err instanceof Error ? err.message : String(err) - }`, - ); - res.status(500).json({ - error: 'Failed to persist voice settings', - code: 'persist_error', - }); + }; + if (aborted) { + onAbort(); + return; + } + signal?.addEventListener('abort', onAbort, { once: true }); + parse(req, res, (err) => { + signal?.removeEventListener('abort', onAbort); + if (aborted || signal?.aborted) { + if (!res.headersSent) sendWorkspaceDraining(res); return; } + next(err); + }); + }; +} - res - .status(200) - .json( - buildWorkspaceVoiceStatus( - deps.boundWorkspace, - loadSettings(deps.boundWorkspace), - ), - ); +export function registerWorkspaceVoiceRoutes( + app: Application, + deps: WorkspaceVoiceRouteDeps, +): void { + app.get('/workspace/voice', (_req: Request, res: Response) => + handleVoiceStatus(res, deps, 'GET /workspace/voice'), + ); + + app.post( + '/workspace/voice', + deps.mutate({ strict: true }), + async (req: Request, res: Response) => { + await handleVoiceUpdate(req, res, deps, 'POST /workspace/voice'); }, ); app.post( '/workspace/voice/transcribe', deps.mutate({ strict: true }), - express.raw({ - type: (req) => - isSupportedAudioContentType( - req.headers['content-type']?.split(';', 1)[0]?.trim().toLowerCase(), - ), - limit: '10mb', - }), + (req: Request, res: Response, next: import('express').NextFunction) => { + if (admitVoiceTranscription(req, res, deps)) next(); + }, + voiceAudioBodyParser(), async (req: Request, res: Response) => { - const contentType = normalizeContentType(req); - if (!isSupportedAudioContentType(contentType)) { - res.status(415).json({ - error: - 'Content-Type must be audio/* or application/octet-stream for voice transcription', - code: 'unsupported_voice_content_type', - }); - return; - } - const audioContentType = contentType; + await handleVoiceTranscription( + req, + res, + deps, + 'POST /workspace/voice/transcribe', + ); + }, + ); +} - const clientId = deps.parseAndValidateClientId(req, res); - if (clientId === null) return; +function createRuntimeVoiceDeps( + runtime: WorkspaceRuntime, + deps: WorkspaceQualifiedVoiceRouteDeps, +): WorkspaceVoiceRouteDeps { + const env = + runtime.env.mode === 'runtime-overlay' + ? (runtime.env.effectiveEnv ?? {}) + : runtime.env.effectiveEnv; + return { + boundWorkspace: runtime.workspaceCwd, + mutate: deps.mutate, + safeBody: deps.safeBody, + persistSetting: deps.persistSetting, + persistSettings: deps.persistSettings, + transcribe: deps.transcribe, + ...(env ? { env } : {}), + scopeOverride: SettingScope.Workspace, + acquireVoiceLease: () => deps.acquireVoiceLease(runtime), + broadcastSettingsChanged: (key, value, scope, clientId) => { + if (runtime.primary) deps.invalidateServeFeaturesCache(); + runtime.bridge.publishWorkspaceEvent({ + type: 'settings_changed', + data: { key, value, scope }, + originatorClientId: clientId, + }); + }, + parseAndValidateClientId: (req, res) => + deps.parseAndValidateClientId(req, res, runtime), + }; +} - const data = readBinaryBody(req); - if (!data || data.byteLength === 0) { - res.status(400).json({ - error: 'Voice audio body must be non-empty binary data', - code: 'invalid_voice_audio', - }); - return; - } +type QualifiedVoiceRequest = Request & { + voiceRouteDeps?: WorkspaceVoiceRouteDeps; +}; - const settings = loadSettings(deps.boundWorkspace); - if (!isVoiceEnabled(settings)) { - res.status(403).json({ - error: 'Voice transcription is disabled for this workspace', - code: 'voice_disabled', - }); - return; - } - const queryVoiceModel = req.query['voiceModel']; - if ( - queryVoiceModel !== undefined && - typeof queryVoiceModel !== 'string' - ) { - res.status(400).json({ - error: '`voiceModel` query parameter must be a string', - code: 'invalid_voice_model', - }); - return; - } - const requestedVoiceModel = - typeof queryVoiceModel === 'string' ? queryVoiceModel.trim() : ''; - if ( - requestedVoiceModel && - requestedVoiceModel.length > MAX_VOICE_MODEL_LENGTH - ) { - res.status(400).json({ - error: `\`voiceModel\` exceeds the ${MAX_VOICE_MODEL_LENGTH}-character limit`, - code: 'invalid_voice_model', - }); - return; - } - const voiceModel = requestedVoiceModel || readVoiceModel(settings); - if (!voiceModel) { - res.status(400).json({ - error: 'A valid voiceModel is required before transcription.', - code: 'voice_model_required', - }); - return; - } +function resolveQualifiedVoiceTarget( + req: Request, + res: Response, + deps: WorkspaceQualifiedVoiceRouteDeps, + includeDraining = false, +): WorkspaceVoiceRouteDeps | undefined { + const runtime = includeDraining + ? resolveManagedWorkspaceRuntimeFromParam(deps.workspaceRegistry, req, res) + : resolveWorkspaceRuntimeFromParam(deps.workspaceRegistry, req, res); + if (!runtime || !requireTrustedWorkspaceRuntime(runtime, res)) return; + return createRuntimeVoiceDeps(runtime, deps); +} - try { - const result = await transcribe({ - data, - mimeType: audioContentType, - voiceModel, - settings, - workspaceCwd: deps.boundWorkspace, - abortSignal: requestAbortSignal(req, res), - }); - res.status(200).json({ v: 1, ...result }); - } catch (err) { - if (sendVoiceError(res, err)) return; - const message = - err instanceof Error - ? sanitizeVoiceErrorMessage(err.message) - : sanitizeVoiceErrorMessage(String(err)); - writeStderrLine( - `qwen serve: POST /workspace/voice/transcribe error (workspace=${deps.boundWorkspace}): ${ - message - }`, - ); - res.status(502).json({ - error: 'Voice transcription failed', - code: 'voice_transcription_failed', - }); - } +export function registerWorkspaceQualifiedVoiceRoutes( + app: Application, + deps: WorkspaceQualifiedVoiceRouteDeps, +): void { + app.get('/workspaces/:workspace/voice', (req: Request, res: Response) => { + const target = resolveQualifiedVoiceTarget(req, res, deps); + if (!target) return; + handleVoiceStatus(res, target, 'GET /workspaces/:workspace/voice'); + }); + + app.post( + '/workspaces/:workspace/voice', + deps.mutate({ strict: true }), + async (req: Request, res: Response) => { + const target = resolveQualifiedVoiceTarget(req, res, deps); + if (!target) return; + await handleVoiceUpdate( + req, + res, + target, + 'POST /workspaces/:workspace/voice', + ); + }, + ); + + app.post( + '/workspaces/:workspace/voice/transcribe', + deps.mutate({ strict: true }), + ( + req: QualifiedVoiceRequest, + res: Response, + next: import('express').NextFunction, + ) => { + const target = resolveQualifiedVoiceTarget(req, res, deps, true); + if (!target) return; + req.voiceRouteDeps = target; + if (admitVoiceTranscription(req, res, target)) next(); + }, + voiceAudioBodyParser(), + async (req: QualifiedVoiceRequest, res: Response) => { + const target = req.voiceRouteDeps; + if (!target) return; + await handleVoiceTranscription( + req, + res, + target, + 'POST /workspaces/:workspace/voice/transcribe', + ); }, ); } diff --git a/packages/cli/src/serve/run-qwen-serve.test.ts b/packages/cli/src/serve/run-qwen-serve.test.ts index 0ed2a0ea654..591071b2e90 100644 --- a/packages/cli/src/serve/run-qwen-serve.test.ts +++ b/packages/cli/src/serve/run-qwen-serve.test.ts @@ -527,9 +527,19 @@ describe('runQwenServe telemetry validation', () => { enabled: false, sensitiveSpanAttributeMaxLength: 1024 * 1024, }); + const shutdownResolvers: Array<() => void> = []; const createBridge = vi .spyOn(acpBridge, 'createAcpSessionBridge') - .mockImplementation(() => makeRuntimeBridge()); + .mockImplementation(() => { + const bridge = makeRuntimeBridge(); + bridge.shutdown = vi.fn( + () => + new Promise((resolve) => { + shutdownResolvers.push(resolve); + }), + ); + return bridge; + }); const handle = await runQwenServe( { @@ -545,6 +555,7 @@ describe('runQwenServe telemetry validation', () => { daemonLogBaseDir: path.join(tmpDir, 'debug'), }, ); + let closing: Promise | undefined; try { const res = await fetch(`${handle.url}/capabilities`); expect(res.status).toBe(200); @@ -574,8 +585,14 @@ describe('runQwenServe telemetry validation', () => { removable: false, }), ]); + + closing = handle.close(); + await vi.waitFor(() => expect(shutdownResolvers).toHaveLength(2)); } finally { - await handle.close(); + closing ??= handle.close(); + await vi.waitFor(() => expect(shutdownResolvers).toHaveLength(2)); + for (const resolve of shutdownResolvers) resolve(); + await closing; } expect(createBridge).toHaveBeenCalledTimes(2); for (const result of createBridge.mock.results) { @@ -585,6 +602,82 @@ describe('runQwenServe telemetry validation', () => { } }); + it('invalidates primary voice capabilities when its workspace service publishes settings changes', async () => { + tmpDir = fs.realpathSync( + fs.mkdtempSync(path.join(os.tmpdir(), 'qws-voice-capability-')), + ); + const workspace = path.join(tmpDir, 'workspace'); + fs.mkdirSync(workspace); + vi.spyOn(qwenCore, 'resolveTelemetrySettings').mockResolvedValue({ + enabled: false, + sensitiveSpanAttributeMaxLength: 1024 * 1024, + }); + vi.spyOn(acpBridge, 'createAcpSessionBridge').mockImplementation(() => + makeRuntimeBridge(), + ); + const originalCreateWorkspaceService = + workspaceServiceRuntime.createDaemonWorkspaceService; + let publishWorkspaceEvent: + | Parameters< + typeof workspaceServiceRuntime.createDaemonWorkspaceService + >[0]['publishWorkspaceEvent'] + | undefined; + vi.spyOn( + workspaceServiceRuntime, + 'createDaemonWorkspaceService', + ).mockImplementation((deps) => { + if (deps.boundWorkspace === canonicalizeWorkspace(workspace)) { + publishWorkspaceEvent = deps.publishWorkspaceEvent; + } + return originalCreateWorkspaceService(deps); + }); + + const handle = await runQwenServe( + { + port: 0, + hostname: '127.0.0.1', + mode: 'http-bridge', + workspace, + serveWebShell: false, + }, + { + preheatBridge: false, + daemonLogBaseDir: path.join(tmpDir, 'debug'), + }, + ); + try { + const before = (await ( + await fetch(`${handle.url}/capabilities`) + ).json()) as { features: string[] }; + expect(before.features).not.toContain('workspace_voice_transcription'); + + fs.mkdirSync(path.join(workspace, '.qwen')); + fs.writeFileSync( + path.join(workspace, '.qwen', 'settings.json'), + JSON.stringify({ + modelProviders: { + openai: [ + { + id: 'qwen3-asr-flash', + baseUrl: 'http://127.0.0.1:65535/v1', + }, + ], + }, + }), + 'utf8', + ); + expect(publishWorkspaceEvent).toBeTypeOf('function'); + publishWorkspaceEvent?.({ type: 'settings_changed', data: {} }); + + const after = (await ( + await fetch(`${handle.url}/capabilities`) + ).json()) as { features: string[] }; + expect(after.features).toContain('workspace_voice_transcription'); + } finally { + await handle.close(); + } + }); + it('adds, advertises, and hot-removes a dynamic workspace runtime', async () => { tmpDir = fs.realpathSync( fs.mkdtempSync(path.join(os.tmpdir(), 'qws-hot-remove-')), diff --git a/packages/cli/src/serve/run-qwen-serve.ts b/packages/cli/src/serve/run-qwen-serve.ts index 269a980cca3..f0fa148f3b7 100644 --- a/packages/cli/src/serve/run-qwen-serve.ts +++ b/packages/cli/src/serve/run-qwen-serve.ts @@ -139,6 +139,7 @@ import type { } from '../commands/channel/pidfile.js'; import { sanitizeLogText } from '@qwen-code/channel-base'; import { isBrowserAutomationMcpAvailable } from './cdp-mcp-command.js'; +import { WorkspaceVoiceCoordinator } from './voice/workspace-voice-coordinator.js'; // Reverse MCP channel; enabled only by explicit option or env opt-in. const QWEN_SERVE_CLIENT_MCP_OVER_WS_ENV = 'QWEN_SERVE_CLIENT_MCP_OVER_WS'; @@ -3255,12 +3256,14 @@ export async function runQwenServe( internalRuntimeBridgesForCleanup.push(bridge); } runtimeBridges.push(bridge); + let invalidatePrimaryServeFeaturesCache = () => {}; const workspaceService = runtime.createDaemonWorkspaceService({ boundWorkspace, contextFilename: contextFilenameForInit ?? 'QWEN.md', statusProvider, workspaceProvidersStatusProvider, workspaceSkillsStatusProvider, + voiceEnv: runtimeEffectiveEnv, isChannelLive: () => bridge.isChannelLive(), persistDisabledTools: persistDisabledToolsFn, persistDisabledSkills: persistDisabledSkillsFn, @@ -3337,7 +3340,15 @@ export async function runQwenServe( bridge.invokeWorkspaceCommand(method, params, invokeOpts), refreshExtensionsForAllSessions: () => bridge.refreshExtensionsForAllSessions(), - publishWorkspaceEvent: (event) => bridge.publishWorkspaceEvent(event), + publishWorkspaceEvent: (event) => { + if ( + event.type === 'settings_changed' || + event.type === 'settings_reloaded' + ) { + invalidatePrimaryServeFeaturesCache(); + } + bridge.publishWorkspaceEvent(event); + }, }); const workspaceRuntimes: WorkspaceRuntime[] = [ @@ -3568,6 +3579,8 @@ export async function runQwenServe( }), workspaceSkillsStatusProvider: runtime.createWorkspaceSkillsStatusProvider(), + voiceEnv: secondaryEnv.effectiveEnv, + voiceSettingsScope: WORKSPACE_SETTING_SCOPE, isChannelLive: () => secondaryBridge.isChannelLive(), preheatAcpChild: () => secondaryBridge.preheat(), persistDisabledTools: persistDisabledToolsFn, @@ -3651,6 +3664,7 @@ export async function runQwenServe( runtime.createWorkspaceRegistry(workspaceRuntimes, { sessionOwnerIndex, }); + const workspaceVoiceCoordinator = new WorkspaceVoiceCoordinator(); core.registerDaemonGaugeCallbacks({ sessionCount: () => @@ -3937,6 +3951,8 @@ export async function runQwenServe( }), workspaceSkillsStatusProvider: runtime.createWorkspaceSkillsStatusProvider(), + voiceEnv: wsEnv.effectiveEnv, + voiceSettingsScope: WORKSPACE_SETTING_SCOPE, isChannelLive: () => wsBridge.isChannelLive(), preheatAcpChild: () => wsBridge.preheat(), persistDisabledTools: persistDisabledToolsFn, @@ -4062,15 +4078,18 @@ export async function runQwenServe( beginDrain(runtimeToDrain: WorkspaceRuntime): void { totalSessionAdmission.beginWorkspaceDrain(runtimeToDrain.workspaceCwd); channelWorkerManager?.beginWorkspaceDrain(runtimeToDrain.workspaceCwd); + workspaceVoiceCoordinator.beginWorkspaceDrain(runtimeToDrain); }, cancelDrain(runtimeToDrain: WorkspaceRuntime): void { channelWorkerManager?.cancelWorkspaceDrain(runtimeToDrain.workspaceCwd); totalSessionAdmission.cancelWorkspaceDrain(runtimeToDrain.workspaceCwd); + workspaceVoiceCoordinator.cancelWorkspaceDrain(runtimeToDrain); }, completeDrain(runtimeToDrain: WorkspaceRuntime): void { totalSessionAdmission.completeWorkspaceDrain( runtimeToDrain.workspaceCwd, ); + workspaceVoiceCoordinator.completeWorkspaceDrain(runtimeToDrain); }, getActivity(runtimeToDrain: WorkspaceRuntime) { return { @@ -4081,6 +4100,8 @@ export async function runQwenServe( channelWorkerManager?.workspaceActivity( runtimeToDrain.workspaceCwd, ) ?? 0, + voiceSessions: + workspaceVoiceCoordinator.getWorkspaceActivity(runtimeToDrain), }; }, disposeRuntime( @@ -4090,6 +4111,10 @@ export async function runQwenServe( const existing = runtimeCleanupPromises.get(runtimeToDrain); if (existing) return existing; const cleanup = (async () => { + await workspaceVoiceCoordinator.disposeRuntime( + runtimeToDrain, + reason, + ); const stopSubSessions = subSessionStoppersByWorkspace.get( runtimeToDrain.workspaceCwd, ); @@ -4174,6 +4199,7 @@ export async function runQwenServe( createWorkspaceRuntime: createDynamicWorkspaceRuntime, workspaceRegistrationStore, workspaceRuntimeRemoval, + voiceCoordinator: workspaceVoiceCoordinator, bridge, webShellDir, boundWorkspace, @@ -4287,6 +4313,12 @@ export async function runQwenServe( }, ), }); + invalidatePrimaryServeFeaturesCache = + ( + app.locals as { + invalidateServeFeaturesCache?: () => void; + } + ).invalidateServeFeaturesCache ?? invalidatePrimaryServeFeaturesCache; // Park the sub-session launcher's stop on app.locals so the close handler // can flip it off before tearing down the bridge it spawns into (symmetric // with stopScheduledTaskKeepalive). Defensive: a launch during drain would @@ -5222,22 +5254,26 @@ export async function runQwenServe( const managedRuntimes = workspaceRegistry.listManaged(); for (const workspaceRuntime of managedRuntimes) { managedRuntimeBridges.add(workspaceRuntime.bridge); - await runtimeRemoval - .disposeRuntime(workspaceRuntime, 'daemon_shutdown') - .catch((err) => { - daemonLog.error( - 'workspace runtime shutdown error', - err instanceof Error ? err : null, - ); - bridgeShutdownError = - err instanceof Error ? err : new Error(String(err)); - try { - workspaceRuntime.bridge.killAllSync(); - } catch { - // Continue shutting down the remaining runtimes. - } - }); } + await Promise.all( + managedRuntimes.map((workspaceRuntime) => + runtimeRemoval + .disposeRuntime(workspaceRuntime, 'daemon_shutdown') + .catch((err) => { + daemonLog.error( + 'workspace runtime shutdown error', + err instanceof Error ? err : null, + ); + bridgeShutdownError = + err instanceof Error ? err : new Error(String(err)); + try { + workspaceRuntime.bridge.killAllSync(); + } catch { + // Continue shutting down the remaining runtimes. + } + }), + ), + ); } for (const bridgeForShutdown of getRuntimeBridgesForCleanup()) { if (managedRuntimeBridges.has(bridgeForShutdown)) continue; diff --git a/packages/cli/src/serve/server.test.ts b/packages/cli/src/serve/server.test.ts index 3b937acdeb7..7df6b977c9d 100644 --- a/packages/cli/src/serve/server.test.ts +++ b/packages/cli/src/serve/server.test.ts @@ -409,6 +409,7 @@ const EXPECTED_REGISTERED_FEATURES = [ 'persistent_workspace_registration', 'workspace_runtime_removal', 'workspace_qualified_rest_core', + 'workspace_qualified_voice', 'extension_management_v2', 'workspace_persisted_transcript', 'workspace_qualified_acp', @@ -2365,6 +2366,34 @@ describe('createServeApp', () => { ); continue; } + if (feature === 'workspace_qualified_voice') { + expect( + predicate({ + multiWorkspaceSessionsEnabled: true, + acpHttpEnabled: true, + }), + ).toBe(true); + expect( + predicate({ + multiWorkspaceSessionsEnabled: true, + acpHttpEnabled: false, + }), + ).toBe(false); + expect(predicate({ multiWorkspaceSessionsEnabled: false })).toBe( + false, + ); + expect(predicate({})).toBe(false); + expect( + getAdvertisedServeFeatures(undefined, { + multiWorkspaceSessionsEnabled: true, + acpHttpEnabled: true, + }), + ).toContain(feature); + expect(getAdvertisedServeFeatures(undefined, {})).not.toContain( + feature, + ); + continue; + } if (feature === 'client_mcp_over_ws') { expect(predicate({ clientMcpOverWsEnabled: true })).toBe(true); expect(predicate({ clientMcpOverWsEnabled: false })).toBe(false); @@ -2891,6 +2920,53 @@ describe('createServeApp', () => { } }); + it('loads workspace environment for direct-embed Voice capability checks', async () => { + const previousQwenHome = process.env['QWEN_HOME']; + const previousWorkspaceAsrKey = process.env['WORKSPACE_ASR_KEY']; + const tempHome = await fsp.mkdtemp( + path.join(os.tmpdir(), 'qwen-voice-capability-env-'), + ); + const workspace = path.join(tempHome, 'workspace'); + try { + await fsp.mkdir(workspace); + await fsp.writeFile( + path.join(workspace, '.env'), + 'WORKSPACE_ASR_KEY=workspace-secret\n', + 'utf8', + ); + process.env['QWEN_HOME'] = tempHome; + resetHomeEnvBootstrapForTesting(); + await fsp.writeFile( + path.join(tempHome, 'settings.json'), + JSON.stringify({ + modelProviders: { + openai: [ + { + id: 'qwen3-asr-flash', + baseUrl: 'https://asr.example/v1', + envKey: 'WORKSPACE_ASR_KEY', + }, + ], + }, + }), + 'utf8', + ); + + const app = createServeApp({ ...baseOpts, workspace }); + const res = await request(app) + .get('/capabilities') + .set('Host', `127.0.0.1:${baseOpts.port}`); + + expect(res.status).toBe(200); + expect(res.body.features).toContain('workspace_voice_transcription'); + } finally { + await fsp.rm(tempHome, { recursive: true, force: true }); + restoreEnv('QWEN_HOME', previousQwenHome); + restoreEnv('WORKSPACE_ASR_KEY', previousWorkspaceAsrKey); + resetHomeEnvBootstrapForTesting(); + } + }); + it('reports disabled prompt queue cap as null in capabilities', async () => { const app = createServeApp({ ...baseOpts, @@ -18168,6 +18244,24 @@ describe('createServeApp ServeAppDeps.fsFactory wiring (#4175 PR 18)', () => { ).not.toThrow(); }); + it('requires the Voice coordinator paired with runtime removal', async () => { + const { createServeApp } = await import('./server.js'); + + expect(() => + createServeApp( + { + port: 0, + hostname: '127.0.0.1', + workspace: '/work/bound', + } as Parameters[0], + () => 0, + { + workspaceRuntimeRemoval: {}, + } as Parameters[2], + ), + ).toThrow(/workspaceRuntimeRemoval requires.*voiceCoordinator/); + }); + it('uses the injected registry sender when client-MCP over WS is enabled', async () => { const { createServeApp } = await import('./server.js'); const runtime = makeInjectedWorkspaceRuntime(); diff --git a/packages/cli/src/serve/server.ts b/packages/cli/src/serve/server.ts index 34f33d67dac..7b3c2d4a64a 100644 --- a/packages/cli/src/serve/server.ts +++ b/packages/cli/src/serve/server.ts @@ -129,10 +129,12 @@ import { registerSseEventsRoutes, } from './routes/sse-events.js'; import { + registerWorkspaceQualifiedVoiceRoutes, registerWorkspaceVoiceRoutes, type WorkspaceVoiceRouteDeps, } from './routes/workspace-voice.js'; import { registerWorkspaceModelsRoutes } from './routes/workspace-models.js'; +import { WorkspaceVoiceCoordinator } from './voice/workspace-voice-coordinator.js'; import { registerA2uiActionRoutes } from './routes/a2ui-action.js'; import { setRateLimiter } from './rate-limit.js'; import { resolveAcpHttpEnabled } from './acp-http-enabled.js'; @@ -313,7 +315,10 @@ function describeRegistryPrimaryForConflict( function getRuntimeEffectiveEnv( metadata: WorkspaceRuntimeEnvMetadata | undefined, ): Readonly> | undefined { - return metadata?.effectiveEnv; + if (!metadata || metadata.mode === 'parent-process') { + return metadata?.effectiveEnv; + } + return metadata.effectiveEnv ?? {}; } export interface ServeAppDeps { @@ -467,6 +472,7 @@ export interface ServeAppDeps { primaryWorkspaceTrusted?: boolean; primaryRuntimeEnv?: WorkspaceRuntimeEnvMetadata; voiceTranscriber?: WorkspaceVoiceRouteDeps['transcribe']; + voiceCoordinator?: WorkspaceVoiceCoordinator; } /** @@ -530,6 +536,11 @@ export function createServeApp( getPort: () => number = () => opts.port, deps: ServeAppDeps = {}, ): Application { + if (deps.workspaceRuntimeRemoval && !deps.voiceCoordinator) { + throw new Error( + 'createServeApp: deps.workspaceRuntimeRemoval requires the matching deps.voiceCoordinator.', + ); + } const app = express(); // Forward `maxSessions` into the default-constructed bridge so // direct callers of `createServeApp` (tests, embeds) get the same @@ -698,6 +709,11 @@ export function createServeApp( deps.workspaceRuntimeRemoval !== undefined, ...(primaryEffectiveEnv ? { env: primaryEffectiveEnv } : {}), }); + ( + app.locals as { + invalidateServeFeaturesCache?: () => void; + } + ).invalidateServeFeaturesCache = invalidateServeFeaturesCache; const statusProvider = deps.statusProvider ?? createDaemonStatusProvider( @@ -804,6 +820,7 @@ export function createServeApp( primaryEffectiveEnv ? { env: primaryEffectiveEnv } : {}, ), workspaceSkillsStatusProvider: createWorkspaceSkillsStatusProvider(), + ...(primaryEffectiveEnv ? { voiceEnv: primaryEffectiveEnv } : {}), isChannelLive: () => bridge.isChannelLive(), persistDisabledTools: deps.persistDisabledTools ?? @@ -864,6 +881,8 @@ export function createServeApp( (app.locals as { workspaceRegistry?: WorkspaceRegistry }).workspaceRegistry = workspaceRegistry; const primaryRuntime = workspaceRegistry.primary; + const voiceCoordinator = + deps.voiceCoordinator ?? new WorkspaceVoiceCoordinator(); const primaryBoundWorkspace = primaryRuntime.workspaceCwd; const primaryBridge = primaryRuntime.bridge; const primaryWorkspace = primaryRuntime.workspaceService; @@ -1314,10 +1333,24 @@ export function createServeApp( persistSetting: deps.persistSetting, persistSettings: deps.persistSettings, transcribe: deps.voiceTranscriber, + env: getRuntimeEffectiveEnv(primaryRuntime.env), + acquireVoiceLease: () => voiceCoordinator.acquire(primaryRuntime), broadcastSettingsChanged, parseAndValidateClientId: (req, res) => parseAndValidateWorkspaceClientId(req, res, primaryBridge), }); + registerWorkspaceQualifiedVoiceRoutes(app, { + workspaceRegistry, + mutate, + safeBody, + persistSetting: deps.persistSetting, + persistSettings: deps.persistSettings, + transcribe: deps.voiceTranscriber, + acquireVoiceLease: (runtime) => voiceCoordinator.acquire(runtime), + parseAndValidateClientId: (req, res, runtime) => + parseAndValidateWorkspaceClientId(req, res, runtime.bridge), + invalidateServeFeaturesCache, + }); if (deps.persistSettings) { registerWorkspaceModelsRoutes(app, { boundWorkspace: primaryBoundWorkspace, @@ -1681,9 +1714,15 @@ export function createServeApp( path: '/voice/stream', onConnection: createVoiceWsConnectionHandler(primaryBoundWorkspace, { env: getRuntimeEffectiveEnv(primaryRuntime.env), + acquireVoiceLease: () => voiceCoordinator.acquire(primaryRuntime), }), }, ], + workspaceVoiceConnection: (runtime, ws, req) => + createVoiceWsConnectionHandler(runtime.workspaceCwd, { + env: getRuntimeEffectiveEnv(runtime.env), + acquireVoiceLease: () => voiceCoordinator.acquire(runtime), + })(ws, req), }); if (acpHandleRef.current) { app.locals['acpHandle'] = acpHandleRef.current; diff --git a/packages/cli/src/serve/server/serve-features.ts b/packages/cli/src/serve/server/serve-features.ts index 6c9e5d515ce..d576d3dff68 100644 --- a/packages/cli/src/serve/server/serve-features.ts +++ b/packages/cli/src/serve/server/serve-features.ts @@ -82,7 +82,11 @@ export function createServeFeatures( }; const getCachedVoiceTranscriptionAvailable = () => { cachedVoiceTranscriptionAvailable ??= - isWorkspaceVoiceTranscriptionAvailable(boundWorkspace); + isWorkspaceVoiceTranscriptionAvailable( + boundWorkspace, + env, + deps.env !== undefined, + ); return cachedVoiceTranscriptionAvailable; }; @@ -130,10 +134,13 @@ export function createServeFeatures( function isWorkspaceVoiceTranscriptionAvailable( boundWorkspace: string, + env: Readonly>, + skipLoadEnvironment: boolean, ): boolean { try { return hasConfiguredBatchVoiceTranscriptionModel( - loadSettings(boundWorkspace), + loadSettings(boundWorkspace, { skipLoadEnvironment }), + { env }, ); } catch (err) { writeStderrLine( diff --git a/packages/cli/src/serve/server/telemetry.test.ts b/packages/cli/src/serve/server/telemetry.test.ts index 606284b4597..d8415adc9c8 100644 --- a/packages/cli/src/serve/server/telemetry.test.ts +++ b/packages/cli/src/serve/server/telemetry.test.ts @@ -258,6 +258,32 @@ describe('daemonTelemetryMiddleware — recordRequest seam', () => { } }); + it('attributes plural workspace voice requests to the selected workspace', () => { + const mw = daemonTelemetryMiddleware(() => '/workspace/secondary'); + for (const [method, path, route] of [ + ['GET', '/workspaces/ws-secondary/voice', 'GET /workspace/voice'], + ['POST', '/workspaces/ws-secondary/voice', 'POST /workspace/voice'], + [ + 'POST', + '/workspaces/ws-secondary/voice/transcribe', + 'POST /workspace/voice/transcribe', + ], + ] as const) { + const res = mockRes(200); + mw(mockReq(method, path), res, vi.fn() as unknown as NextFunction); + res.emit('finish'); + + expect(coreMocks.withDaemonRequestSpan).toHaveBeenLastCalledWith( + expect.objectContaining({ + method, + route, + workspaceHash: 'hash:/workspace/secondary', + }), + expect.any(Function), + ); + } + }); + it('excludes the dashboard status poll (GET /daemon/status) from recordRequest', () => { const recordRequest = vi.fn(); const mw = daemonTelemetryMiddleware(() => '/ws', recordRequest); diff --git a/packages/cli/src/serve/server/telemetry.ts b/packages/cli/src/serve/server/telemetry.ts index 06366d4f11f..ae9b8841b3f 100644 --- a/packages/cli/src/serve/server/telemetry.ts +++ b/packages/cli/src/serve/server/telemetry.ts @@ -154,6 +154,7 @@ export function resolveDaemonTelemetryRoute( suffix === '/workspace/preflight' || suffix === '/workspace/hooks' || suffix === '/workspace/settings' || + suffix === '/workspace/voice' || suffix === '/workspace/permissions' || suffix === '/workspace/trust' || suffix === '/workspace/memory' || @@ -181,6 +182,8 @@ export function resolveDaemonTelemetryRoute( if (req.method === 'POST') { if ( suffix === '/workspace/settings' || + suffix === '/workspace/voice' || + suffix === '/workspace/voice/transcribe' || suffix === '/workspace/permissions' || suffix === '/workspace/trust/request' || suffix === '/workspace/init' || diff --git a/packages/cli/src/serve/voice/voice-ws.test.ts b/packages/cli/src/serve/voice/voice-ws.test.ts index 616f9d2fdab..c22a929789f 100644 --- a/packages/cli/src/serve/voice/voice-ws.test.ts +++ b/packages/cli/src/serve/voice/voice-ws.test.ts @@ -8,8 +8,10 @@ import { afterEach, describe, it, expect, vi } from 'vitest'; import { createVoiceWsConnectionHandler } from './voice-ws.js'; +import { WorkspaceVoiceCoordinator } from './workspace-voice-coordinator.js'; import type { DaemonVoiceContext } from './resolve-voice-config.js'; import type { VoiceStreamSession } from '../../ui/voice/voice-stream-session.js'; +import type { WorkspaceRuntime } from '../workspace-registry.js'; /** Minimal stand-in for a `ws` WebSocket the handler attaches to. */ class FakeWs { @@ -17,6 +19,7 @@ class FakeWs { readyState = 1; readonly sent: Array = []; closeCode: number | undefined; + closeReason: string | undefined; private handlers: Record void>> = {}; constructor(private readonly emitCloseOnClose = true) {} @@ -28,8 +31,9 @@ class FakeWs { send(data: string | Uint8Array): void { this.sent.push(typeof data === 'string' ? data : 'binary'); } - close(code?: number): void { + close(code?: number, reason?: string): void { this.closeCode = code; + this.closeReason = reason; this.readyState = 3; if (this.emitCloseOnClose) this.emit('close'); } @@ -335,6 +339,38 @@ describe('createVoiceWsConnectionHandler', () => { expect(session.abort).toHaveBeenCalledOnce(); }); + it('aborts an in-flight upstream open and releases the lease on disconnect', async () => { + const coordinator = new WorkspaceVoiceCoordinator(); + const runtime = { workspaceId: 'secondary' } as WorkspaceRuntime; + let openSignal: AbortSignal | undefined; + const ws = new FakeWs(); + const handler = createVoiceWsConnectionHandler('/ws', { + loadContext: () => streamingCtx(), + acquireVoiceLease: () => coordinator.acquire(runtime), + openStream: async (_ctx, _callbacks, abortSignal) => { + openSignal = abortSignal; + return await new Promise((_resolve, reject) => { + abortSignal?.addEventListener( + 'abort', + () => reject(new Error('aborted')), + { once: true }, + ); + }); + }, + }); + handler(ws as never, {} as never); + ws.text({ type: 'start' }); + await tick(); + expect(coordinator.getWorkspaceActivity(runtime)).toBe(1); + + ws.close(); + await tick(); + await tick(); + + expect(openSignal?.aborted).toBe(true); + expect(coordinator.getWorkspaceActivity(runtime)).toBe(0); + }); + it('aborts the streaming session if finalization times out', async () => { vi.useFakeTimers(); const session: VoiceStreamSession = { @@ -472,4 +508,80 @@ describe('createVoiceWsConnectionHandler', () => { handler(next as never, {} as never); expect(next.closeCode).not.toBe(1013); }); + + it('releases a finalized session even when the socket never emits close', async () => { + const handler = createVoiceWsConnectionHandler('/ws', { + loadContext: () => streamingCtx(), + openStream: async () => ({ + pushAudio: vi.fn(), + finish: vi.fn(async () => 'done'), + abort: vi.fn(), + }), + }); + const finalized = new FakeWs(false); + handler(finalized as never, {} as never); + finalized.text({ type: 'stop' }); + await tick(); + await tick(); + expect(finalized.closeCode).toBe(1000); + + const next = Array.from({ length: 8 }, () => new FakeWs()); + for (const ws of next) handler(ws as never, {} as never); + expect(next.every((ws) => ws.closeCode !== 1013)).toBe(true); + }); + + it('closes with restart semantics when the target runtime is draining', () => { + const ws = new FakeWs(); + const handler = createVoiceWsConnectionHandler('/ws', { + acquireVoiceLease: () => ({ kind: 'rejected', reason: 'draining' }), + }); + + handler(ws as never, {} as never); + + expect(ws.closeCode).toBe(1012); + expect(ws.frames()).toEqual([]); + }); + + it('closes with busy semantics when Voice capacity is exhausted', () => { + const ws = new FakeWs(); + const handler = createVoiceWsConnectionHandler('/ws', { + acquireVoiceLease: () => ({ kind: 'rejected', reason: 'capacity' }), + }); + + handler(ws as never, {} as never); + + expect(ws.closeCode).toBe(1013); + expect(ws.frames()).toEqual([expect.objectContaining({ type: 'error' })]); + }); + + it('closes only the disposed runtime session and releases its lease', async () => { + const coordinator = new WorkspaceVoiceCoordinator(); + const runtime = { workspaceId: 'secondary' } as WorkspaceRuntime; + const ws = new FakeWs(false); + const handler = createVoiceWsConnectionHandler('/ws', { + acquireVoiceLease: () => coordinator.acquire(runtime), + }); + handler(ws as never, {} as never); + expect(coordinator.getWorkspaceActivity(runtime)).toBe(1); + + await coordinator.disposeRuntime(runtime, 'workspace_removed'); + + expect(ws.closeCode).toBe(1012); + expect(coordinator.getWorkspaceActivity(runtime)).toBe(0); + }); + + it('reports daemon shutdown when its runtime lease is aborted on shutdown', async () => { + const coordinator = new WorkspaceVoiceCoordinator(); + const runtime = { workspaceId: 'primary' } as WorkspaceRuntime; + const ws = new FakeWs(false); + const handler = createVoiceWsConnectionHandler('/ws', { + acquireVoiceLease: () => coordinator.acquire(runtime), + }); + handler(ws as never, {} as never); + + await coordinator.disposeRuntime(runtime, 'daemon_shutdown'); + + expect(ws.closeCode).toBe(1012); + expect(ws.closeReason).toBe('Server shutting down'); + }); }); diff --git a/packages/cli/src/serve/voice/voice-ws.ts b/packages/cli/src/serve/voice/voice-ws.ts index 924f52196de..7ef9c166a0d 100644 --- a/packages/cli/src/serve/voice/voice-ws.ts +++ b/packages/cli/src/serve/voice/voice-ws.ts @@ -24,6 +24,12 @@ import type { VoiceStreamCallbacks, VoiceStreamSession, } from '../../ui/voice/voice-stream-session.js'; +import { + MAX_CONCURRENT_VOICE_SESSIONS, + type VoiceAdmissionLease, + type VoiceAdmissionResult, + VoiceLeaseAbortError, +} from './workspace-voice-coordinator.js'; const debugLogger = createDebugLogger('VOICE_WS'); @@ -34,10 +40,6 @@ const MAX_QUEUED_AUDIO_BYTES = MAX_BATCH_AUDIO_BYTES * 2; // Hard cap on a single voice connection so a client that opens the socket and // never sends `stop` can't pin an upstream ASR session indefinitely. const MAX_CONNECTION_MS = 6 * 60_000; -// Voice WS bypasses the ACP connection registry; cap concurrent sessions so a -// client can't open unbounded sockets (each opens an upstream ASR connection). -// Generous for one interactive user across a few tabs. -const MAX_CONCURRENT_VOICE_SESSIONS = 8; const GENERIC_TRANSCRIPTION_ERROR = 'Voice transcription failed. Please try again.'; const NO_VOICE_MODEL_ERROR = 'No voice model is configured for this workspace.'; @@ -72,13 +74,20 @@ export interface VoiceWsDeps { openStream?: ( ctx: DaemonVoiceContext, callbacks: VoiceStreamCallbacks, + abortSignal?: AbortSignal, ) => Promise; - transcribe?: (ctx: DaemonVoiceContext, pcm: Uint8Array) => Promise; + transcribe?: ( + ctx: DaemonVoiceContext, + pcm: Uint8Array, + abortSignal?: AbortSignal, + ) => Promise; + acquireVoiceLease?: () => VoiceAdmissionResult; } async function defaultOpenStream( ctx: DaemonVoiceContext, callbacks: VoiceStreamCallbacks, + abortSignal?: AbortSignal, ): Promise { try { const cfg = resolveVoiceStreamConfig({ @@ -87,11 +96,13 @@ async function defaultOpenStream( voiceModel: ctx.voiceModel, env: ctx.env, }); - await assertVoiceBaseUrlNetworkAllowed(cfg); - return await openVoiceStreamWithRetry(() => - cfg.transport === 'qwen-asr-realtime' - ? openQwenAsrRealtimeStream(cfg, callbacks) - : openVoiceStream(cfg, callbacks), + await assertVoiceBaseUrlNetworkAllowed(cfg, undefined, abortSignal); + return await openVoiceStreamWithRetry( + () => + cfg.transport === 'qwen-asr-realtime' + ? openQwenAsrRealtimeStream(cfg, callbacks, { abortSignal }) + : openVoiceStream(cfg, callbacks, { abortSignal }), + { abortSignal }, ); } catch (error) { debugLogger.debug(`[voice-ws] stream open error: ${errMessage(error)}`); @@ -102,6 +113,7 @@ async function defaultOpenStream( function defaultTranscribe( ctx: DaemonVoiceContext, pcm: Uint8Array, + abortSignal?: AbortSignal, ): Promise { return transcribeVoiceAudio( { data: encodeWav(pcm), mimeType: 'audio/wav' }, @@ -110,6 +122,7 @@ function defaultTranscribe( settings: ctx.settings, voiceModel: ctx.voiceModel, env: ctx.env, + abortSignal, }, ).catch((error: unknown) => { debugLogger.debug( @@ -123,6 +136,13 @@ function errMessage(error: unknown): string { return error instanceof Error ? error.message : String(error); } +function voiceLeaseCloseReason(signal: AbortSignal): string { + return signal.reason instanceof VoiceLeaseAbortError && + signal.reason.kind === 'daemon_shutdown' + ? 'Server shutting down' + : 'Workspace removed'; +} + function voiceConfigErrorMessage(error: unknown): string { const message = errMessage(error); return message === NO_VOICE_MODEL_ERROR @@ -184,15 +204,36 @@ export function createVoiceWsConnectionHandler( loadDaemonVoiceContext(workspaceCwd, { env: deps.env })); const openStream = deps.openStream ?? defaultOpenStream; const transcribe = deps.transcribe ?? defaultTranscribe; - // Shared across all connections from this daemon (factory closure). + // Kept only for direct embeds that have not supplied the daemon-wide + // coordinator. Production always uses `acquireVoiceLease`. let activeSessions = 0; + const acquireLocalLease = (): VoiceAdmissionResult => { + if (activeSessions >= MAX_CONCURRENT_VOICE_SESSIONS) { + return { kind: 'rejected', reason: 'capacity' }; + } + activeSessions++; + const controller = new AbortController(); + let released = false; + const lease: VoiceAdmissionLease = { + signal: controller.signal, + release: () => { + if (released) return; + released = true; + activeSessions--; + }, + }; + return { kind: 'admitted', lease }; + }; return (ws: WebSocket) => { - if (activeSessions >= MAX_CONCURRENT_VOICE_SESSIONS) { - writeStderrLine( - `qwen serve: voice websocket rejected; activeSessions=${activeSessions}`, - ); + const admission = (deps.acquireVoiceLease ?? acquireLocalLease)(); + if (admission.kind === 'rejected') { + writeStderrLine('qwen serve: voice websocket rejected'); try { + if (admission.reason === 'draining') { + ws.close(1012, 'Workspace removed'); + return; + } ws.send( JSON.stringify({ type: 'error', @@ -205,18 +246,14 @@ export function createVoiceWsConnectionHandler( } return; } - activeSessions++; - writeStderrLine( - `qwen serve: voice websocket accepted; activeSessions=${activeSessions}`, - ); + const lease = admission.lease; + writeStderrLine('qwen serve: voice websocket accepted'); let released = false; const releaseSlot = () => { if (!released) { released = true; - activeSessions--; - writeStderrLine( - `qwen serve: voice websocket slot released; activeSessions=${activeSessions}`, - ); + lease.release(); + writeStderrLine('qwen serve: voice websocket slot released'); } }; @@ -228,6 +265,7 @@ export function createVoiceWsConnectionHandler( let bufferedBytes = 0; let queuedBytes = 0; let pendingOperations = 0; + const operationController = new AbortController(); // Serialize message handling so async start/push/finalize never interleave. let chain: Promise = Promise.resolve(); @@ -258,6 +296,9 @@ export function createVoiceWsConnectionHandler( function cleanup(): void { state = 'closed'; clearTimeout(hardTimer); + if (!operationController.signal.aborted) { + operationController.abort(new Error('Voice connection closed.')); + } if (session) { try { session.abort(); @@ -272,6 +313,21 @@ export function createVoiceWsConnectionHandler( queuedBytes = 0; } + lease.signal.addEventListener( + 'abort', + () => { + if (state === 'closed') return; + cleanup(); + try { + ws.close(1012, voiceLeaseCloseReason(lease.signal)); + } catch { + // ignore + } + releaseSlotWhenIdle(); + }, + { once: true }, + ); + function fail(message: string): void { if (state === 'closed') return; writeStderrLine(`qwen serve: voice websocket failed: ${message}`); @@ -311,7 +367,7 @@ export function createVoiceWsConnectionHandler( fail(GENERIC_TRANSCRIPTION_ERROR); }, }; - const opening = openStream(ctx, callbacks); + const opening = openStream(ctx, callbacks, operationController.signal); sessionPromise = opening; const opened = await opening; if (state === 'closed') { @@ -345,11 +401,16 @@ export function createVoiceWsConnectionHandler( } } } else if (pcmChunks.length > 0) { - transcript = await transcribe(ctx!, Buffer.concat(pcmChunks)); + transcript = await transcribe( + ctx!, + Buffer.concat(pcmChunks), + operationController.signal, + ); } sendJson({ type: 'final', text: transcript }); writeStderrLine('qwen serve: voice websocket finalized successfully'); cleanup(); + releaseSlotWhenIdle(); try { ws.close(1000, 'done'); } catch { diff --git a/packages/cli/src/serve/voice/workspace-voice-coordinator.test.ts b/packages/cli/src/serve/voice/workspace-voice-coordinator.test.ts new file mode 100644 index 00000000000..1426459c2cb --- /dev/null +++ b/packages/cli/src/serve/voice/workspace-voice-coordinator.test.ts @@ -0,0 +1,173 @@ +/** + * @license + * Copyright 2026 Qwen Team + * SPDX-License-Identifier: Apache-2.0 + */ + +import { describe, expect, it, vi } from 'vitest'; +import type { WorkspaceRuntime } from '../workspace-registry.js'; +import { + MAX_CONCURRENT_VOICE_SESSIONS, + VoiceLeaseAbortError, + WorkspaceVoiceCoordinator, +} from './workspace-voice-coordinator.js'; + +function runtime(id: string): WorkspaceRuntime { + return { workspaceId: id } as WorkspaceRuntime; +} + +describe('WorkspaceVoiceCoordinator', () => { + it('shares capacity across runtimes and releases it exactly once', () => { + const coordinator = new WorkspaceVoiceCoordinator(); + const first = runtime('first'); + const second = runtime('second'); + const leases = Array.from({ length: MAX_CONCURRENT_VOICE_SESSIONS }, () => + coordinator.acquire(first), + ); + expect(leases.every((item) => item.kind === 'admitted')).toBe(true); + expect(coordinator.acquire(second)).toEqual({ + kind: 'rejected', + reason: 'capacity', + }); + const lease = leases[0]; + if (!lease || lease.kind !== 'admitted') throw new Error('expected lease'); + lease.lease.release(); + lease.lease.release(); + expect(coordinator.acquire(second).kind).toBe('admitted'); + }); + + it('drains new work without aborting existing work and aborts on disposal', async () => { + const coordinator = new WorkspaceVoiceCoordinator(); + const target = runtime('target'); + const admitted = coordinator.acquire(target); + if (admitted.kind !== 'admitted') throw new Error('expected lease'); + coordinator.beginWorkspaceDrain(target); + expect(admitted.lease.signal.aborted).toBe(false); + expect(coordinator.getWorkspaceActivity(target)).toBe(1); + expect(coordinator.acquire(target)).toEqual({ + kind: 'rejected', + reason: 'draining', + }); + const dispose = coordinator.disposeRuntime(target, 'workspace_removed'); + expect(admitted.lease.signal.aborted).toBe(true); + expect(admitted.lease.signal.reason).toMatchObject({ + name: 'VoiceLeaseAbortError', + kind: 'workspace_removed', + }); + admitted.lease.release(); + await dispose; + }); + + it('stops disposal waiting after five seconds without releasing an active lease', async () => { + vi.useFakeTimers(); + try { + const coordinator = new WorkspaceVoiceCoordinator(); + const target = runtime('target'); + const admitted = coordinator.acquire(target); + if (admitted.kind !== 'admitted') throw new Error('expected lease'); + + const dispose = coordinator.disposeRuntime(target, 'workspace_removed'); + await vi.advanceTimersByTimeAsync(5_000); + await dispose; + + expect(coordinator.getWorkspaceActivity(target)).toBe(1); + expect(vi.getTimerCount()).toBe(0); + admitted.lease.release(); + expect(coordinator.getWorkspaceActivity(target)).toBe(0); + expect(coordinator['states'].has(target)).toBe(false); + } finally { + vi.useRealTimers(); + } + }); + + it('aborts active leases when a workspace drain completes', () => { + const coordinator = new WorkspaceVoiceCoordinator(); + const target = runtime('target'); + const admitted = coordinator.acquire(target); + if (admitted.kind !== 'admitted') throw new Error('expected lease'); + + coordinator.completeWorkspaceDrain(target); + + expect(admitted.lease.signal.aborted).toBe(true); + expect(coordinator.getWorkspaceActivity(target)).toBe(1); + admitted.lease.release(); + expect(coordinator.getWorkspaceActivity(target)).toBe(0); + expect(coordinator['states'].has(target)).toBe(false); + }); + + it('cancels an unfinished drain but cannot reactivate a completed drain', () => { + const coordinator = new WorkspaceVoiceCoordinator(); + const target = runtime('target'); + + coordinator.beginWorkspaceDrain(target); + expect(coordinator.acquire(target)).toEqual({ + kind: 'rejected', + reason: 'draining', + }); + + coordinator.cancelWorkspaceDrain(target); + expect(coordinator['states'].has(target)).toBe(false); + const admitted = coordinator.acquire(target); + if (admitted.kind !== 'admitted') throw new Error('expected lease'); + + coordinator.completeWorkspaceDrain(target); + coordinator.cancelWorkspaceDrain(target); + expect(coordinator['states'].get(target)?.draining).toBe(true); + expect(coordinator.acquire(target)).toEqual({ + kind: 'rejected', + reason: 'draining', + }); + admitted.lease.release(); + }); + + it('keeps a re-added runtime generation independent from old leases', async () => { + const coordinator = new WorkspaceVoiceCoordinator(); + const oldRuntime = runtime('same-id'); + const newRuntime = runtime('same-id'); + const oldAdmission = coordinator.acquire(oldRuntime); + if (oldAdmission.kind !== 'admitted') throw new Error('expected lease'); + + coordinator.beginWorkspaceDrain(oldRuntime); + const newAdmission = coordinator.acquire(newRuntime); + expect(newAdmission.kind).toBe('admitted'); + const disposal = coordinator.disposeRuntime( + oldRuntime, + 'workspace_removed', + ); + oldAdmission.lease.release(); + await disposal; + + expect(coordinator.getWorkspaceActivity(oldRuntime)).toBe(0); + expect(coordinator.getWorkspaceActivity(newRuntime)).toBe(1); + if (newAdmission.kind === 'admitted') newAdmission.lease.release(); + }); + + it('rejects a disposed runtime that never acquired a lease', async () => { + const coordinator = new WorkspaceVoiceCoordinator(); + const target = runtime('failed-construction'); + + await coordinator.disposeRuntime(target, 'workspace_removed'); + + expect(coordinator.getWorkspaceActivity(target)).toBe(0); + expect(coordinator.acquire(target)).toEqual({ + kind: 'rejected', + reason: 'draining', + }); + }); + + it('tags daemon shutdown aborts independently of their message', async () => { + const coordinator = new WorkspaceVoiceCoordinator(); + const target = runtime('target'); + const admitted = coordinator.acquire(target); + if (admitted.kind !== 'admitted') throw new Error('expected lease'); + + const disposal = coordinator.disposeRuntime(target, 'daemon_shutdown'); + + expect(admitted.lease.signal.reason).toBeInstanceOf(VoiceLeaseAbortError); + expect(admitted.lease.signal.reason).toMatchObject({ + kind: 'daemon_shutdown', + }); + admitted.lease.release(); + await disposal; + }); +}); diff --git a/packages/cli/src/serve/voice/workspace-voice-coordinator.ts b/packages/cli/src/serve/voice/workspace-voice-coordinator.ts new file mode 100644 index 00000000000..ff67a471634 --- /dev/null +++ b/packages/cli/src/serve/voice/workspace-voice-coordinator.ts @@ -0,0 +1,186 @@ +/** + * @license + * Copyright 2026 Qwen Team + * SPDX-License-Identifier: Apache-2.0 + */ + +import type { WorkspaceRuntime } from '../workspace-registry.js'; + +export const MAX_CONCURRENT_VOICE_SESSIONS = 8; +const DISPOSE_WAIT_MS = 5_000; + +export class VoiceLeaseAbortError extends Error { + constructor( + readonly kind: 'daemon_shutdown' | 'workspace_removed', + message: string, + ) { + super(message); + this.name = 'VoiceLeaseAbortError'; + } +} + +export interface VoiceAdmissionLease { + readonly signal: AbortSignal; + release(): void; +} + +export type VoiceAdmissionResult = + | { readonly kind: 'admitted'; readonly lease: VoiceAdmissionLease } + | { readonly kind: 'rejected'; readonly reason: 'draining' | 'capacity' }; + +interface RuntimeVoiceState { + draining: boolean; + completed: boolean; + leases: Set; + idleWaiters: Set<() => void>; +} + +class Lease implements VoiceAdmissionLease { + readonly controller = new AbortController(); + readonly signal = this.controller.signal; + private released = false; + + constructor(private readonly onRelease: (lease: Lease) => void) {} + + release(): void { + if (this.released) return; + this.released = true; + this.onRelease(this); + } + + abort(reason: Error): void { + if (!this.signal.aborted) this.controller.abort(reason); + } +} + +export class WorkspaceVoiceCoordinator { + private readonly states = new Map(); + private readonly disposed = new WeakSet(); + private active = 0; + + acquire(runtime: WorkspaceRuntime): VoiceAdmissionResult { + if (this.disposed.has(runtime)) { + return { kind: 'rejected', reason: 'draining' }; + } + const state = this.stateFor(runtime); + if (state.draining) return { kind: 'rejected', reason: 'draining' }; + if (this.active >= MAX_CONCURRENT_VOICE_SESSIONS) { + return { kind: 'rejected', reason: 'capacity' }; + } + + const lease = new Lease((current) => this.release(runtime, current)); + state.leases.add(lease); + this.active++; + return { kind: 'admitted', lease }; + } + + beginWorkspaceDrain(runtime: WorkspaceRuntime): void { + this.stateFor(runtime).draining = true; + } + + cancelWorkspaceDrain(runtime: WorkspaceRuntime): void { + const state = this.states.get(runtime); + if (!state || state.completed) return; + state.draining = false; + if (state.leases.size === 0) this.states.delete(runtime); + } + + completeWorkspaceDrain(runtime: WorkspaceRuntime): void { + this.disposed.add(runtime); + const state = this.states.get(runtime); + if (!state) return; + state.draining = true; + state.completed = true; + const abortReason = new VoiceLeaseAbortError( + 'workspace_removed', + 'Workspace drain completed', + ); + for (const lease of state.leases) lease.abort(abortReason); + this.deleteIfIdle(runtime, state); + } + + getWorkspaceActivity(runtime: WorkspaceRuntime): number { + return this.states.get(runtime)?.leases.size ?? 0; + } + + async disposeRuntime( + runtime: WorkspaceRuntime, + reason: 'daemon_shutdown' | 'workspace_removed', + ): Promise { + this.disposed.add(runtime); + const state = this.states.get(runtime); + if (!state) return; + state.draining = true; + state.completed = true; + const abortReason = new VoiceLeaseAbortError( + reason, + reason === 'workspace_removed' + ? 'Workspace runtime was removed' + : 'Daemon is shutting down', + ); + for (const lease of state.leases) lease.abort(abortReason); + if (state.leases.size === 0) { + this.deleteIfIdle(runtime, state); + return; + } + let resolveIdle: (() => void) | undefined; + const idle = new Promise((resolve) => { + resolveIdle = resolve; + state.idleWaiters.add(resolve); + }); + let timeout: ReturnType | undefined; + let timedOut = false; + try { + await Promise.race([ + idle, + new Promise((resolve) => { + timeout = setTimeout(() => { + timedOut = true; + resolve(); + }, DISPOSE_WAIT_MS); + timeout.unref?.(); + }), + ]); + } finally { + if (timeout) clearTimeout(timeout); + if (resolveIdle) state.idleWaiters.delete(resolveIdle); + } + if (timedOut && state.leases.size > 0) { + process.stderr.write( + `qwen serve: Voice runtime disposal timed out with ${state.leases.size} active lease(s).\n`, + ); + } + } + + private stateFor(runtime: WorkspaceRuntime): RuntimeVoiceState { + let state = this.states.get(runtime); + if (!state) { + state = { + draining: false, + completed: false, + leases: new Set(), + idleWaiters: new Set(), + }; + this.states.set(runtime, state); + } + return state; + } + + private release(runtime: WorkspaceRuntime, lease: Lease): void { + const state = this.states.get(runtime); + if (!state || !state.leases.delete(lease)) return; + this.active--; + if (state.leases.size === 0) { + for (const resolve of state.idleWaiters) resolve(); + state.idleWaiters.clear(); + } + this.deleteIfIdle(runtime, state); + } + + private deleteIfIdle( + runtime: WorkspaceRuntime, + state: RuntimeVoiceState, + ): void { + if (state.completed && state.leases.size === 0) this.states.delete(runtime); + } +} diff --git a/packages/cli/src/serve/workspace-service/__tests__/facade.test.ts b/packages/cli/src/serve/workspace-service/__tests__/facade.test.ts index 0cec1c0d4ca..fae84ed1a07 100644 --- a/packages/cli/src/serve/workspace-service/__tests__/facade.test.ts +++ b/packages/cli/src/serve/workspace-service/__tests__/facade.test.ts @@ -274,6 +274,43 @@ describe('createDaemonWorkspaceService', () => { }); }); + it('forces qualified ACP Voice writes into the workspace scope', async () => { + await withIsolatedQwenHome(async () => { + const persistSettings = vi.fn(async () => {}); + const svc = createDaemonWorkspaceService( + makeDeps({ + persistSettings, + voiceSettingsScope: SettingScope.Workspace, + voiceEnv: {}, + }), + ); + + await svc.setWorkspaceVoiceSettings(makeCtx(), { + enabled: false, + mode: 'tap', + language: 'english', + }); + + expect(persistSettings).toHaveBeenCalledWith('/workspace', [ + { + scope: SettingScope.Workspace, + key: 'general.voice.mode', + value: 'tap', + }, + { + scope: SettingScope.Workspace, + key: 'general.voice.language', + value: 'english', + }, + { + scope: SettingScope.Workspace, + key: 'general.voice.enabled', + value: false, + }, + ]); + }); + }); + it('rejects invalid voice settings before persisting', async () => { await withIsolatedQwenHome(async () => { const persistSettings = vi.fn(async () => {}); @@ -296,7 +333,7 @@ describe('createDaemonWorkspaceService', () => { }); }); - it('does not publish fallback voice writes when a later write fails', async () => { + it('publishes committed fallback voice writes when a later write fails', async () => { await withIsolatedQwenHome(async () => { const persistSetting = vi.fn( async ( @@ -332,7 +369,16 @@ describe('createDaemonWorkspaceService', () => { cause: expect.objectContaining({ message: 'disk full' }), }); - expect(publishWorkspaceEvent).not.toHaveBeenCalled(); + expect(publishWorkspaceEvent).toHaveBeenCalledOnce(); + expect(publishWorkspaceEvent).toHaveBeenCalledWith({ + type: 'settings_changed', + data: { + key: 'general.voice.mode', + value: 'tap', + scope: 'user', + }, + originatorClientId: 'voice-client', + }); expect(mockWriteStderrLine).toHaveBeenCalledWith( expect.stringContaining('partial persist error'), ); diff --git a/packages/cli/src/serve/workspace-service/index.ts b/packages/cli/src/serve/workspace-service/index.ts index f847b9d5ed2..21782c5b7c9 100644 --- a/packages/cli/src/serve/workspace-service/index.ts +++ b/packages/cli/src/serve/workspace-service/index.ts @@ -207,6 +207,8 @@ export function createDaemonWorkspaceService( persistDisabledSkills, persistSetting, persistSettings, + voiceEnv, + voiceSettingsScope, preheatAcpChild: preheatAcpChildOnBridge, queryWorkspaceStatus, invokeWorkspaceCommand, @@ -484,7 +486,10 @@ export function createDaemonWorkspaceService( async getWorkspaceVoiceStatus(_ctx: WorkspaceRequestContext) { return buildWorkspaceVoiceStatus( boundWorkspace, - loadSettings(boundWorkspace), + loadSettings( + boundWorkspace, + voiceEnv ? { skipLoadEnvironment: true } : true, + ), ); }, @@ -551,13 +556,17 @@ export function createDaemonWorkspaceService( ); } - const settings = loadSettings(boundWorkspace); - validateWorkspaceVoiceState(settings, request); + const settings = loadSettings( + boundWorkspace, + voiceEnv ? { skipLoadEnvironment: true } : true, + ); + validateWorkspaceVoiceState(settings, request, { env: voiceEnv }); const workspaceTrusted = getWorkspaceTrustStatus(settings.merged, boundWorkspace).effective .state === 'trusted'; const writes = buildWorkspaceVoiceSettingsWrites(settings, request, { workspaceTrusted, + ...(voiceSettingsScope ? { scopeOverride: voiceSettingsScope } : {}), }); const publishWrite = (write: WorkspaceVoiceSettingsWrite) => { @@ -602,6 +611,9 @@ export function createDaemonWorkspaceService( err instanceof Error ? err.message : String(err) }`, ); + for (const committedWrite of committed) { + publishWrite(committedWrite); + } throw new WorkspaceSettingsPartialPersistError( `Voice settings partial persist failed: committed=${committed.length}/${writes.length}`, committed, @@ -617,7 +629,10 @@ export function createDaemonWorkspaceService( return buildWorkspaceVoiceStatus( boundWorkspace, - loadSettings(boundWorkspace), + loadSettings( + boundWorkspace, + voiceEnv ? { skipLoadEnvironment: true } : true, + ), ); }, diff --git a/packages/cli/src/serve/workspace-service/types.ts b/packages/cli/src/serve/workspace-service/types.ts index e16ab1fe31b..58e2d36983f 100644 --- a/packages/cli/src/serve/workspace-service/types.ts +++ b/packages/cli/src/serve/workspace-service/types.ts @@ -443,6 +443,12 @@ export interface DaemonWorkspaceServiceDeps { writes: WorkspaceSettingsWrite[], ) => Promise; + /** Runtime-local environment used by workspace Voice operations. */ + voiceEnv?: Readonly>; + + /** Force Voice settings writes into this scope for workspace-qualified ACP. */ + voiceSettingsScope?: SettingScope; + /** Reload daemon-side process.env from .env / settings.env. */ reloadDaemonEnv?: (workspace: string) => Promise; diff --git a/packages/cli/src/services/voice-service.test.ts b/packages/cli/src/services/voice-service.test.ts index c1b5ec433d4..fa64a0d6968 100644 --- a/packages/cli/src/services/voice-service.test.ts +++ b/packages/cli/src/services/voice-service.test.ts @@ -4,7 +4,7 @@ * SPDX-License-Identifier: Apache-2.0 */ -import { describe, expect, it } from 'vitest'; +import { afterEach, describe, expect, it } from 'vitest'; import { LoadedSettings, SettingScope, @@ -56,6 +56,16 @@ function expectWorkspaceVoiceError(action: () => unknown, code: string): void { } describe('voice service', () => { + const originalDashscopeKey = process.env['DASHSCOPE_API_KEY']; + + afterEach(() => { + if (originalDashscopeKey === undefined) { + delete process.env['DASHSCOPE_API_KEY']; + } else { + process.env['DASHSCOPE_API_KEY'] = originalDashscopeKey; + } + }); + it('builds settings writes using voice and model persistence scopes', () => { const settings = makeSettings({ workspace: { @@ -161,6 +171,44 @@ describe('voice service', () => { ]); }); + it('applies a workspace scope override to every voice setting', () => { + const settings = makeSettings({ user: {} }); + + expect( + buildWorkspaceVoiceSettingsWrites( + settings, + { + voiceModel: 'qwen3-asr-flash', + mode: 'tap', + language: 'english', + enabled: true, + }, + { scopeOverride: SettingScope.Workspace }, + ), + ).toEqual([ + { + scope: SettingScope.Workspace, + key: 'voiceModel', + value: 'qwen3-asr-flash', + }, + { + scope: SettingScope.Workspace, + key: 'general.voice.mode', + value: 'tap', + }, + { + scope: SettingScope.Workspace, + key: 'general.voice.language', + value: 'english', + }, + { + scope: SettingScope.Workspace, + key: 'general.voice.enabled', + value: true, + }, + ]); + }); + it('requires an effective voice model before enabling voice', () => { const settings = makeSettings({ user: {} }); @@ -256,6 +304,44 @@ describe('voice service', () => { ); }); + it('uses the supplied runtime environment for validation and transcription', async () => { + process.env['DASHSCOPE_API_KEY'] = 'process-secret'; + const settings = makeSettings({ + user: { + modelProviders: { + openai: [ + { + id: 'qwen3-asr-flash', + baseUrl: 'https://dashscope.aliyuncs.com/compatible-mode/v1', + envKey: 'DASHSCOPE_API_KEY', + }, + ], + }, + }, + }); + const runtimeEnv = { DASHSCOPE_API_KEY: undefined }; + + expectWorkspaceVoiceError( + () => + validateWorkspaceVoiceState( + settings, + { enabled: true, voiceModel: 'qwen3-asr-flash' }, + { env: runtimeEnv }, + ), + 'invalid_voice_model', + ); + await expect( + transcribeWorkspaceVoiceAudio({ + workspaceCwd: '/workspace', + settings, + voiceModel: 'qwen3-asr-flash', + data: new Uint8Array([1, 2, 3]), + mimeType: 'audio/wav', + env: runtimeEnv, + }), + ).rejects.toThrow('requires DASHSCOPE_API_KEY'); + }); + it('rejects realtime-only models for batch daemon transcription', async () => { const settings = makeSettings({ user: { diff --git a/packages/cli/src/services/voice-service.ts b/packages/cli/src/services/voice-service.ts index fb4c901b4c7..5504329244f 100644 --- a/packages/cli/src/services/voice-service.ts +++ b/packages/cli/src/services/voice-service.ts @@ -56,6 +56,7 @@ export interface WorkspaceVoiceTranscriptionInput extends RecordedVoiceAudio { voiceModel: string; settings: LoadedSettings; workspaceCwd: string; + env?: Readonly>; abortSignal?: AbortSignal; } @@ -101,16 +102,15 @@ export function voiceSettingsScopeToWire( export function buildWorkspaceVoiceSettingsWrites( settings: LoadedSettings, update: WorkspaceVoiceStateUpdate, - opts: { workspaceTrusted?: boolean } = {}, + opts: { workspaceTrusted?: boolean; scopeOverride?: SettingScope } = {}, ): WorkspaceVoiceSettingsWrite[] { - const voiceSettingsScope = getVoiceSettingsScope( - settings, - opts.workspaceTrusted, - ); + const voiceSettingsScope = + opts.scopeOverride ?? + getVoiceSettingsScope(settings, opts.workspaceTrusted); const writes: WorkspaceVoiceSettingsWrite[] = []; if (update.voiceModel !== undefined) { writes.push({ - scope: getPersistScopeForModelSelection(settings), + scope: opts.scopeOverride ?? getPersistScopeForModelSelection(settings), key: 'voiceModel', value: update.voiceModel, }); @@ -181,13 +181,14 @@ export function listAvailableVoiceModels( export function hasConfiguredBatchVoiceTranscriptionModel( settings: LoadedSettings, + opts: { env?: Readonly> } = {}, ): boolean { for (const model of listAvailableVoiceModels(settings)) { if (model.transport !== 'qwen-asr-chat') { continue; } try { - validateWorkspaceVoiceConfig(settings, model.id); + validateWorkspaceVoiceConfig(settings, model.id, opts); return true; } catch (err) { debugLogger.debug( @@ -257,14 +258,25 @@ export function validateWorkspaceVoiceModel( export function validateWorkspaceVoiceConfig( settings: LoadedSettings, voiceModel: string, + opts: { env?: Readonly> } = {}, ): WorkspaceVoiceModelDescriptor { const descriptor = validateWorkspaceVoiceModel(settings, voiceModel); try { const config = createVoiceModelSource(settings); if (descriptor.transport === 'qwen-asr-chat') { - resolveVoiceTranscriptionConfig({ config, settings, voiceModel }); + resolveVoiceTranscriptionConfig({ + config, + settings, + voiceModel, + env: opts.env, + }); } else { - resolveVoiceStreamConfig({ config, settings, voiceModel }); + resolveVoiceStreamConfig({ + config, + settings, + voiceModel, + env: opts.env, + }); } } catch (err) { throw new WorkspaceVoiceError( @@ -281,6 +293,7 @@ export function validateWorkspaceVoiceConfig( export function validateWorkspaceVoiceState( settings: LoadedSettings, update: WorkspaceVoiceStateUpdate, + opts: { env?: Readonly> } = {}, ): void { if ( update.enabled === undefined && @@ -309,7 +322,7 @@ export function validateWorkspaceVoiceState( 'A valid voiceModel is required before enabling voice.', ); } - validateWorkspaceVoiceConfig(settings, nextVoiceModel); + validateWorkspaceVoiceConfig(settings, nextVoiceModel, opts); } export async function transcribeWorkspaceVoiceAudio( @@ -332,6 +345,7 @@ export async function transcribeWorkspaceVoiceAudio( config: createVoiceModelSource(input.settings), settings: input.settings, voiceModel: input.voiceModel, + env: input.env, abortSignal: input.abortSignal, }, ); diff --git a/packages/cli/src/services/voice-transcriber.ts b/packages/cli/src/services/voice-transcriber.ts index 853b811c16f..43a289ac128 100644 --- a/packages/cli/src/services/voice-transcriber.ts +++ b/packages/cli/src/services/voice-transcriber.ts @@ -210,6 +210,7 @@ async function defaultLookupHost( export async function assertVoiceBaseUrlNetworkAllowed( voiceConfig: VoiceTranscriptionConfig, lookupHost?: VoiceHostLookup, + abortSignal?: AbortSignal, ): Promise { const hostname = normalizeHostname(new URL(voiceConfig.baseUrl).hostname); if (isLoopbackHost(hostname)) { @@ -224,12 +225,33 @@ export async function assertVoiceBaseUrlNetworkAllowed( return; } let result: { address: string } | Array<{ address: string }>; + let onAbort: (() => void) | undefined; try { - result = await (lookupHost ?? defaultLookupHost)(hostname); + if (abortSignal?.aborted) { + throw abortSignal.reason; + } + const lookup = (lookupHost ?? defaultLookupHost)(hostname); + result = abortSignal + ? await Promise.race([ + lookup, + new Promise((_resolve, reject) => { + onAbort = () => reject(abortSignal.reason); + if (abortSignal.aborted) onAbort(); + else abortSignal.addEventListener('abort', onAbort, { once: true }); + }), + ]) + : await lookup; } catch { + if (abortSignal?.aborted) { + throw abortSignal.reason instanceof Error + ? abortSignal.reason + : new Error('Voice request was aborted.'); + } throw new Error( `Voice model '${voiceConfig.model}': DNS lookup failed for ${hostname}. Cannot verify network safety.`, ); + } finally { + if (onAbort) abortSignal?.removeEventListener('abort', onAbort); } const records = Array.isArray(result) ? result : [result]; if (records.some((record) => isPrivateNetworkIp(record.address))) { @@ -616,7 +638,11 @@ export async function transcribeVoiceAudio( args: TranscribeVoiceAudioArgs, ): Promise { const voiceConfig = resolveVoiceTranscriptionConfig(args); - await assertVoiceBaseUrlNetworkAllowed(voiceConfig, args.lookupHost); + await assertVoiceBaseUrlNetworkAllowed( + voiceConfig, + args.lookupHost, + args.abortSignal, + ); const fetchFn = args.fetchFn ?? fetch; const language = resolveLanguageCode(readVoiceLanguage(args.settings)); const keytermsContext = buildKeytermsContext(args.settings); diff --git a/packages/cli/src/ui/voice/qwen-asr-realtime-session.test.ts b/packages/cli/src/ui/voice/qwen-asr-realtime-session.test.ts index bc3272fb89d..5bb90e3fb1b 100644 --- a/packages/cli/src/ui/voice/qwen-asr-realtime-session.test.ts +++ b/packages/cli/src/ui/voice/qwen-asr-realtime-session.test.ts @@ -69,9 +69,66 @@ describe('qwen-asr-realtime-session', () => { ); }); + it('aborts an upstream connection while it is opening', async () => { + const socket = new FakeSocket(); + const controller = new AbortController(); + const sessionPromise = openQwenAsrRealtimeStream( + { + baseUrl: 'https://dashscope.example/v1', + model: 'qwen3-asr-flash-realtime', + }, + {}, + { + createWebSocket: () => socket, + abortSignal: controller.signal, + }, + ); + + controller.abort(); + + await expect(sessionPromise).rejects.toThrow( + 'Voice stream opening was aborted.', + ); + expect(socket.readyState).toBe(3); + }); + + it('aborts an established upstream connection', async () => { + const socket = new FakeSocket(); + const controller = new AbortController(); + const sessionPromise = openQwenAsrRealtimeStream( + { + baseUrl: 'https://dashscope.example/v1', + model: 'qwen3-asr-flash-realtime', + }, + {}, + { + createWebSocket: () => socket, + abortSignal: controller.signal, + }, + ); + socket.emit( + 'message', + JSON.stringify({ type: 'session.updated', event_id: 'updated' }), + false, + ); + const session = await sessionPromise; + + controller.abort(); + + expect(socket.readyState).toBe(3); + await expect(session.finish()).rejects.toThrow( + 'Voice stream opening was aborted.', + ); + }); + it('streams PCM chunks as base64 events and resolves the completed transcript', async () => { const socket = new FakeSocket(); const createWebSocket = vi.fn(() => socket); + const controller = new AbortController(); + const removeAbortListener = vi.spyOn( + controller.signal, + 'removeEventListener', + ); const interim = vi.fn(); const sessionPromise = openQwenAsrRealtimeStream( { @@ -82,7 +139,7 @@ describe('qwen-asr-realtime-session', () => { keytermsContext: 'grep regex OAuth', }, { onInterim: interim }, - { createWebSocket }, + { createWebSocket, abortSignal: controller.signal }, ); expect(createWebSocket).toHaveBeenCalledWith( @@ -155,6 +212,10 @@ describe('qwen-asr-realtime-session', () => { ); await expect(transcriptPromise).resolves.toBe('hello world'); + expect(removeAbortListener).toHaveBeenCalledWith( + 'abort', + expect.any(Function), + ); }); it('drops realtime audio chunks when the socket buffer is backed up', async () => { diff --git a/packages/cli/src/ui/voice/qwen-asr-realtime-session.ts b/packages/cli/src/ui/voice/qwen-asr-realtime-session.ts index e371d0e2770..3d647bc6675 100644 --- a/packages/cli/src/ui/voice/qwen-asr-realtime-session.ts +++ b/packages/cli/src/ui/voice/qwen-asr-realtime-session.ts @@ -21,6 +21,7 @@ export interface QwenRealtimeDeps { url: string, options: { headers: Record }, ) => SocketLike; + abortSignal?: AbortSignal; } const CONNECT_TIMEOUT_MS = 8000; @@ -61,6 +62,10 @@ export function openQwenAsrRealtimeStream( }) as unknown as SocketLike); return new Promise((resolve, reject) => { + if (deps.abortSignal?.aborted) { + reject(new Error('Voice stream opening was aborted.')); + return; + } const ws = createWebSocket( deriveQwenRealtimeUrl(config.baseUrl, config.model), { @@ -81,6 +86,13 @@ export function openQwenAsrRealtimeStream( let terminalError: Error | null = null; let failed = false; let backpressureWarned = false; + let onAbort: (() => void) | undefined; + + const removeAbortListener = () => { + if (!onAbort) return; + deps.abortSignal?.removeEventListener('abort', onAbort); + onAbort = undefined; + }; const sendJson = (body: Record) => { ws.send(JSON.stringify({ event_id: randomUUID(), ...body })); @@ -111,6 +123,7 @@ export function openQwenAsrRealtimeStream( const fail = (error: unknown) => { if (failed) return; failed = true; + removeAbortListener(); const normalized = toError(error); clearConnectTimer(); clearFinishTimer(); @@ -130,6 +143,15 @@ export function openQwenAsrRealtimeStream( callbacks.onError?.(normalized); }; + if (deps.abortSignal) { + onAbort = () => fail(new Error('Voice stream opening was aborted.')); + if (deps.abortSignal.aborted) { + onAbort(); + return; + } + deps.abortSignal.addEventListener('abort', onAbort, { once: true }); + } + connectTimer = setTimeout(() => { if (!opened) fail(new Error('Qwen ASR realtime connection timed out.')); }, CONNECT_TIMEOUT_MS); @@ -260,6 +282,7 @@ export function openQwenAsrRealtimeStream( break; } failed = true; + removeAbortListener(); clearFinishTimer(); finishedTranscript = committed.trim(); finishResolve?.(finishedTranscript); @@ -285,6 +308,7 @@ export function openQwenAsrRealtimeStream( ws.on('error', fail); ws.on('close', () => { + removeAbortListener(); clearConnectTimer(); clearFinishTimer(); if (failed) return; diff --git a/packages/cli/src/ui/voice/voice-stream-retry.test.ts b/packages/cli/src/ui/voice/voice-stream-retry.test.ts index 0189597ddd0..87d845abe17 100644 --- a/packages/cli/src/ui/voice/voice-stream-retry.test.ts +++ b/packages/cli/src/ui/voice/voice-stream-retry.test.ts @@ -95,4 +95,35 @@ describe('openVoiceStreamWithRetry', () => { await expect(openVoiceStreamWithRetry(open)).rejects.toBe(error); expect(open).toHaveBeenCalledTimes(1); }); + + it('cancels the retry delay when the runtime is disposed', async () => { + vi.useFakeTimers(); + try { + const controller = new AbortController(); + const open = vi + .fn<() => Promise>() + .mockRejectedValueOnce(new Error('early connect failed')) + .mockResolvedValueOnce(session()); + + const result = openVoiceStreamWithRetry(open, { + abortSignal: controller.signal, + }); + const assertion = result.then( + () => { + throw new Error('Expected aborted retry to reject.'); + }, + (error: unknown) => { + expect(error).toEqual(new Error('Voice stream opening was aborted.')); + }, + ); + await vi.waitFor(() => expect(open).toHaveBeenCalledOnce()); + controller.abort(); + + await assertion; + expect(open).toHaveBeenCalledOnce(); + expect(vi.getTimerCount()).toBe(0); + } finally { + vi.useRealTimers(); + } + }); }); diff --git a/packages/cli/src/ui/voice/voice-stream-retry.ts b/packages/cli/src/ui/voice/voice-stream-retry.ts index f73c3baf7e5..e9db1a57787 100644 --- a/packages/cli/src/ui/voice/voice-stream-retry.ts +++ b/packages/cli/src/ui/voice/voice-stream-retry.ts @@ -10,8 +10,22 @@ import { createDebugLogger } from '@qwen-code/qwen-code-core'; const RETRY_DELAY_MS = 200; const debugLogger = createDebugLogger('VOICE_STREAM'); -function delay(ms: number): Promise { - return new Promise((resolve) => setTimeout(resolve, ms)); +function delay(ms: number, abortSignal?: AbortSignal): Promise { + return new Promise((resolve, reject) => { + if (abortSignal?.aborted) { + reject(new Error('Voice stream opening was aborted.')); + return; + } + const timer = setTimeout(() => { + abortSignal?.removeEventListener('abort', onAbort); + resolve(); + }, ms); + const onAbort = () => { + clearTimeout(timer); + reject(new Error('Voice stream opening was aborted.')); + }; + abortSignal?.addEventListener('abort', onAbort, { once: true }); + }); } function isRetryable(error: unknown): boolean { @@ -28,7 +42,11 @@ function isRetryable(error: unknown): boolean { export async function openVoiceStreamWithRetry( open: () => Promise, + opts: { abortSignal?: AbortSignal } = {}, ): Promise { + if (opts.abortSignal?.aborted) { + throw new Error('Voice stream opening was aborted.'); + } try { return await open(); } catch (error) { @@ -36,7 +54,10 @@ export async function openVoiceStreamWithRetry( throw error; } debugLogger.debug('[voice] stream open failed, retrying:', error); - await delay(RETRY_DELAY_MS); + await delay(RETRY_DELAY_MS, opts.abortSignal); + if (opts.abortSignal?.aborted) { + throw new Error('Voice stream opening was aborted.'); + } return open(); } } diff --git a/packages/cli/src/ui/voice/voice-stream-session.test.ts b/packages/cli/src/ui/voice/voice-stream-session.test.ts index c73775acd89..60eab884a71 100644 --- a/packages/cli/src/ui/voice/voice-stream-session.test.ts +++ b/packages/cli/src/ui/voice/voice-stream-session.test.ts @@ -145,9 +145,81 @@ describe('voice-stream-session', () => { ); }); + it('aborts an upstream connection while it is opening', async () => { + const socket = new FakeSocket(); + const controller = new AbortController(); + const sessionPromise = openVoiceStream( + { + baseUrl: 'https://dashscope.example/v1', + model: 'paraformer-realtime-v2', + }, + {}, + { + createWebSocket: () => socket, + abortSignal: controller.signal, + }, + ); + + controller.abort(); + + await expect(sessionPromise).rejects.toThrow( + 'Voice stream opening was aborted.', + ); + expect(socket.readyState).toBe(3); + }); + + it('aborts an established upstream connection', async () => { + const socket = new FakeSocket(); + const controller = new AbortController(); + const sessionPromise = openVoiceStream( + { + baseUrl: 'https://dashscope.example/v1', + model: 'paraformer-realtime-v2', + }, + {}, + { + createWebSocket: () => socket, + abortSignal: controller.signal, + }, + ); + socket.emit('open'); + socket.emit( + 'message', + JSON.stringify({ header: { event: 'task-started' } }), + false, + ); + const session = await sessionPromise; + + controller.abort(); + + expect(socket.readyState).toBe(3); + await expect(session.finish()).rejects.toThrow( + 'Voice stream opening was aborted.', + ); + }); + it('resolves finish when task-finished arrives before finish is called', async () => { const socket = new FakeSocket(); - const session = await startSession(socket); + const controller = new AbortController(); + const removeAbortListener = vi.spyOn( + controller.signal, + 'removeEventListener', + ); + const sessionPromise = openVoiceStream( + { + baseUrl: 'https://dashscope.example/v1', + model: 'paraformer-realtime-v2', + }, + {}, + { createWebSocket: () => socket, abortSignal: controller.signal }, + ); + socket.emit('open'); + socket.emit( + 'message', + JSON.stringify({ header: { event: 'task-started' } }), + false, + ); + const session = await sessionPromise; socket.emit( 'message', @@ -166,6 +238,10 @@ describe('voice-stream-session', () => { ); await expect(session.finish()).resolves.toBe('hello world'); + expect(removeAbortListener).toHaveBeenCalledWith( + 'abort', + expect.any(Function), + ); }); it('drops audio chunks when the socket buffer is backed up', async () => { diff --git a/packages/cli/src/ui/voice/voice-stream-session.ts b/packages/cli/src/ui/voice/voice-stream-session.ts index 0cc43d297ea..323e440714e 100644 --- a/packages/cli/src/ui/voice/voice-stream-session.ts +++ b/packages/cli/src/ui/voice/voice-stream-session.ts @@ -54,6 +54,7 @@ export interface VoiceStreamDeps { url: string, options: { headers: Record }, ) => SocketLike; + abortSignal?: AbortSignal; } const CONNECT_TIMEOUT_MS = 8000; @@ -96,6 +97,10 @@ export function openVoiceStream( }) as unknown as SocketLike); return new Promise((resolve, reject) => { + if (deps.abortSignal?.aborted) { + reject(new Error('Voice stream opening was aborted.')); + return; + } const streamUrl = deriveStreamUrl(config.baseUrl); const ws = createWebSocket(streamUrl, { headers: config.apiKey @@ -114,6 +119,13 @@ export function openVoiceStream( let terminalError: Error | null = null; let finishedTranscript: string | null = null; let backpressureWarned = false; + let onAbort: (() => void) | undefined; + + const removeAbortListener = () => { + if (!onAbort) return; + deps.abortSignal?.removeEventListener('abort', onAbort); + onAbort = undefined; + }; const clearFinishTimer = () => { if (finishTimer) { @@ -132,6 +144,7 @@ export function openVoiceStream( const fail = (error: unknown) => { if (settled) return; settled = true; + removeAbortListener(); const normalized = error instanceof Error ? error : new Error(String(error)); clearConnectTimer(); @@ -155,6 +168,15 @@ export function openVoiceStream( } }; + if (deps.abortSignal) { + onAbort = () => fail(new Error('Voice stream opening was aborted.')); + if (deps.abortSignal.aborted) { + onAbort(); + return; + } + deps.abortSignal.addEventListener('abort', onAbort, { once: true }); + } + connectTimer = setTimeout(() => { if (!started) fail(new Error('Voice stream connection timed out.')); }, CONNECT_TIMEOUT_MS); @@ -294,6 +316,7 @@ export function openVoiceStream( } finishedTranscript = committed.trim(); settled = true; + removeAbortListener(); clearConnectTimer(); clearFinishTimer(); try { @@ -323,6 +346,7 @@ export function openVoiceStream( }); ws.on('close', () => { + removeAbortListener(); clearConnectTimer(); clearFinishTimer(); if (settled) return; diff --git a/packages/cli/src/ui/voice/voice-transcriber.test.ts b/packages/cli/src/ui/voice/voice-transcriber.test.ts index af5644e1589..511d41c3a7f 100644 --- a/packages/cli/src/ui/voice/voice-transcriber.test.ts +++ b/packages/cli/src/ui/voice/voice-transcriber.test.ts @@ -612,6 +612,22 @@ describe('voice-transcriber', () => { ).rejects.toThrow(/DNS lookup failed for asr\.example/); }); + it('aborts a pending DNS safety lookup', async () => { + const controller = new AbortController(); + const check = assertVoiceBaseUrlNetworkAllowed( + { + model: 'qwen3-asr-flash', + baseUrl: 'https://asr.example/v1', + }, + () => new Promise(() => undefined), + controller.signal, + ); + + controller.abort(new Error('Workspace runtime was removed')); + + await expect(check).rejects.toThrow('Workspace runtime was removed'); + }); + it('allows localhost voice URLs for development', () => { const config = createConfig([ { diff --git a/packages/sdk-typescript/src/daemon/DaemonClient.ts b/packages/sdk-typescript/src/daemon/DaemonClient.ts index ec7269b4a0f..906aaf87a98 100644 --- a/packages/sdk-typescript/src/daemon/DaemonClient.ts +++ b/packages/sdk-typescript/src/daemon/DaemonClient.ts @@ -2632,12 +2632,40 @@ export class DaemonClient { async transcribeWorkspaceVoice( audio: DaemonVoiceAudioInput, opts: DaemonWorkspaceVoiceTranscribeOptions, + ): Promise { + return await this.voiceTranscriptionRequest( + '/workspace/voice/transcribe', + 'POST /workspace/voice/transcribe', + audio, + opts, + ); + } + + /** @internal */ + async workspaceVoiceTranscriptionRequest( + workspaceSelector: string, + audio: DaemonVoiceAudioInput, + opts: DaemonWorkspaceVoiceTranscribeOptions, + ): Promise { + return await this.voiceTranscriptionRequest( + `/workspaces/${workspaceSelector}/voice/transcribe`, + 'POST /workspaces/:workspace/voice/transcribe', + audio, + opts, + ); + } + + private async voiceTranscriptionRequest( + path: string, + label: string, + audio: DaemonVoiceAudioInput, + opts: DaemonWorkspaceVoiceTranscribeOptions, ): Promise { const query = opts.voiceModel ? `?${new URLSearchParams({ voiceModel: opts.voiceModel }).toString()}` : ''; return await this.fetchWithTimeout( - `${this.baseUrl}/workspace/voice/transcribe${query}`, + `${this.baseUrl}${path}${query}`, { method: 'POST', headers: this.headers({ 'Content-Type': opts.mimeType }, opts.clientId), @@ -2645,11 +2673,11 @@ export class DaemonClient { }, async (res) => { if (!res.ok) { - throw await this.failOnError(res, 'POST /workspace/voice/transcribe'); + throw await this.failOnError(res, label); } return (await res.json()) as DaemonWorkspaceVoiceTranscriptionResult; }, - VOICE_TRANSCRIPTION_DEFAULT_TIMEOUT_MS, + opts.timeoutMs ?? VOICE_TRANSCRIPTION_DEFAULT_TIMEOUT_MS, 'rest', ); } @@ -3830,6 +3858,38 @@ export class WorkspaceDaemonClient { return this.get('/mcp', 'GET /workspaces/:workspace/mcp'); } + workspaceVoice(clientId?: string): Promise { + return this.client.workspaceJsonRequest( + this.workspaceSelector, + '/voice', + 'GET /workspaces/:workspace/voice', + { clientId, mode: 'rest' }, + ); + } + + setWorkspaceVoice( + update: DaemonWorkspaceVoiceUpdate, + clientId?: string, + ): Promise { + return this.client.workspaceJsonRequest( + this.workspaceSelector, + '/voice', + 'POST /workspaces/:workspace/voice', + { method: 'POST', body: update, clientId, mode: 'rest' }, + ); + } + + transcribeWorkspaceVoice( + audio: DaemonVoiceAudioInput, + opts: DaemonWorkspaceVoiceTranscribeOptions, + ): Promise { + return this.client.workspaceVoiceTranscriptionRequest( + this.workspaceSelector, + audio, + opts, + ); + } + workspaceGit(): Promise { return this.client.workspaceJsonRequest( this.workspaceSelector, diff --git a/packages/sdk-typescript/src/daemon/types.ts b/packages/sdk-typescript/src/daemon/types.ts index 668d357d422..8a2050a9b0d 100644 --- a/packages/sdk-typescript/src/daemon/types.ts +++ b/packages/sdk-typescript/src/daemon/types.ts @@ -42,6 +42,7 @@ export interface DaemonWorkspaceRemovalActivity { acpConnections: number; memoryTasks: number; channelWorkers: number; + voiceSessions?: number; } export interface DaemonWorkspaceRemovalResult { @@ -2086,6 +2087,7 @@ export interface DaemonWorkspaceVoiceTranscribeOptions { mimeType: string; voiceModel?: string; clientId?: string; + timeoutMs?: number; } export interface DaemonWorkspaceVoiceTranscriptionResult { diff --git a/packages/sdk-typescript/test/unit/DaemonClientVoice.test.ts b/packages/sdk-typescript/test/unit/DaemonClientVoice.test.ts index bfef221b80f..ebc3594edf0 100644 --- a/packages/sdk-typescript/test/unit/DaemonClientVoice.test.ts +++ b/packages/sdk-typescript/test/unit/DaemonClientVoice.test.ts @@ -131,6 +131,86 @@ describe('DaemonClient voice helpers', () => { expect(calls[0]?.body).toBe(audio); }); + it('uses the encoded workspace selector for qualified Voice requests', async () => { + const status = { + v: 1, + workspaceCwd: '/work with space', + enabled: true, + mode: 'hold', + language: '', + voiceModel: 'qwen3-asr-flash', + availableVoiceModels: [], + }; + const transcription = { + v: 1, + text: 'hello', + model: 'qwen3-asr-flash', + transport: 'qwen-asr-chat', + }; + const { fetch, calls } = recordingFetch((request) => + request.url.includes('/transcribe') + ? jsonResponse(200, transcription) + : jsonResponse(200, status), + ); + const client = new DaemonClient({ baseUrl: 'http://daemon', fetch }); + const workspace = client.workspaceByCwd('/work with space'); + + await expect(workspace.workspaceVoice('client-1')).resolves.toEqual(status); + await expect( + workspace.setWorkspaceVoice({ enabled: true }, 'client-2'), + ).resolves.toEqual(status); + await expect( + workspace.transcribeWorkspaceVoice(new Uint8Array([1, 2]), { + mimeType: 'audio/wav', + voiceModel: 'qwen3 asr', + clientId: 'client-3', + }), + ).resolves.toEqual(transcription); + + expect(calls.map((call) => call.url)).toEqual([ + 'http://daemon/workspaces/%2Fwork%20with%20space/voice', + 'http://daemon/workspaces/%2Fwork%20with%20space/voice', + 'http://daemon/workspaces/%2Fwork%20with%20space/voice/transcribe?voiceModel=qwen3+asr', + ]); + expect(calls[1]?.headers['x-qwen-client-id']).toBe('client-2'); + expect(calls[2]?.headers['content-type']).toBe('audio/wav'); + expect(calls[2]?.headers['x-qwen-client-id']).toBe('client-3'); + }); + + it('applies a custom timeout to qualified Voice transcription', async () => { + vi.useFakeTimers(); + try { + let requestSignal: AbortSignal | null | undefined; + const fetch = vi.fn( + async (_input: RequestInfo | URL, init?: RequestInit) => { + requestSignal = init?.signal; + return await new Promise((_resolve, reject) => { + init?.signal?.addEventListener( + 'abort', + () => reject(init.signal?.reason), + { once: true }, + ); + }); + }, + ) as unknown as typeof globalThis.fetch; + const client = new DaemonClient({ baseUrl: 'http://daemon', fetch }); + const result = client + .workspaceById('secondary-id') + .transcribeWorkspaceVoice(new Uint8Array([1]), { + mimeType: 'audio/wav', + timeoutMs: 25, + }); + const settled = result.catch(() => undefined); + + await vi.advanceTimersByTimeAsync(25); + await settled; + + expect(requestSignal?.aborted).toBe(true); + } finally { + vi.useRealTimers(); + } + }); + it('uses the REST endpoint for binary voice audio when ACP transport is configured', async () => { const response = { v: 1, @@ -166,6 +246,62 @@ describe('DaemonClient voice helpers', () => { expect(calls[0]?.method).toBe('POST'); }); + it('uses REST for every qualified Voice method when ACP transport is configured', async () => { + const status = { + v: 1, + workspaceCwd: '/secondary', + enabled: false, + mode: 'hold', + language: '', + voiceModel: null, + availableVoiceModels: [], + }; + const transcription = { + v: 1, + text: 'hello', + model: 'qwen3-asr-flash', + transport: 'qwen-asr-chat', + }; + const { fetch, calls } = recordingFetch((request) => + jsonResponse( + 200, + request.url.endsWith('/transcribe') ? transcription : status, + ), + ); + const acpTransport: DaemonTransport = { + type: 'acp-http', + supportsReplay: false, + connected: true, + fetch: vi.fn(async () => { + throw new Error('qualified Voice must not use ACP route mapping'); + }), + subscribeEvents: vi.fn(() => emptyAsyncEvents()), + dispose: vi.fn(), + }; + const workspace = new DaemonClient({ + baseUrl: 'http://daemon', + fetch, + transport: acpTransport, + }).workspaceById('secondary-id'); + + await expect(workspace.workspaceVoice()).resolves.toEqual(status); + await expect( + workspace.setWorkspaceVoice({ enabled: false }), + ).resolves.toEqual(status); + await expect( + workspace.transcribeWorkspaceVoice(new Uint8Array([1]), { + mimeType: 'audio/wav', + }), + ).resolves.toEqual(transcription); + + expect(acpTransport.fetch).not.toHaveBeenCalled(); + expect(calls.map((call) => `${call.method} ${call.url}`)).toEqual([ + 'GET http://daemon/workspaces/secondary-id/voice', + 'POST http://daemon/workspaces/secondary-id/voice', + 'POST http://daemon/workspaces/secondary-id/voice/transcribe', + ]); + }); + it('allows voice transcription to run longer than the client default timeout', async () => { const response = { v: 1, diff --git a/packages/web-shell/client/components/sidebar/WebShellSidebar.tsx b/packages/web-shell/client/components/sidebar/WebShellSidebar.tsx index b8d2fecd257..c251b8e4ef3 100644 --- a/packages/web-shell/client/components/sidebar/WebShellSidebar.tsx +++ b/packages/web-shell/client/components/sidebar/WebShellSidebar.tsx @@ -2753,6 +2753,11 @@ export function WebShellSidebar({ count: workspaceRemovalActivity.channelWorkers, })} +
  • + {t('sidebar.removeWorkspaceVoiceSessions', { + count: workspaceRemovalActivity.voiceSessions ?? 0, + })} +
  • )} {workspaceRemovalActivity && diff --git a/packages/web-shell/client/components/sidebar/WebShellSidebar.workspace-removal.test.tsx b/packages/web-shell/client/components/sidebar/WebShellSidebar.workspace-removal.test.tsx index 6647faeb41c..5fd57b79400 100644 --- a/packages/web-shell/client/components/sidebar/WebShellSidebar.workspace-removal.test.tsx +++ b/packages/web-shell/client/components/sidebar/WebShellSidebar.workspace-removal.test.tsx @@ -275,4 +275,32 @@ describe('WebShellSidebar workspace removal', () => { ); expect(dialogButton('Force remove').disabled).toBe(true); }); + + it('shows Voice-only activity before offering force removal', async () => { + workspaceActions.removeWorkspace.mockRejectedValueOnce( + new DaemonHttpError( + 409, + { + code: 'workspace_busy', + activity: { + sessions: 0, + activePrompts: 0, + pendingSessionStarts: 0, + acpConnections: 0, + memoryTasks: 0, + channelWorkers: 0, + voiceSessions: 1, + }, + }, + 'busy', + ), + ); + renderSidebar(); + openRemoval('/tmp/other'); + + await act(async () => click(dialogButton('Remove workspace'))); + + expect(document.body.textContent).toContain('Voice sessions: 1'); + expect(dialogButton('Force remove').disabled).toBe(false); + }); }); diff --git a/packages/web-shell/client/i18n.tsx b/packages/web-shell/client/i18n.tsx index 440ea454e1f..5c6efeaab22 100644 --- a/packages/web-shell/client/i18n.tsx +++ b/packages/web-shell/client/i18n.tsx @@ -875,6 +875,8 @@ const EN: Messages = { `ACP connections: ${v?.count ?? 0}`, 'sidebar.removeWorkspaceMemoryTasks': (v) => `Memory tasks: ${v?.count ?? 0}`, 'sidebar.removeWorkspaceWorkers': (v) => `Channel workers: ${v?.count ?? 0}`, + 'sidebar.removeWorkspaceVoiceSessions': (v) => + `Voice sessions: ${v?.count ?? 0}`, 'sidebar.noSessions': 'No sessions.', 'sidebar.projectFallback': 'Project', 'sidebar.sessionsOverview': 'Session Overview', @@ -2781,6 +2783,7 @@ const ZH: Messages = { 'sidebar.removeWorkspaceConnections': (v) => `ACP 连接:${v?.count ?? 0}`, 'sidebar.removeWorkspaceMemoryTasks': (v) => `Memory 任务:${v?.count ?? 0}`, 'sidebar.removeWorkspaceWorkers': (v) => `Channel worker:${v?.count ?? 0}`, + 'sidebar.removeWorkspaceVoiceSessions': (v) => `语音会话:${v?.count ?? 0}`, 'sidebar.noSessions': '暂无会话', 'sidebar.projectFallback': '项目', 'sidebar.sessionsOverview': '会话总览',