From 5ebb9abae69516ccb4ef8addb96520c2b12c21e6 Mon Sep 17 00:00:00 2001 From: Bob Lee Date: Thu, 27 Aug 2026 19:38:55 -0700 Subject: [PATCH 1/5] feat(voice): add client-level realtime assistant --- scripts/check-core-boundaries.test.mjs | 6 +- .../cargo-dependency-boundaries.mjs | 2 + .../core-boundaries/rules/feature-rules.mjs | 3 +- src/apps/cli/src/peer_host/deny.rs | 18 + src/apps/desktop/Cargo.toml | 2 +- src/apps/desktop/src/api/peer_host_invoke.rs | 18 + .../src/api/remote_workspace_policy.rs | 36 + src/apps/desktop/src/api/speech_api.rs | 201 +++- src/apps/desktop/src/lib.rs | 9 + .../assembly/core/src/service/config/types.rs | 52 + src/crates/contracts/core-types/src/speech.rs | 138 +++ src/crates/contracts/events/src/lib.rs | 4 +- src/crates/contracts/events/src/speech.rs | 1 + .../services/services-integrations/Cargo.toml | 7 + .../services-integrations/src/speech/mod.rs | 8 + .../src/speech/realtime.rs | 1019 +++++++++++++++++ src/web-ui/src/app/App.tsx | 7 +- .../src/app/components/NavPanel/MainNav.tsx | 27 +- .../NavPanel/unifiedSessionCreation.test.ts | 16 +- .../src/flow_chat/components/ChatInput.scss | 24 +- .../src/flow_chat/components/ChatInput.tsx | 6 +- .../voice/RealtimeVoiceCall.appearance.ts | 41 + .../components/voice/RealtimeVoiceCall.scss | 278 +++++ .../voice/RealtimeVoiceCallButton.tsx | 191 +++ .../voice/RealtimeVoiceCallContext.tsx | 35 + .../voice/realtimeVoiceAudio.test.ts | 113 ++ .../components/voice/realtimeVoiceAudio.ts | 143 +++ .../voice/realtimeVoiceTranscript.test.ts | 16 + .../voice/realtimeVoiceTranscript.ts | 16 + .../components/voice/useRealtimeVoiceCall.ts | 923 +++++++++++++++ .../components/voice/voiceClientContext.ts | 172 +++ .../components/voice/voiceTaskBridge.test.ts | 128 +++ .../components/voice/voiceTaskBridge.ts | 358 ++++++ .../api/adapters/peer-device-adapter.test.ts | 6 + .../api/adapters/peer-device-adapter.ts | 9 + .../api/service-api/SpeechAPI.ts | 132 +++ .../registry/defaultAppearanceRegistry.ts | 2 + .../config/components/VoiceInputConfig.scss | 4 + .../config/components/VoiceInputConfig.tsx | 122 +- .../src/infrastructure/config/types/index.ts | 12 + .../infrastructure/speech/voiceInputAudio.ts | 22 +- .../locales/en-US/settings/voice-input.json | 67 ++ .../locales/zh-CN/settings/voice-input.json | 67 ++ .../locales/zh-TW/settings/voice-input.json | 67 ++ 44 files changed, 4479 insertions(+), 49 deletions(-) create mode 100644 src/crates/services/services-integrations/src/speech/realtime.rs create mode 100644 src/web-ui/src/flow_chat/components/voice/RealtimeVoiceCall.appearance.ts create mode 100644 src/web-ui/src/flow_chat/components/voice/RealtimeVoiceCall.scss create mode 100644 src/web-ui/src/flow_chat/components/voice/RealtimeVoiceCallButton.tsx create mode 100644 src/web-ui/src/flow_chat/components/voice/RealtimeVoiceCallContext.tsx create mode 100644 src/web-ui/src/flow_chat/components/voice/realtimeVoiceAudio.test.ts create mode 100644 src/web-ui/src/flow_chat/components/voice/realtimeVoiceAudio.ts create mode 100644 src/web-ui/src/flow_chat/components/voice/realtimeVoiceTranscript.test.ts create mode 100644 src/web-ui/src/flow_chat/components/voice/realtimeVoiceTranscript.ts create mode 100644 src/web-ui/src/flow_chat/components/voice/useRealtimeVoiceCall.ts create mode 100644 src/web-ui/src/flow_chat/components/voice/voiceClientContext.ts create mode 100644 src/web-ui/src/flow_chat/components/voice/voiceTaskBridge.test.ts create mode 100644 src/web-ui/src/flow_chat/components/voice/voiceTaskBridge.ts diff --git a/scripts/check-core-boundaries.test.mjs b/scripts/check-core-boundaries.test.mjs index 6045bdc4bd..c7b74ccd78 100644 --- a/scripts/check-core-boundaries.test.mjs +++ b/scripts/check-core-boundaries.test.mjs @@ -2993,7 +2993,7 @@ test('services integrations image codecs stay attached to exact product owners', assert.match(messages, /image shared Image activation alias must not select capabilities: png/); }); -test('services integrations WebSocket TLS stays attached to remote-connect', () => { +test('services integrations WebSocket TLS stays attached to reviewed realtime owners', () => { const pkg = { ...packageAt('bitfun-services-integrations', 'src/crates/services/services-integrations/Cargo.toml', [{ name: 'tokio-tungstenite', @@ -3007,6 +3007,10 @@ test('services integrations WebSocket TLS stays attached to remote-connect', () 'dep:tokio-tungstenite', 'tokio-tungstenite?/rustls-tls-native-roots', ], + 'speech-realtime': [ + 'dep:tokio-tungstenite', + 'tokio-tungstenite?/rustls-tls-native-roots', + ], }, }; diff --git a/scripts/core-boundaries/cargo-dependency-boundaries.mjs b/scripts/core-boundaries/cargo-dependency-boundaries.mjs index f6cbae069e..88d71b1d73 100644 --- a/scripts/core-boundaries/cargo-dependency-boundaries.mjs +++ b/scripts/core-boundaries/cargo-dependency-boundaries.mjs @@ -150,6 +150,7 @@ const SERVICES_INTEGRATIONS_TOKIO_FEATURES = new Map([ ['remote-ssh-concrete', ['fs', 'io-util', 'macros', 'net', 'process', 'rt', 'sync', 'time']], ['review-platform', ['fs', 'io-util', 'sync']], ['speech', ['fs', 'io-util', 'macros', 'rt', 'sync']], + ['speech-realtime', ['fs', 'io-util', 'macros', 'net', 'rt', 'sync', 'time']], ['workspace-search', ['io-util', 'rt', 'sync', 'time']], ['script-tool-runtime', ['io-util', 'process', 'rt', 'sync', 'time']], ]); @@ -560,6 +561,7 @@ const THIRD_PARTY_CAPABILITY_PROFILES = new Map([ optional: true, ownerFeatureCapabilities: new Map([ ['remote-connect', ['rustls-tls-native-roots']], + ['speech-realtime', ['rustls-tls-native-roots']], ]), })], ]), diff --git a/scripts/core-boundaries/rules/feature-rules.mjs b/scripts/core-boundaries/rules/feature-rules.mjs index f4720269c2..8f426491f6 100644 --- a/scripts/core-boundaries/rules/feature-rules.mjs +++ b/scripts/core-boundaries/rules/feature-rules.mjs @@ -252,6 +252,7 @@ export const optionalDependencyFeatureOwnerRules = [ crateName: 'services-integrations', reason: 'services-integrations optional runtime dependencies must stay owned by explicit integration features', + reviewedAggregateFeatures: ['speech-realtime'], dependencies: [ { depName: 'aes', ownerFeatures: ['remote-connect'] }, { depName: 'aes-gcm', ownerFeatures: ['mcp', 'remote-connect', 'remote-ssh-concrete'] }, @@ -328,7 +329,7 @@ export const optionalDependencyFeatureOwnerRules = [ { depName: 'terminal-core', ownerFeatures: ['remote-ssh', 'remote-ssh-concrete'] }, { depName: 'tar', ownerFeatures: ['speech'] }, { depName: 'thiserror', ownerFeatures: ['browser-control', 'git', 'hook-import', 'miniapp-market', 'plugin-source', 'remote-ssh', 'remote-ssh-concrete', 'review-platform', 'speech', 'web-tools', 'workspace-search'] }, - { depName: 'tokio-tungstenite', ownerFeatures: ['remote-connect'] }, + { depName: 'tokio-tungstenite', ownerFeatures: ['remote-connect', 'speech-realtime'] }, { depName: 'tokio-util', ownerFeatures: ['remote-ssh', 'speech'] }, { depName: 'urlencoding', ownerFeatures: ['canvas-runtime', 'miniapp-market', 'remote-connect', 'review-platform'] }, { depName: 'uuid', ownerFeatures: ['canvas-runtime', 'hook-import', 'miniapp-runtime', 'plugin-source', 'remote-connect', 'remote-ssh-concrete', 'speech'] }, diff --git a/src/apps/cli/src/peer_host/deny.rs b/src/apps/cli/src/peer_host/deny.rs index 8b0fcfb067..eaf4645395 100644 --- a/src/apps/cli/src/peer_host/deny.rs +++ b/src/apps/cli/src/peer_host/deny.rs @@ -114,6 +114,15 @@ static LOCAL_ONLY_COMMANDS: &[&str] = &[ "speech_append_audio_chunk", "speech_finish_input_session", "speech_cancel_input_session", + "speech_start_realtime_session", + "speech_append_realtime_audio", + "speech_commit_realtime_audio", + "speech_send_realtime_tool_result", + "speech_speak_realtime_text", + "speech_cancel_realtime_response", + "speech_close_realtime_session", + "speech_get_realtime_config", + "speech_save_realtime_config", // Granting Git ownership trust writes this user's global Git configuration // and tells Git to run hooks from a tree they do not own. That decision // belongs to the person at this machine, so refuse it explicitly rather @@ -266,6 +275,15 @@ mod tests { "speech_append_audio_chunk", "speech_finish_input_session", "speech_cancel_input_session", + "speech_start_realtime_session", + "speech_append_realtime_audio", + "speech_commit_realtime_audio", + "speech_send_realtime_tool_result", + "speech_speak_realtime_text", + "speech_cancel_realtime_response", + "speech_close_realtime_session", + "speech_get_realtime_config", + "speech_save_realtime_config", ] { assert!(is_local_only_command(command), "{command}"); } diff --git a/src/apps/desktop/Cargo.toml b/src/apps/desktop/Cargo.toml index 6b44ddf8f3..52bcdfb676 100644 --- a/src/apps/desktop/Cargo.toml +++ b/src/apps/desktop/Cargo.toml @@ -24,7 +24,7 @@ bitfun-relay-service = { path = "../../crates/services/relay-service" } bitfun-agent-runtime = { path = "../../crates/execution/agent-runtime", features = ["agent-runtime"] } bitfun-runtime-ports = { path = "../../crates/contracts/runtime-ports", features = ["agent-api", "permission", "workspace-ports"] } bitfun-product-domains = { path = "../../crates/contracts/product-domains", features = ["appearance-market"] } -bitfun-services-integrations = { path = "../../crates/services/services-integrations", features = ["canvas-runtime", "miniapp-market", "speech"] } +bitfun-services-integrations = { path = "../../crates/services/services-integrations", features = ["canvas-runtime", "miniapp-market", "speech-realtime"] } bitfun-core-types = { path = "../../crates/contracts/core-types" } bitfun-agent-tools = { path = "../../crates/execution/tool-contracts", features = ["element-token"] } bitfun-transport = { path = "../../crates/adapters/transport", features = ["tauri-adapter"] } diff --git a/src/apps/desktop/src/api/peer_host_invoke.rs b/src/apps/desktop/src/api/peer_host_invoke.rs index 0247b7f8d3..97a87dabde 100644 --- a/src/apps/desktop/src/api/peer_host_invoke.rs +++ b/src/apps/desktop/src/api/peer_host_invoke.rs @@ -151,6 +151,15 @@ static LOCAL_ONLY_COMMANDS: &[&str] = &[ "speech_append_audio_chunk", "speech_finish_input_session", "speech_cancel_input_session", + "speech_start_realtime_session", + "speech_append_realtime_audio", + "speech_commit_realtime_audio", + "speech_send_realtime_tool_result", + "speech_speak_realtime_text", + "speech_cancel_realtime_response", + "speech_close_realtime_session", + "speech_get_realtime_config", + "speech_save_realtime_config", // Granting Git ownership trust writes to the peer user's global Git // configuration and tells Git to run hooks from a tree they do not own. // That decision stays with the person at that machine; a controller can @@ -655,6 +664,15 @@ mod tests { "speech_append_audio_chunk", "speech_finish_input_session", "speech_cancel_input_session", + "speech_start_realtime_session", + "speech_append_realtime_audio", + "speech_commit_realtime_audio", + "speech_send_realtime_tool_result", + "speech_speak_realtime_text", + "speech_cancel_realtime_response", + "speech_close_realtime_session", + "speech_get_realtime_config", + "speech_save_realtime_config", // Same controller-owned observer/credential family as the other // dispatch verbs already denied here. "dispatch_continue", diff --git a/src/apps/desktop/src/api/remote_workspace_policy.rs b/src/apps/desktop/src/api/remote_workspace_policy.rs index dc18ce8432..73fdfe1a28 100644 --- a/src/apps/desktop/src/api/remote_workspace_policy.rs +++ b/src/apps/desktop/src/api/remote_workspace_policy.rs @@ -1704,6 +1704,14 @@ pub const REMOTE_WORKSPACE_COMMAND_POLICIES: &[(&str, RemoteWorkspacePolicy)] = "speech_append_audio_chunk", RemoteWorkspacePolicy::LocalOnly, ), + ( + "speech_append_realtime_audio", + RemoteWorkspacePolicy::LocalOnly, + ), + ( + "speech_cancel_realtime_response", + RemoteWorkspacePolicy::LocalOnly, + ), ( "speech_cancel_input_session", RemoteWorkspacePolicy::LocalOnly, @@ -1719,10 +1727,38 @@ pub const REMOTE_WORKSPACE_COMMAND_POLICIES: &[(&str, RemoteWorkspacePolicy)] = RemoteWorkspacePolicy::LocalOnly, ), ("speech_list_models", RemoteWorkspacePolicy::LocalOnly), + ( + "speech_get_realtime_config", + RemoteWorkspacePolicy::LocalOnly, + ), + ( + "speech_save_realtime_config", + RemoteWorkspacePolicy::LocalOnly, + ), + ( + "speech_close_realtime_session", + RemoteWorkspacePolicy::LocalOnly, + ), + ( + "speech_commit_realtime_audio", + RemoteWorkspacePolicy::LocalOnly, + ), + ( + "speech_send_realtime_tool_result", + RemoteWorkspacePolicy::LocalOnly, + ), + ( + "speech_speak_realtime_text", + RemoteWorkspacePolicy::LocalOnly, + ), ( "speech_start_input_session", RemoteWorkspacePolicy::LocalOnly, ), + ( + "speech_start_realtime_session", + RemoteWorkspacePolicy::LocalOnly, + ), ("speech_verify_model", RemoteWorkspacePolicy::LocalOnly), ("ssh_connect", RemoteWorkspacePolicy::WorkspaceAgnostic), ( diff --git a/src/apps/desktop/src/api/speech_api.rs b/src/apps/desktop/src/api/speech_api.rs index 9edc88ff9c..63122faab2 100644 --- a/src/apps/desktop/src/api/speech_api.rs +++ b/src/apps/desktop/src/api/speech_api.rs @@ -2,13 +2,20 @@ use crate::api::AppState; use bitfun_core_types::speech::{ - SpeechAppendAudioChunkRequest, SpeechAppendAudioChunkResponse, SpeechCancelInputSessionRequest, + SpeechAppendAudioChunkRequest, SpeechAppendAudioChunkResponse, + SpeechAppendRealtimeAudioRequest, SpeechCancelInputSessionRequest, SpeechCancelModelDownloadRequest, SpeechDeleteModelRequest, SpeechDownloadModelRequest, - SpeechFinishInputSessionRequest, SpeechInputSession, SpeechListModelsResponse, - SpeechModelProgressEvent, SpeechModelStatus, SpeechStartInputSessionRequest, - SpeechTranscriptionResult, SpeechVerifyModelRequest, + SpeechFinishInputSessionRequest, SpeechGetRealtimeConfigRequest, SpeechInputSession, + SpeechListModelsResponse, SpeechModelProgressEvent, SpeechModelStatus, SpeechRealtimeConfig, + SpeechRealtimeEvent, SpeechRealtimeSession, SpeechRealtimeSessionRequest, + SpeechRealtimeSpeakRequest, SpeechRealtimeToolResultRequest, SpeechSaveRealtimeConfigRequest, + SpeechStartInputSessionRequest, SpeechStartRealtimeSessionRequest, SpeechTranscriptionResult, + SpeechVerifyModelRequest, }; -use bitfun_events::{SPEECH_MODEL_PROGRESS_EVENT, SPEECH_MODEL_STATUS_CHANGED_EVENT}; +use bitfun_events::{ + SPEECH_MODEL_PROGRESS_EVENT, SPEECH_MODEL_STATUS_CHANGED_EVENT, SPEECH_REALTIME_EVENT, +}; +use bitfun_services_integrations::speech::VolcengineRealtimeSpeechConfig; use tauri::{AppHandle, Emitter, State}; #[tauri::command] @@ -135,6 +142,190 @@ pub async fn speech_cancel_input_session( .map_err(|error| format!("Failed to cancel speech input session: {error}")) } +#[tauri::command] +pub async fn speech_start_realtime_session( + state: State<'_, AppState>, + app: AppHandle, + request: SpeechStartRealtimeSessionRequest, +) -> Result { + let global_config: bitfun_core::service::config::GlobalConfig = state + .config_service + .get_config(None) + .await + .map_err(|error| format!("Failed to load realtime voice configuration: {error}"))?; + let voice_call = global_config.app.voice_call; + if !voice_call.enabled { + return Err("Realtime voice calling is disabled in settings".to_string()); + } + if voice_call.provider != "volcengine" { + return Err(format!( + "Unsupported realtime voice provider: {}", + voice_call.provider + )); + } + + let event_app = app.clone(); + state + .speech_service + .start_realtime_session( + VolcengineRealtimeSpeechConfig { + api_key: voice_call.api_key, + voice: voice_call.voice, + speed: voice_call.speed, + loudness: voice_call.loudness, + client_context: request.client_context, + }, + move |event: SpeechRealtimeEvent| { + if let Err(error) = event_app.emit(SPEECH_REALTIME_EVENT, &event) { + log::warn!("Failed to emit realtime speech event: {error}"); + } + }, + ) + .await + .map_err(|error| format!("Failed to start realtime speech session: {error}")) +} + +#[tauri::command] +pub async fn speech_append_realtime_audio( + state: State<'_, AppState>, + request: SpeechAppendRealtimeAudioRequest, +) -> Result<(), String> { + state + .speech_service + .append_realtime_audio(request) + .await + .map_err(|error| format!("Failed to append realtime speech audio: {error}")) +} + +#[tauri::command] +pub async fn speech_commit_realtime_audio( + state: State<'_, AppState>, + request: SpeechRealtimeSessionRequest, +) -> Result<(), String> { + state + .speech_service + .commit_realtime_audio(request) + .await + .map_err(|error| format!("Failed to commit realtime speech audio: {error}")) +} + +#[tauri::command] +pub async fn speech_send_realtime_tool_result( + state: State<'_, AppState>, + request: SpeechRealtimeToolResultRequest, +) -> Result<(), String> { + state + .speech_service + .send_realtime_tool_result(request) + .await + .map_err(|error| format!("Failed to send realtime speech tool result: {error}")) +} + +#[tauri::command] +pub async fn speech_speak_realtime_text( + state: State<'_, AppState>, + request: SpeechRealtimeSpeakRequest, +) -> Result<(), String> { + state + .speech_service + .speak_realtime_text(request) + .await + .map_err(|error| format!("Failed to speak realtime progress: {error}")) +} + +#[tauri::command] +pub async fn speech_cancel_realtime_response( + state: State<'_, AppState>, + request: SpeechRealtimeSessionRequest, +) -> Result<(), String> { + state + .speech_service + .cancel_realtime_response(request) + .await + .map_err(|error| format!("Failed to cancel realtime speech response: {error}")) +} + +#[tauri::command] +pub async fn speech_close_realtime_session( + state: State<'_, AppState>, + request: SpeechRealtimeSessionRequest, +) -> Result<(), String> { + state + .speech_service + .close_realtime_session(request) + .await + .map_err(|error| format!("Failed to close realtime speech session: {error}")) +} + +#[tauri::command] +pub async fn speech_get_realtime_config( + state: State<'_, AppState>, + request: SpeechGetRealtimeConfigRequest, +) -> Result { + let _ = request; + let global_config: bitfun_core::service::config::GlobalConfig = state + .config_service + .get_config(None) + .await + .map_err(|error| format!("Failed to load controller realtime voice settings: {error}"))?; + let config = global_config.app.voice_call; + Ok(SpeechRealtimeConfig { + enabled: config.enabled, + provider: config.provider, + api_key: config.api_key, + voice: config.voice, + speed: config.speed, + loudness: config.loudness, + microphone_device_id: config.microphone_device_id, + }) +} + +#[tauri::command] +pub async fn speech_save_realtime_config( + state: State<'_, AppState>, + request: SpeechSaveRealtimeConfigRequest, +) -> Result { + let api_key = request.api_key.trim().to_string(); + let voice = request.voice.trim().to_string(); + if request.enabled && api_key.is_empty() { + return Err( + "Volcengine realtime speech API key is required when voice calling is enabled" + .to_string(), + ); + } + if voice.is_empty() { + return Err("Volcengine realtime speech voice cannot be empty".to_string()); + } + if !(-50..=100).contains(&request.speed) || !(-50..=100).contains(&request.loudness) { + return Err("Realtime speech speed and loudness must be between -50 and 100".to_string()); + } + let config = bitfun_core::service::config::VoiceCallConfig { + enabled: request.enabled, + provider: "volcengine".to_string(), + api_key, + voice, + speed: request.speed, + loudness: request.loudness, + microphone_device_id: request.microphone_device_id, + }; + state + .config_service + .set_config("app.voice_call", &config) + .await + .map_err(|error| format!("Failed to save controller realtime voice settings: {error}"))?; + crate::api::remote_connect_api::notify_settings_changed(); + + Ok(SpeechRealtimeConfig { + enabled: config.enabled, + provider: config.provider, + api_key: config.api_key, + voice: config.voice, + speed: config.speed, + loudness: config.loudness, + microphone_device_id: config.microphone_device_id, + }) +} + fn emit_status(app: &AppHandle, status: &SpeechModelStatus) { if let Err(error) = app.emit(SPEECH_MODEL_STATUS_CHANGED_EVENT, status) { log::warn!("Failed to emit speech model status event: {error}"); diff --git a/src/apps/desktop/src/lib.rs b/src/apps/desktop/src/lib.rs index 8be1547450..7aec71c6ed 100644 --- a/src/apps/desktop/src/lib.rs +++ b/src/apps/desktop/src/lib.rs @@ -1583,6 +1583,15 @@ pub async fn run() { speech_append_audio_chunk, speech_finish_input_session, speech_cancel_input_session, + speech_start_realtime_session, + speech_append_realtime_audio, + speech_commit_realtime_audio, + speech_send_realtime_tool_result, + speech_speak_realtime_text, + speech_cancel_realtime_response, + speech_close_realtime_session, + speech_get_realtime_config, + speech_save_realtime_config, get_agent_profile_configs, get_agent_profile_config, set_agent_profile_config, diff --git a/src/crates/assembly/core/src/service/config/types.rs b/src/crates/assembly/core/src/service/config/types.rs index fdd736b6a5..204d721db1 100644 --- a/src/crates/assembly/core/src/service/config/types.rs +++ b/src/crates/assembly/core/src/service/config/types.rs @@ -167,6 +167,13 @@ pub struct AppConfig { #[serde(default)] pub flow_chat: AppFlowChatConfig, pub ai_experience: AIExperienceConfig, + /// Controller-owned end-to-end realtime voice conversation settings. + /// + /// This is deliberately a sibling of `ai_experience`: generic AI + /// experience mutations can be routed to a peer host, while the dedicated + /// realtime speech commands always resolve this value on the controller. + #[serde(default)] + pub voice_call: VoiceCallConfig, /// User-defined keyboard shortcut overrides. /// Stored as opaque JSON so the backend remains schema-agnostic; /// the frontend owns the versioned format (StoredKeybindingsV1). @@ -381,6 +388,39 @@ impl Default for VoiceInputConfig { } } +/// Controller-local full-duplex voice-call preferences. +/// +/// The API key crosses only the controller-local settings IPC; starting a +/// session resolves it in the Desktop adapter, so it is never forwarded to a +/// remote workspace or peer HostInvoke request. Session execution selected by +/// a voice tool call can still target a remote workspace or peer through the +/// normal Agent Runtime adapters. +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(default)] +pub struct VoiceCallConfig { + pub enabled: bool, + pub provider: String, + pub api_key: String, + pub voice: String, + pub speed: i32, + pub loudness: i32, + pub microphone_device_id: String, +} + +impl Default for VoiceCallConfig { + fn default() -> Self { + Self { + enabled: true, + provider: "volcengine".to_string(), + api_key: String::new(), + voice: "zh_female_vv_jupiter_bigtts".to_string(), + speed: 0, + loudness: 0, + microphone_device_id: String::new(), + } + } +} + /// Domain request for atomically saving a cloud speech-recognition model and /// selecting it for voice input. Text-generation fields are intentionally not /// part of this contract. @@ -1629,6 +1669,7 @@ impl Default for AppConfig { }, flow_chat: AppFlowChatConfig::default(), ai_experience: AIExperienceConfig::default(), + voice_call: VoiceCallConfig::default(), keybindings: None, user_tool_groups: UserToolGroupsConfig::default(), user_skill_groups: UserSkillGroupsConfig::default(), @@ -2047,6 +2088,17 @@ mod tests { }; use bitfun_runtime_ports::ToolPermissionConfig; + #[test] + fn legacy_app_config_defaults_realtime_voice_call() { + let config: AppConfig = serde_json::from_value(serde_json::json!({})) + .expect("legacy app config should deserialize"); + + assert!(config.voice_call.enabled); + assert_eq!(config.voice_call.provider, "volcengine"); + assert!(config.voice_call.api_key.is_empty()); + assert_eq!(config.voice_call.voice, "zh_female_vv_jupiter_bigtts"); + } + #[test] fn prevent_sleep_defaults_to_disabled() { assert!(!AppConfig::default().prevent_sleep); diff --git a/src/crates/contracts/core-types/src/speech.rs b/src/crates/contracts/core-types/src/speech.rs index 1c08c8e6d8..401babf95b 100644 --- a/src/crates/contracts/core-types/src/speech.rs +++ b/src/crates/contracts/core-types/src/speech.rs @@ -134,3 +134,141 @@ pub struct SpeechTranscriptionResult { pub duration_ms: u64, pub audio_duration_seconds: f64, } + +/// Request to start a controller-local full-duplex voice conversation. +/// +/// Provider credentials and model policy are intentionally not part of this +/// frontend command contract. The Desktop adapter resolves them from the +/// controller's persisted configuration before calling the integration +/// service. +#[derive(Debug, Clone, Default, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct SpeechStartRealtimeSessionRequest { + /// Compact controller-side context captured when the client-level voice + /// assistant starts. Older clients omit it; the realtime model can refresh + /// the snapshot later through its client-context tool. + #[serde(default)] + pub client_context: Option, +} + +#[derive(Debug, Clone, Default, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct SpeechGetRealtimeConfigRequest {} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct SpeechRealtimeConfig { + pub enabled: bool, + pub provider: String, + pub api_key: String, + pub voice: String, + pub speed: i32, + pub loudness: i32, + pub microphone_device_id: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct SpeechSaveRealtimeConfigRequest { + pub enabled: bool, + pub api_key: String, + pub voice: String, + pub speed: i32, + pub loudness: i32, + pub microphone_device_id: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct SpeechRealtimeSession { + pub session_id: String, + pub input_sample_rate: u32, + pub output_sample_rate: u32, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct SpeechAppendRealtimeAudioRequest { + pub session_id: String, + /// Base64-encoded PCM16 little-endian mono audio at 16 kHz. + pub pcm16_base64: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct SpeechRealtimeSessionRequest { + pub session_id: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct SpeechRealtimeToolResultRequest { + pub session_id: String, + pub call_id: String, + pub result: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct SpeechRealtimeSpeakRequest { + pub session_id: String, + pub text: String, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum SpeechRealtimeEventKind { + Connected, + Ready, + UserSpeechStarted, + UserTranscriptDelta, + UserTranscriptCompleted, + AssistantTextDelta, + AssistantTextCompleted, + AssistantAudioStarted, + AssistantAudioDelta, + AssistantAudioCompleted, + FunctionCall, + Closed, + Error, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct SpeechRealtimeFunctionCall { + pub call_id: String, + pub name: String, + /// Provider-generated JSON arguments. Consumers must validate the decoded + /// object again before starting any BitFun task. + pub arguments: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct SpeechRealtimeEvent { + pub session_id: String, + pub kind: SpeechRealtimeEventKind, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub text: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub audio_base64: Option, + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub function_calls: Vec, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub provider_session_id: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub status_code: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub message: Option, +} + +#[cfg(test)] +mod tests { + use super::SpeechStartRealtimeSessionRequest; + + #[test] + fn realtime_start_request_accepts_legacy_empty_payload() { + let request: SpeechStartRealtimeSessionRequest = serde_json::from_str("{}").unwrap(); + assert_eq!(request.client_context, None); + } +} diff --git a/src/crates/contracts/events/src/lib.rs b/src/crates/contracts/events/src/lib.rs index 54a2d35aed..9890e3acf4 100644 --- a/src/crates/contracts/events/src/lib.rs +++ b/src/crates/contracts/events/src/lib.rs @@ -25,5 +25,7 @@ pub use bitfun_core_types::ToolImageAttachment; pub use catalog::{AIModelCatalogUpdatedEvent, AI_MODEL_CATALOG_UPDATED_EVENT}; pub use emitter::EventEmitter; pub use frontend_projection::{project_agentic_frontend_event, AgenticFrontendEvent}; -pub use speech::{SPEECH_MODEL_PROGRESS_EVENT, SPEECH_MODEL_STATUS_CHANGED_EVENT}; +pub use speech::{ + SPEECH_MODEL_PROGRESS_EVENT, SPEECH_MODEL_STATUS_CHANGED_EVENT, SPEECH_REALTIME_EVENT, +}; pub use types::*; diff --git a/src/crates/contracts/events/src/speech.rs b/src/crates/contracts/events/src/speech.rs index e8bec2e376..6f28688cb7 100644 --- a/src/crates/contracts/events/src/speech.rs +++ b/src/crates/contracts/events/src/speech.rs @@ -1,2 +1,3 @@ pub const SPEECH_MODEL_PROGRESS_EVENT: &str = "speech://model-download-progress"; pub const SPEECH_MODEL_STATUS_CHANGED_EVENT: &str = "speech://model-status-changed"; +pub const SPEECH_REALTIME_EVENT: &str = "speech://realtime-event"; diff --git a/src/crates/services/services-integrations/Cargo.toml b/src/crates/services/services-integrations/Cargo.toml index 4b08feccbb..5972b0d84a 100644 --- a/src/crates/services/services-integrations/Cargo.toml +++ b/src/crates/services/services-integrations/Cargo.toml @@ -388,6 +388,13 @@ speech = [ "tokio-util", "uuid", ] +speech-realtime = [ + "speech", + "tokio/net", + "tokio/time", + "dep:tokio-tungstenite", + "tokio-tungstenite?/rustls-tls-native-roots", +] workspace-search = [ "async-trait", "bitfun-services-core/filesystem", diff --git a/src/crates/services/services-integrations/src/speech/mod.rs b/src/crates/services/services-integrations/src/speech/mod.rs index e5dca07949..9838029850 100644 --- a/src/crates/services/services-integrations/src/speech/mod.rs +++ b/src/crates/services/services-integrations/src/speech/mod.rs @@ -6,6 +6,8 @@ mod error; mod model_catalog; mod model_store; mod qwen3_asr_int8; +#[cfg(feature = "speech-realtime")] +mod realtime; mod recognizer; mod recognizer_router; mod sensevoice_int8; @@ -27,6 +29,8 @@ use base64::Engine; pub use bitfun_core_types::speech::*; pub use error::{SpeechError, SpeechResult}; use error::{SpeechError as BitFunError, SpeechResult as BitFunResult}; +#[cfg(feature = "speech-realtime")] +pub use realtime::VolcengineRealtimeSpeechConfig; use std::collections::HashMap; use std::path::PathBuf; use std::sync::Arc; @@ -71,6 +75,8 @@ pub struct SpeechService { recognizer: Arc, downloads: Arc>>>, sessions: Arc>>, + #[cfg(feature = "speech-realtime")] + realtime: realtime::RealtimeSpeechRegistry, } #[derive(Debug)] @@ -93,6 +99,8 @@ impl SpeechService { recognizer: Arc::new(SpeechRecognizerRouter::new()), downloads: Arc::new(Mutex::new(HashMap::new())), sessions: Arc::new(Mutex::new(HashMap::new())), + #[cfg(feature = "speech-realtime")] + realtime: realtime::RealtimeSpeechRegistry::default(), } } diff --git a/src/crates/services/services-integrations/src/speech/realtime.rs b/src/crates/services/services-integrations/src/speech/realtime.rs new file mode 100644 index 0000000000..853bc118df --- /dev/null +++ b/src/crates/services/services-integrations/src/speech/realtime.rs @@ -0,0 +1,1019 @@ +use super::{BitFunError, BitFunResult, SpeechService}; +use base64::engine::general_purpose::STANDARD as BASE64_STANDARD; +use base64::Engine; +use bitfun_core_types::speech::{ + SpeechAppendRealtimeAudioRequest, SpeechRealtimeEvent, SpeechRealtimeEventKind, + SpeechRealtimeFunctionCall, SpeechRealtimeSession, SpeechRealtimeSessionRequest, + SpeechRealtimeSpeakRequest, SpeechRealtimeToolResultRequest, +}; +use futures_util::{SinkExt, StreamExt}; +use serde_json::{json, Value}; +use std::collections::HashMap; +use std::sync::Arc; +use std::time::Duration; +use tokio::sync::{mpsc, oneshot, Mutex}; +use tokio::time::timeout; +use tokio_tungstenite::tungstenite::{ + client::IntoClientRequest, + http::{header::HeaderName, HeaderValue}, + Message, +}; +use tokio_util::sync::CancellationToken; +use uuid::Uuid; + +const VOLCENGINE_REALTIME_URL: &str = + "wss://openspeech.bytedance.com/api/v3/duplex/realtime/dialogue"; +const VOLCENGINE_REALTIME_MODEL: &str = "1.2.6.1"; +const INPUT_SAMPLE_RATE: u32 = 16_000; +const OUTPUT_SAMPLE_RATE: u32 = 24_000; +const CONNECT_TIMEOUT: Duration = Duration::from_secs(15); +const COMMAND_QUEUE_CAPACITY: usize = 256; +const MAX_AUDIO_CHUNK_BYTES: usize = 64 * 1024; +const MAX_TOOL_RESULT_BYTES: usize = 16 * 1024; +const MAX_SPOKEN_PROGRESS_CHARS: usize = 600; +const MAX_CLIENT_CONTEXT_CHARS: usize = 16 * 1024; + +const BITFUN_VOICE_INSTRUCTIONS: &str = r#"You are BitFun's client-level realtime voice assistant. Reply naturally and concisely in the user's language. Your voice call belongs to the whole BitFun client, not to one chat session or one workspace. + +Use get_bitfun_client_context whenever the user asks about the current client, open workspaces/projects, visible sessions, running tasks, or names a workspace whose exact id is not already known from a fresh context result. Never guess a workspace id. Use switch_bitfun_workspace for navigation-only requests. When the user asks you to inspect, create, change, run, debug, research, or otherwise complete work, call run_bitfun_task with a complete standalone task description and the intended workspace_id. Omit workspace_id only when the user clearly means the active workspace. Set activate_workspace to true when the user asks to enter, switch to, or visibly work in that workspace; use false only for an explicit background request. + +If the user asks to stop, cancel, abort, or interrupt the BitFun task currently running through this client voice assistant, call stop_bitfun_task immediately. A stop request is a control operation, not a new task: never pass it to run_bitfun_task and never claim the task stopped before the stop_bitfun_task result confirms it. Do not claim that work is complete before the tool result arrives. BitFun will speak brief public progress summaries while the Agent task is running; do not expose private reasoning, raw logs, or tool payloads. After the tool result arrives, summarize the outcome clearly and mention any user action still required. Never invent client state or task results."#; + +#[derive(Debug, Clone)] +pub struct VolcengineRealtimeSpeechConfig { + pub api_key: String, + pub voice: String, + pub speed: i32, + pub loudness: i32, + pub client_context: Option, +} + +#[derive(Clone, Default)] +pub(crate) struct RealtimeSpeechRegistry { + sessions: Arc>>, +} + +#[derive(Clone)] +struct RealtimeSpeechSessionHandle { + sender: mpsc::Sender, + cancel: CancellationToken, +} + +enum RealtimeSpeechCommand { + Send(Value), + Close, +} + +type RealtimeEventHandler = Arc; + +impl SpeechService { + pub async fn start_realtime_session( + &self, + config: VolcengineRealtimeSpeechConfig, + on_event: F, + ) -> BitFunResult + where + F: Fn(SpeechRealtimeEvent) + Send + Sync + 'static, + { + let config = validate_config(config)?; + let session_id = Uuid::new_v4().to_string(); + let connect_id = Uuid::new_v4().to_string(); + let (sender, receiver) = mpsc::channel(COMMAND_QUEUE_CAPACITY); + let cancel = CancellationToken::new(); + let handle = RealtimeSpeechSessionHandle { + sender, + cancel: cancel.clone(), + }; + + self.realtime + .sessions + .lock() + .await + .insert(session_id.clone(), handle); + + let (ready_sender, ready_receiver) = oneshot::channel(); + let event_handler: RealtimeEventHandler = Arc::new(on_event); + let registry = self.realtime.clone(); + let actor_session_id = session_id.clone(); + tokio::spawn(async move { + run_realtime_actor( + actor_session_id.clone(), + connect_id, + config, + receiver, + cancel, + ready_sender, + event_handler, + ) + .await; + registry.sessions.lock().await.remove(&actor_session_id); + }); + + match timeout(CONNECT_TIMEOUT + Duration::from_secs(2), ready_receiver).await { + Ok(Ok(Ok(()))) => Ok(SpeechRealtimeSession { + session_id, + input_sample_rate: INPUT_SAMPLE_RATE, + output_sample_rate: OUTPUT_SAMPLE_RATE, + }), + Ok(Ok(Err(message))) => { + self.realtime.sessions.lock().await.remove(&session_id); + Err(BitFunError::service(message)) + } + Ok(Err(_)) => { + self.realtime.sessions.lock().await.remove(&session_id); + Err(BitFunError::service( + "Realtime speech connection ended before it became ready", + )) + } + Err(_) => { + if let Some(handle) = self.realtime.sessions.lock().await.remove(&session_id) { + handle.cancel.cancel(); + } + Err(BitFunError::service( + "Timed out connecting to the realtime speech service", + )) + } + } + } + + pub async fn append_realtime_audio( + &self, + request: SpeechAppendRealtimeAudioRequest, + ) -> BitFunResult<()> { + let audio = BASE64_STANDARD + .decode(request.pcm16_base64.as_bytes()) + .map_err(|error| { + BitFunError::validation(format!("Invalid base64 realtime audio chunk: {error}")) + })?; + if audio.is_empty() || audio.len() % 2 != 0 { + return Err(BitFunError::validation( + "Realtime PCM16 audio must contain complete samples", + )); + } + if audio.len() > MAX_AUDIO_CHUNK_BYTES { + return Err(BitFunError::validation(format!( + "Realtime audio chunk exceeds {MAX_AUDIO_CHUNK_BYTES} bytes" + ))); + } + self.send_realtime_command( + &request.session_id, + json!({ + "type": "input_audio_buffer.append", + "event_id": Uuid::new_v4().simple().to_string(), + "audio": request.pcm16_base64, + }), + ) + .await + } + + pub async fn commit_realtime_audio( + &self, + request: SpeechRealtimeSessionRequest, + ) -> BitFunResult<()> { + self.send_realtime_command( + &request.session_id, + json!({ + "type": "input_audio_buffer.commit", + "event_id": Uuid::new_v4().simple().to_string(), + }), + ) + .await + } + + pub async fn send_realtime_tool_result( + &self, + request: SpeechRealtimeToolResultRequest, + ) -> BitFunResult<()> { + let call_id = request.call_id.trim(); + if call_id.is_empty() { + return Err(BitFunError::validation("Function call id cannot be empty")); + } + if request.result.len() > MAX_TOOL_RESULT_BYTES { + return Err(BitFunError::validation(format!( + "Realtime function result exceeds {MAX_TOOL_RESULT_BYTES} bytes" + ))); + } + self.send_realtime_command( + &request.session_id, + tool_result_payload(call_id, &request.result), + ) + .await + } + + pub async fn speak_realtime_text( + &self, + request: SpeechRealtimeSpeakRequest, + ) -> BitFunResult<()> { + let text = request.text.trim(); + if text.is_empty() { + return Err(BitFunError::validation( + "Spoken progress text cannot be empty", + )); + } + if text.chars().count() > MAX_SPOKEN_PROGRESS_CHARS { + return Err(BitFunError::validation(format!( + "Spoken progress text exceeds {MAX_SPOKEN_PROGRESS_CHARS} characters" + ))); + } + self.send_realtime_command(&request.session_id, spoken_text_payload(text)) + .await + } + + pub async fn cancel_realtime_response( + &self, + request: SpeechRealtimeSessionRequest, + ) -> BitFunResult<()> { + self.send_realtime_command( + &request.session_id, + json!({ + "type": "response.cancel", + "event_id": Uuid::new_v4().simple().to_string(), + }), + ) + .await + } + + pub async fn close_realtime_session( + &self, + request: SpeechRealtimeSessionRequest, + ) -> BitFunResult<()> { + let handle = self + .realtime + .sessions + .lock() + .await + .remove(&request.session_id) + .ok_or_else(|| { + BitFunError::NotFound("Realtime speech session not found".to_string()) + })?; + handle + .sender + .send(RealtimeSpeechCommand::Close) + .await + .map_err(|_| BitFunError::service("Realtime speech session is already closed")) + } + + async fn send_realtime_command(&self, session_id: &str, payload: Value) -> BitFunResult<()> { + let handle = self + .realtime + .sessions + .lock() + .await + .get(session_id) + .cloned() + .ok_or_else(|| { + BitFunError::NotFound("Realtime speech session not found".to_string()) + })?; + handle + .sender + .send(RealtimeSpeechCommand::Send(payload)) + .await + .map_err(|_| BitFunError::service("Realtime speech session is already closed")) + } +} + +fn tool_result_payload(call_id: &str, result: &str) -> Value { + json!({ + "type": "conversation.item.create", + "event_id": Uuid::new_v4().simple().to_string(), + "items": [{ + "type": "message", + "call_id": call_id, + "role": "tool", + "content": [{ "type": "input_text", "text": result }], + }], + }) +} + +fn spoken_text_payload(text: &str) -> Value { + json!({ + "type": "speech_text_buffer.commit", + "event_id": Uuid::new_v4().simple().to_string(), + "speech_id": Uuid::new_v4().to_string(), + "text": text, + }) +} + +fn validate_config( + mut config: VolcengineRealtimeSpeechConfig, +) -> BitFunResult { + config.api_key = config.api_key.trim().to_string(); + config.voice = config.voice.trim().to_string(); + if config.api_key.is_empty() { + return Err(BitFunError::validation( + "Volcengine realtime speech API key is not configured", + )); + } + if config.voice.is_empty() { + return Err(BitFunError::validation( + "Volcengine realtime speech voice is not configured", + )); + } + if !(-50..=100).contains(&config.speed) || !(-50..=100).contains(&config.loudness) { + return Err(BitFunError::validation( + "Realtime speech speed and loudness must be between -50 and 100", + )); + } + config.client_context = config.client_context.and_then(|context| { + let trimmed = context.trim(); + if trimmed.is_empty() { + None + } else { + Some(trim_chars(trimmed, MAX_CLIENT_CONTEXT_CHARS)) + } + }); + Ok(config) +} + +fn trim_chars(value: &str, max_chars: usize) -> String { + if value.chars().count() <= max_chars { + return value.to_string(); + } + value.chars().take(max_chars).collect() +} + +async fn run_realtime_actor( + session_id: String, + connect_id: String, + config: VolcengineRealtimeSpeechConfig, + mut receiver: mpsc::Receiver, + cancel: CancellationToken, + ready_sender: oneshot::Sender>, + on_event: RealtimeEventHandler, +) { + bitfun_services_core::tls_provider::ensure_ring_crypto_provider(); + let mut request = match VOLCENGINE_REALTIME_URL.into_client_request() { + Ok(request) => request, + Err(error) => { + let message = format!("Failed to build realtime speech request: {error}"); + let _ = ready_sender.send(Err(message.clone())); + emit_error(&on_event, &session_id, message); + return; + } + }; + let api_key_header = HeaderName::from_static("x-api-key"); + let connect_id_header = HeaderName::from_static("x-api-connect-id"); + let api_key = match HeaderValue::from_str(&config.api_key) { + Ok(value) => value, + Err(_) => { + let message = + "Volcengine realtime speech API key contains invalid header characters".to_string(); + let _ = ready_sender.send(Err(message.clone())); + emit_error(&on_event, &session_id, message); + return; + } + }; + let connect_header = match HeaderValue::from_str(&connect_id) { + Ok(value) => value, + Err(_) => { + let message = "Failed to construct realtime speech connection id".to_string(); + let _ = ready_sender.send(Err(message.clone())); + emit_error(&on_event, &session_id, message); + return; + } + }; + request.headers_mut().insert(api_key_header, api_key); + request + .headers_mut() + .insert(connect_id_header, connect_header); + + let websocket = match timeout(CONNECT_TIMEOUT, tokio_tungstenite::connect_async(request)).await + { + Ok(Ok((stream, _response))) => stream, + Ok(Err(error)) => { + let message = format!("Failed to connect to realtime speech service: {error}"); + let _ = ready_sender.send(Err(message.clone())); + emit_error(&on_event, &session_id, message); + return; + } + Err(_) => { + let message = "Timed out connecting to the realtime speech service".to_string(); + let _ = ready_sender.send(Err(message.clone())); + emit_error(&on_event, &session_id, message); + return; + } + }; + + emit_event( + &on_event, + SpeechRealtimeEvent { + session_id: session_id.clone(), + kind: SpeechRealtimeEventKind::Connected, + text: None, + audio_base64: None, + function_calls: Vec::new(), + provider_session_id: None, + status_code: None, + message: None, + }, + ); + + let (mut sink, mut stream) = websocket.split(); + let create_payload = session_create_payload(&config, &session_id); + if let Err(error) = send_json(&mut sink, create_payload).await { + let message = format!("Failed to create realtime speech session: {error}"); + let _ = ready_sender.send(Err(message.clone())); + emit_error(&on_event, &session_id, message); + return; + } + let mut ready_sender = Some(ready_sender); + let mut closed_emitted = false; + + loop { + tokio::select! { + _ = cancel.cancelled() => { + let _ = send_json(&mut sink, json!({ + "type": "session.close", + "event_id": Uuid::new_v4().simple().to_string(), + })).await; + break; + } + command = receiver.recv() => { + match command { + Some(RealtimeSpeechCommand::Send(payload)) => { + if let Err(error) = send_json(&mut sink, payload).await { + emit_error( + &on_event, + &session_id, + format!("Failed to send realtime speech request: {error}"), + ); + break; + } + } + Some(RealtimeSpeechCommand::Close) => { + let _ = send_json(&mut sink, json!({ + "type": "session.close", + "event_id": Uuid::new_v4().simple().to_string(), + })).await; + break; + } + None => break, + } + } + incoming = stream.next() => { + match incoming { + Some(Ok(Message::Text(text))) => { + let session_created = provider_frame_is_type(text.as_ref(), "session.created"); + if let Some(was_closed) = process_provider_frame( + &session_id, + text.as_ref(), + &on_event, + ) { + if let Some(sender) = ready_sender.take() { + let _ = sender.send(Err(provider_startup_error(text.as_ref()))); + } + closed_emitted = was_closed; + break; + } + if session_created { + if let Some(sender) = ready_sender.take() { + let _ = sender.send(Ok(())); + } + } + } + Some(Ok(Message::Binary(bytes))) => { + let text = match std::str::from_utf8(bytes.as_ref()) { + Ok(text) => text, + Err(error) => { + emit_error( + &on_event, + &session_id, + format!("Realtime speech service returned a non-UTF-8 binary event: {error}"), + ); + break; + } + }; + let session_created = provider_frame_is_type(text, "session.created"); + if let Some(was_closed) = process_provider_frame( + &session_id, + text, + &on_event, + ) { + if let Some(sender) = ready_sender.take() { + let _ = sender.send(Err(provider_startup_error(text))); + } + closed_emitted = was_closed; + break; + } + if session_created { + if let Some(sender) = ready_sender.take() { + let _ = sender.send(Ok(())); + } + } + } + Some(Ok(Message::Ping(payload))) => { + if sink.send(Message::Pong(payload)).await.is_err() { + break; + } + } + Some(Ok(Message::Close(_))) | None => break, + Some(Ok(_)) => {} + Some(Err(error)) => { + emit_error( + &on_event, + &session_id, + format!("Realtime speech connection failed: {error}"), + ); + break; + } + } + } + } + } + + if let Some(sender) = ready_sender.take() { + let _ = sender.send(Err( + "Realtime speech connection ended before session.created".to_string(), + )); + } + let _ = sink.close().await; + if !closed_emitted { + emit_event( + &on_event, + SpeechRealtimeEvent { + session_id, + kind: SpeechRealtimeEventKind::Closed, + text: None, + audio_base64: None, + function_calls: Vec::new(), + provider_session_id: None, + status_code: None, + message: None, + }, + ); + } +} + +fn provider_frame_is_type(text: &str, expected: &str) -> bool { + serde_json::from_str::(text) + .ok() + .and_then(|payload| { + payload + .get("type") + .and_then(Value::as_str) + .map(str::to_string) + }) + .as_deref() + == Some(expected) +} + +fn provider_startup_error(text: &str) -> String { + let payload = serde_json::from_str::(text).ok(); + payload + .as_ref() + .and_then(|value| { + value + .get("message") + .or_else(|| value.pointer("/error/message")) + .and_then(Value::as_str) + }) + .map(|message| format!("Realtime speech session failed: {message}")) + .unwrap_or_else(|| "Realtime speech session ended before it became ready".to_string()) +} + +fn session_create_payload( + config: &VolcengineRealtimeSpeechConfig, + provider_session_id: &str, +) -> Value { + let instructions = match config.client_context.as_deref() { + Some(context) => format!( + "{BITFUN_VOICE_INSTRUCTIONS}\n\nCLIENT CONTEXT SNAPSHOT AT CALL START (JSON; refresh with get_bitfun_client_context before relying on changing state):\n{context}" + ), + None => BITFUN_VOICE_INSTRUCTIONS.to_string(), + }; + json!({ + "type": "session.create", + "event_id": Uuid::new_v4().simple().to_string(), + "session": { + "id": provider_session_id, + "model": VOLCENGINE_REALTIME_MODEL, + "instructions": instructions, + "audio": { + "input": { "format": { "type": "pcm", "rate": INPUT_SAMPLE_RATE } }, + "output": { + "format": { "type": "pcm_s16le", "rate": OUTPUT_SAMPLE_RATE }, + "voice": config.voice, + "speed": config.speed, + "loudness": config.loudness, + }, + }, + "tools": [ + { + "type": "function", + "name": "get_bitfun_client_context", + "description": "Return a fresh snapshot of the BitFun client: active scene, active and opened workspaces, visible sessions, running Agent tasks, and the task owned by this voice assistant. Call this before resolving a workspace name or answering questions about current client state.", + "parameters": { + "type": "object", + "additionalProperties": false, + "properties": {} + } + }, + { + "type": "function", + "name": "switch_bitfun_workspace", + "description": "Activate one currently opened BitFun workspace in the client. Obtain the exact workspace_id from get_bitfun_client_context and never guess it.", + "parameters": { + "type": "object", + "additionalProperties": false, + "properties": { + "workspace_id": { + "type": "string", + "description": "Exact id of an opened workspace from get_bitfun_client_context." + } + }, + "required": ["workspace_id"] + } + }, + { + "type": "function", + "name": "run_bitfun_task", + "description": "Start a new BitFun Agent session in the intended opened workspace, autonomously complete the requested task with normal tools and permissions, and return the final result. Use get_bitfun_client_context first when the user names a workspace or project.", + "parameters": { + "type": "object", + "additionalProperties": false, + "properties": { + "task": { + "type": "string", + "description": "A complete standalone description of the work BitFun should perform." + }, + "workspace_id": { + "type": "string", + "description": "Exact opened workspace id from get_bitfun_client_context. Omit only when the user clearly means the active workspace." + }, + "activate_workspace": { + "type": "boolean", + "description": "Whether to activate the target workspace and show the new Agent session. Defaults to true; use false only for an explicit background request." + } + }, + "required": ["task"] + } + }, + { + "type": "function", + "name": "stop_bitfun_task", + "description": "Stop the BitFun Agent task currently owned by this client-level voice assistant. Use this for any user request to stop, cancel, abort, or interrupt the current task.", + "parameters": { + "type": "object", + "additionalProperties": false, + "properties": {} + } + } + ], + }, + "extension": { + "asr": { "extra": {} }, + "tts": { "extra": {} }, + "dialog": { + "extra": { + "enable_loudness_norm": true, + "enable_music": false, + } + }, + "extra": { "enable_proactive_speak": true }, + }, + }) +} + +/// Process either a WebSocket text frame or a UTF-8 JSON payload carried in a +/// binary frame. The official Volcengine clients accept both frame kinds. +/// `Some(true)` means a provider close event was emitted; `Some(false)` means +/// the actor should terminate after an error event. +fn process_provider_frame( + session_id: &str, + text: &str, + on_event: &RealtimeEventHandler, +) -> Option { + let payload = match serde_json::from_str::(text) { + Ok(payload) => payload, + Err(error) => { + emit_error( + on_event, + session_id, + format!("Realtime speech service returned invalid JSON: {error}"), + ); + return Some(false); + } + }; + let event = parse_provider_event(session_id, &payload)?; + let was_closed = event.kind == SpeechRealtimeEventKind::Closed; + let was_error = event.kind == SpeechRealtimeEventKind::Error; + emit_event(on_event, event); + if was_closed { + Some(true) + } else if was_error { + Some(false) + } else { + None + } +} + +async fn send_json(sink: &mut S, payload: Value) -> Result<(), String> +where + S: futures_util::Sink + Unpin, + S::Error: std::fmt::Display, +{ + let text = serde_json::to_string(&payload).map_err(|error| error.to_string())?; + sink.send(Message::Text(text.into())) + .await + .map_err(|error| error.to_string()) +} + +fn parse_provider_event(session_id: &str, payload: &Value) -> Option { + let event_type = payload.get("type")?.as_str()?; + let mut event = SpeechRealtimeEvent { + session_id: session_id.to_string(), + kind: SpeechRealtimeEventKind::Ready, + text: None, + audio_base64: None, + function_calls: Vec::new(), + provider_session_id: None, + status_code: status_code(payload), + message: payload + .get("message") + .or_else(|| payload.pointer("/error/message")) + .and_then(Value::as_str) + .map(str::to_string), + }; + + match event_type { + "session.created" => { + event.kind = SpeechRealtimeEventKind::Ready; + event.provider_session_id = payload + .pointer("/session/id") + .or_else(|| payload.get("id")) + .and_then(Value::as_str) + .map(str::to_string); + } + "conversation.item.input_audio_transcription.started" => { + event.kind = SpeechRealtimeEventKind::UserSpeechStarted; + } + "conversation.item.input_audio_transcription.delta" => { + event.kind = SpeechRealtimeEventKind::UserTranscriptDelta; + event.text = event_text(payload); + } + "conversation.item.input_audio_transcription.completed" => { + event.kind = SpeechRealtimeEventKind::UserTranscriptCompleted; + event.text = event_text(payload); + } + "response.output_text.delta" => { + event.kind = SpeechRealtimeEventKind::AssistantTextDelta; + event.text = event_text(payload); + } + "response.output_text.done" => { + event.kind = SpeechRealtimeEventKind::AssistantTextCompleted; + event.text = event_text(payload); + } + "response.output_audio.started" => { + event.kind = SpeechRealtimeEventKind::AssistantAudioStarted; + } + "response.output_audio.delta" => { + event.kind = SpeechRealtimeEventKind::AssistantAudioDelta; + event.audio_base64 = payload + .get("audio") + .or_else(|| payload.get("delta")) + .and_then(Value::as_str) + .map(str::to_string); + } + "response.output_audio.done" => { + event.kind = SpeechRealtimeEventKind::AssistantAudioCompleted; + } + "response.function_call_arguments.done" => { + event.kind = SpeechRealtimeEventKind::FunctionCall; + let items = match payload.get("items") { + Some(Value::Array(items)) => items.clone(), + Some(Value::Object(_)) => vec![payload["items"].clone()], + _ => Vec::new(), + }; + event.function_calls = items + .iter() + .filter_map(|item| { + let call_id = item.get("call_id")?.as_str()?.trim(); + let name = item + .get("name") + .or_else(|| item.pointer("/function/name"))? + .as_str()? + .trim(); + if call_id.is_empty() || name.is_empty() { + return None; + } + let arguments = item + .get("arguments") + .or_else(|| item.pointer("/function/arguments")) + .map(|value| { + value + .as_str() + .map(str::to_string) + .unwrap_or_else(|| value.to_string()) + }) + .unwrap_or_else(|| "{}".to_string()); + Some(SpeechRealtimeFunctionCall { + call_id: call_id.to_string(), + name: name.to_string(), + arguments, + }) + }) + .collect(); + } + "session.closed" => { + event.kind = SpeechRealtimeEventKind::Closed; + } + "error" => { + event.kind = SpeechRealtimeEventKind::Error; + if event.message.is_none() { + event.message = Some("Realtime speech service reported an error".to_string()); + } + } + _ => return None, + } + + Some(event) +} + +fn event_text(payload: &Value) -> Option { + ["delta", "transcript", "text", "content"] + .into_iter() + .find_map(|key| payload.get(key).and_then(Value::as_str)) + .or_else(|| payload.pointer("/item/transcript").and_then(Value::as_str)) + .map(str::to_string) +} + +fn status_code(payload: &Value) -> Option { + payload.get("status_code").and_then(|value| { + value + .as_i64() + .or_else(|| value.as_str().and_then(|code| code.parse().ok())) + }) +} + +fn emit_event(handler: &RealtimeEventHandler, event: SpeechRealtimeEvent) { + handler(event); +} + +fn emit_error(handler: &RealtimeEventHandler, session_id: &str, message: String) { + emit_event( + handler, + SpeechRealtimeEvent { + session_id: session_id.to_string(), + kind: SpeechRealtimeEventKind::Error, + text: None, + audio_base64: None, + function_calls: Vec::new(), + provider_session_id: None, + status_code: None, + message: Some(message), + }, + ); +} + +#[cfg(test)] +mod tests { + use super::*; + + fn config() -> VolcengineRealtimeSpeechConfig { + VolcengineRealtimeSpeechConfig { + api_key: "fixture-key".to_string(), + voice: "zh_female_vv_jupiter_bigtts".to_string(), + speed: 0, + loudness: 0, + client_context: Some(r#"{"active_workspace":{"id":"workspace-1"}}"#.to_string()), + } + } + + #[test] + fn session_create_uses_documented_pcm_contract_and_bitfun_tool() { + let payload = session_create_payload(&config(), "local-session"); + assert_eq!( + payload.pointer("/session/id"), + Some(&json!("local-session")) + ); + assert_eq!(payload.pointer("/session/model"), Some(&json!("1.2.6.1"))); + assert_eq!( + payload.pointer("/session/audio/input/format"), + Some(&json!({"type": "pcm", "rate": 16000})) + ); + assert_eq!( + payload.pointer("/session/audio/output/format"), + Some(&json!({"type": "pcm_s16le", "rate": 24000})) + ); + assert_eq!( + payload.pointer("/session/tools/0/name"), + Some(&json!("get_bitfun_client_context")) + ); + assert_eq!( + payload.pointer("/session/tools/1/name"), + Some(&json!("switch_bitfun_workspace")) + ); + assert_eq!( + payload.pointer("/session/tools/2/name"), + Some(&json!("run_bitfun_task")) + ); + assert_eq!( + payload.pointer("/session/tools/3/name"), + Some(&json!("stop_bitfun_task")) + ); + let instructions = payload + .pointer("/session/instructions") + .and_then(Value::as_str) + .unwrap(); + assert!(instructions.contains("never claim the task stopped")); + assert!(instructions.contains("workspace-1")); + assert_eq!( + payload.pointer("/extension/extra/enable_proactive_speak"), + Some(&json!(true)) + ); + assert_eq!( + payload.pointer("/extension/dialog/extra/enable_loudness_norm"), + Some(&json!(true)) + ); + } + + #[test] + fn function_call_event_preserves_call_identity_and_arguments() { + let event = parse_provider_event( + "local-session", + &json!({ + "type": "response.function_call_arguments.done", + "items": [{ + "call_id": "call-1", + "name": "run_bitfun_task", + "arguments": "{\"task\":\"run tests\"}" + }] + }), + ) + .unwrap(); + + assert_eq!(event.kind, SpeechRealtimeEventKind::FunctionCall); + assert_eq!(event.function_calls.len(), 1); + assert_eq!(event.function_calls[0].call_id, "call-1"); + assert_eq!( + event.function_calls[0].arguments, + "{\"task\":\"run tests\"}" + ); + + let single_item_event = parse_provider_event( + "local-session", + &json!({ + "type": "response.function_call_arguments.done", + "items": { + "call_id": "call-2", + "function": { + "name": "run_bitfun_task", + "arguments": {"task": "inspect the workspace"} + } + } + }), + ) + .unwrap(); + assert_eq!(single_item_event.function_calls[0].call_id, "call-2"); + assert_eq!( + single_item_event.function_calls[0].arguments, + "{\"task\":\"inspect the workspace\"}" + ); + } + + #[test] + fn documented_audio_delta_is_forwarded_from_json_frames() { + let events = Arc::new(std::sync::Mutex::new(Vec::new())); + let captured = events.clone(); + let handler: RealtimeEventHandler = Arc::new(move |event| { + captured.lock().unwrap().push(event); + }); + + let outcome = process_provider_frame( + "local-session", + r#"{"type":"response.output_audio.delta","delta":"AQIDBA=="}"#, + &handler, + ); + + assert_eq!(outcome, None); + let events = events.lock().unwrap(); + assert_eq!(events.len(), 1); + assert_eq!(events[0].kind, SpeechRealtimeEventKind::AssistantAudioDelta); + assert_eq!(events[0].audio_base64.as_deref(), Some("AQIDBA==")); + } + + #[test] + fn tool_results_and_spoken_progress_match_the_provider_contract() { + let tool_result = tool_result_payload("call-1", r#"{"ok":true}"#); + assert_eq!( + tool_result.pointer("/items/0/type"), + Some(&json!("message")) + ); + assert_eq!(tool_result.pointer("/items/0/role"), Some(&json!("tool"))); + assert_eq!( + tool_result.pointer("/items/0/call_id"), + Some(&json!("call-1")) + ); + + let spoken = spoken_text_payload("A short progress update."); + assert_eq!( + spoken.get("type"), + Some(&json!("speech_text_buffer.commit")) + ); + assert!(spoken + .get("speech_id") + .and_then(Value::as_str) + .is_some_and(|value| !value.is_empty())); + assert_eq!(spoken.get("text"), Some(&json!("A short progress update."))); + } +} diff --git a/src/web-ui/src/app/App.tsx b/src/web-ui/src/app/App.tsx index e630808d32..bb8884f172 100644 --- a/src/web-ui/src/app/App.tsx +++ b/src/web-ui/src/app/App.tsx @@ -29,6 +29,7 @@ import { isStartupOverlayPresent, } from './startup/startupOverlay'; import { ToolbarModeProvider } from '../flow_chat/components/toolbar-mode/ToolbarModeProvider'; +import { RealtimeVoiceCallProvider } from '../flow_chat/components/voice/RealtimeVoiceCallContext'; import type { AgentCompanionPetCommand } from './services/agentCompanionPetCommands'; import AskUserAnnouncer from './components/NavPanel/AskUserAnnouncer'; import { shouldBlockBrowserShortcut } from './browserShortcutPolicy'; @@ -914,7 +915,8 @@ function App() { - + + {/* One shell-owned command/search surface for every scene and nav mode. */} @@ -947,7 +949,8 @@ function App() { Mounted here (inside ToolbarModeProvider, outside LazyAppLayout) so it persists across both normal and Toolbar Mode. */} - + + diff --git a/src/web-ui/src/app/components/NavPanel/MainNav.tsx b/src/web-ui/src/app/components/NavPanel/MainNav.tsx index 1d7cc85e37..be380cfb9e 100644 --- a/src/web-ui/src/app/components/NavPanel/MainNav.tsx +++ b/src/web-ui/src/app/components/NavPanel/MainNav.tsx @@ -2,7 +2,7 @@ * MainNav — primary product navigation sidebar. * * Layout (top to bottom): - * 1. Search and New Session + * 1. Search and client voice assistant * 2. AI Assistant, Task Board, Mini Apps, then Extensions & Compatibility * 3. Unified Sessions (all or grouped by project / assistant) * @@ -15,7 +15,7 @@ import { createPortal } from 'react-dom'; import { KeyHint } from '@bitfun/ui'; import { getAppearanceOverlayHost } from '@/infrastructure/appearance/runtime/AppearanceOverlayHost'; import { isImeOwnedKeyboardEvent } from '@/shared/utils/ime'; -import { Plus, FolderOpen, FolderPlus, History, Check, User, Users, Puzzle, ChevronDown, Network, Search, CalendarClock } from 'lucide-react'; +import { FolderOpen, FolderPlus, History, Check, User, Users, Puzzle, ChevronDown, Network, Search, CalendarClock } from 'lucide-react'; // import { PanelsTopLeft } from 'lucide-react'; // temporarily hidden: Pages nav entry import { BITFUN_ICON_SIZE, @@ -49,6 +49,7 @@ import { subscribeGlobalSearchShortcut, } from '@/app/global-search/globalSearchShortcut'; import { useExternalAppAwareness } from '@/infrastructure/config/components/external-sources/useExternalAppAwareness'; +import { RealtimeVoiceCallButton } from '@/flow_chat/components/voice/RealtimeVoiceCallButton'; import './NavPanel.scss'; @@ -245,10 +246,6 @@ const MainNav: React.FC = ({ }; }, [workspaceMenuOpen, updateWorkspaceMenuPos]); - const handleCreateSession = useCallback(() => { - void activateProductAction('session.new'); - }, []); - const handleOpenAgents = useCallback(() => { void activateProductAction('surface.agents.open'); }, []); @@ -381,7 +378,6 @@ const MainNav: React.FC = ({ getAppearanceOverlayHost() ) : null; - const createSessionLabel = t('nav.sessions.newSession'); const addSessionGroupTooltip = t('nav.tooltips.addSessionGroup'); const agentsTooltip = t('nav.tooltips.agents'); const skillsTooltip = t('nav.tooltips.skills'); @@ -395,7 +391,7 @@ const MainNav: React.FC = ({ const isTaskBoardActive = activeTabId === 'todos'; return ( <> - {/* ── Search and New Session ─────────────────── */} + {/* ── Search and client voice assistant ─────── */}
@@ -427,20 +423,7 @@ const MainNav: React.FC = ({
- - - +
diff --git a/src/web-ui/src/app/components/NavPanel/unifiedSessionCreation.test.ts b/src/web-ui/src/app/components/NavPanel/unifiedSessionCreation.test.ts index 3af8c87290..ff1b62dd7b 100644 --- a/src/web-ui/src/app/components/NavPanel/unifiedSessionCreation.test.ts +++ b/src/web-ui/src/app/components/NavPanel/unifiedSessionCreation.test.ts @@ -7,20 +7,20 @@ function source(relativePath: string): string { } describe('unified project session creation', () => { - it('exposes one icon-only unified action beside search', () => { + it('places the client voice action beside search instead of the new-session shortcut', () => { const mainNav = source('./MainNav.tsx'); const workspaceItem = source('./sections/workspaces/WorkspaceItem.tsx'); const utilityRowIndex = mainNav.indexOf('data-bf-part="utilityRow"'); - const newSessionIndex = mainNav.indexOf('data-testid="nav-new-session-btn"'); + const voiceIndex = mainNav.indexOf(''); const sectionsIndex = mainNav.indexOf('data-testid="nav-sections"'); const sessionsSectionIndex = mainNav.indexOf('data-bf-section="sessions"'); - expect(newSessionIndex).toBeGreaterThan(utilityRowIndex); - expect(sectionsIndex).toBeGreaterThan(newSessionIndex); - expect(sessionsSectionIndex).toBeGreaterThan(newSessionIndex); - expect(mainNav).toContain('