From 113a6ccdc2aed7df2d0aa6b711eba8069388eb0e Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Fri, 21 Aug 2026 08:38:04 -0700 Subject: [PATCH] feat(talk): integrate realtime relay with macOS Talk Co-authored-by: Zhilong Zheng --- apps/.i18n/native-source.json | 99 ++ .../Sources/OpenClaw/AppState+Preview.swift | 1 + .../Sources/OpenClaw/AppState+Talk.swift | 10 + apps/macos/Sources/OpenClaw/AppState.swift | 6 +- apps/macos/Sources/OpenClaw/Constants.swift | 5 + .../OpenClaw/TalkModeGatewayConfig.swift | 82 +- .../OpenClaw/TalkModeRuntime+Realtime.swift | 764 +++++++++++ .../Sources/OpenClaw/TalkModeRuntime.swift | 602 +++++---- .../TalkRecognitionCaptureLifecycle.swift | 57 + .../Sources/OpenClaw/VoiceWakeSettings.swift | 14 + .../TalkModeGatewayConfigTests.swift | 112 ++ .../TalkModeRuntimeSpeechTests.swift | 1167 ++++++++++++++++- docs/nodes/talk.md | 69 + 13 files changed, 2764 insertions(+), 224 deletions(-) create mode 100644 apps/macos/Sources/OpenClaw/AppState+Talk.swift create mode 100644 apps/macos/Sources/OpenClaw/TalkModeRuntime+Realtime.swift create mode 100644 apps/macos/Sources/OpenClaw/TalkRecognitionCaptureLifecycle.swift diff --git a/apps/.i18n/native-source.json b/apps/.i18n/native-source.json index d237bf739417..06321e08c6ec 100644 --- a/apps/.i18n/native-source.json +++ b/apps/.i18n/native-source.json @@ -39322,6 +39322,17 @@ } ] }, + { + "id": "native.apple.11f6f4fa597f9df4", + "source": "Realtime Talk requested an invalid audio sample rate", + "surface": "apple", + "sites": [ + { + "kind": "ui-localized-call", + "path": "apps/macos/Sources/OpenClaw/MacRealtimeTalkAudioCapture.swift" + } + ] + }, { "id": "native.apple.400963cb08a3c253", "source": "Realtime Voice", @@ -39440,6 +39451,28 @@ } ] }, + { + "id": "native.apple.48e407d88c565c02", + "source": "Realtime disconnected repeatedly — using native speech", + "surface": "apple", + "sites": [ + { + "kind": "ui-localized-call", + "path": "apps/macos/Sources/OpenClaw/TalkModeRuntime+Realtime.swift" + } + ] + }, + { + "id": "native.apple.ab73f15b8147eb41", + "source": "Realtime disconnected — reconnecting…", + "surface": "apple", + "sites": [ + { + "kind": "ui-localized-call", + "path": "apps/macos/Sources/OpenClaw/TalkModeRuntime+Realtime.swift" + } + ] + }, { "id": "native.apple.ca3a762269827f37", "source": "Realtime failed", @@ -39462,6 +39495,17 @@ } ] }, + { + "id": "native.apple.b5b843b6a0a96cc5", + "source": "Realtime microphone became unavailable: %@", + "surface": "apple", + "sites": [ + { + "kind": "ui-localized-call", + "path": "apps/macos/Sources/OpenClaw/MacRealtimeTalkAudioCapture.swift" + } + ] + }, { "id": "native.apple.96ece1b431ac3f96", "source": "Realtime output cancellation failed: %@", @@ -39495,6 +39539,17 @@ } ] }, + { + "id": "native.apple.ddb5344f621cd627", + "source": "Realtime unavailable — using native speech", + "surface": "apple", + "sites": [ + { + "kind": "ui-localized-call", + "path": "apps/macos/Sources/OpenClaw/TalkModeRuntime+Realtime.swift" + } + ] + }, { "id": "native.apple.4c945b4b74b4de2d", "source": "Realtime voice did not start.", @@ -42847,6 +42902,28 @@ } ] }, + { + "id": "native.apple.d40597f2d391f8cf", + "source": "Selected audio input has no usable Float32 format", + "surface": "apple", + "sites": [ + { + "kind": "ui-localized-call", + "path": "apps/macos/Sources/OpenClaw/MacRealtimeTalkAudioCapture.swift" + } + ] + }, + { + "id": "native.apple.8a4aa2d658ca95a3", + "source": "Selected input and system default are unavailable", + "surface": "apple", + "sites": [ + { + "kind": "ui-localized-call", + "path": "apps/macos/Sources/OpenClaw/MacRealtimeTalkAudioCapture.swift" + } + ] + }, { "id": "native.apple.209377172e2dd2aa", "source": "Selected reasoning is unavailable for this target.", @@ -48268,6 +48345,28 @@ } ] }, + { + "id": "native.apple.050cddb7d6d24e52", + "source": "Use realtime Gateway relay", + "surface": "apple", + "sites": [ + { + "kind": "ui-named-argument", + "path": "apps/macos/Sources/OpenClaw/VoiceWakeSettings.swift" + } + ] + }, + { + "id": "native.apple.d7c89897a6c2da7f", + "source": "Use the Gateway's configured realtime voice session on this Mac. Requires realtime, gateway-relay, and agent-consult in Talk settings.", + "surface": "apple", + "sites": [ + { + "kind": "ui-named-argument-multiline", + "path": "apps/macos/Sources/OpenClaw/VoiceWakeSettings.swift" + } + ] + }, { "id": "native.apple.58566125ab806740", "source": "Use the panel + Canvas", diff --git a/apps/macos/Sources/OpenClaw/AppState+Preview.swift b/apps/macos/Sources/OpenClaw/AppState+Preview.swift index b7d3efc254e0..199b1352cf84 100644 --- a/apps/macos/Sources/OpenClaw/AppState+Preview.swift +++ b/apps/macos/Sources/OpenClaw/AppState+Preview.swift @@ -20,6 +20,7 @@ extension AppState { state.voiceWakeAdditionalLocaleIDs = ["en-US", "de-DE"] state.voicePushToTalkEnabled = false state.talkEnabled = false + state.talkRealtimeRelayEnabled = false state.talkPhaseSoundsEnabled = true state.talkShiftToStopEnabled = true state.iconOverride = .system diff --git a/apps/macos/Sources/OpenClaw/AppState+Talk.swift b/apps/macos/Sources/OpenClaw/AppState+Talk.swift new file mode 100644 index 000000000000..0fa2abaff6a1 --- /dev/null +++ b/apps/macos/Sources/OpenClaw/AppState+Talk.swift @@ -0,0 +1,10 @@ +import Foundation + +extension AppState { + func persistTalkRealtimeRelayPreference(previousValue: Bool) { + guard !self.isPreview else { return } + AppDefaults.standard.set(self.talkRealtimeRelayEnabled, forKey: talkRealtimeRelayEnabledKey) + guard self.talkEnabled, self.talkRealtimeRelayEnabled != previousValue else { return } + Task { await TalkModeRuntime.shared.realtimeRelayPreferenceDidChange() } + } +} diff --git a/apps/macos/Sources/OpenClaw/AppState.swift b/apps/macos/Sources/OpenClaw/AppState.swift index 4bf2a944bc1e..6eaa4a686576 100644 --- a/apps/macos/Sources/OpenClaw/AppState.swift +++ b/apps/macos/Sources/OpenClaw/AppState.swift @@ -81,7 +81,7 @@ final class AppState { private static let logger = Logger(subsystem: "ai.openclaw", category: "app-state") private static let execApprovalsReadRetryAttempts = 5 - private let isPreview: Bool + let isPreview: Bool @ObservationIgnored private let execApprovalsDefaultsAsyncResolver: @MainActor () async -> Result @ObservationIgnored private let execApprovalsReadRetryDelay: Duration @@ -260,6 +260,10 @@ final class AppState { } } + var talkRealtimeRelayEnabled = isTalkRealtimeRelayEnabled() { + didSet { self.persistTalkRealtimeRelayPreference(previousValue: oldValue) } + } + var talkPhaseSoundsEnabled: Bool { didSet { self.ifNotPreview { diff --git a/apps/macos/Sources/OpenClaw/Constants.swift b/apps/macos/Sources/OpenClaw/Constants.swift index d9543d1a9e24..b74455eb86fd 100644 --- a/apps/macos/Sources/OpenClaw/Constants.swift +++ b/apps/macos/Sources/OpenClaw/Constants.swift @@ -31,6 +31,7 @@ let voiceWakeAdditionalLocalesKey = "openclaw.voiceWakeAdditionalLocaleIDs" let voicePushToTalkEnabledKey = "openclaw.voicePushToTalkEnabled" let voiceWakeTriggersTalkModeKey = "openclaw.voiceWakeTriggersTalkMode" let talkEnabledKey = "openclaw.talkEnabled" +let talkRealtimeRelayEnabledKey = "openclaw.talkRealtimeRelayEnabled" let talkPhaseSoundsEnabledKey = "openclaw.talkPhaseSoundsEnabled" let talkShiftToStopEnabledKey = "openclaw.talkShiftToStopEnabled" let iconOverrideKey = "openclaw.iconOverride" @@ -48,6 +49,10 @@ let cookieSyncEnabledKey = "openclaw.cookieSyncEnabled" let cookieSyncIntoProfileKey = "openclaw.cookieSyncIntoProfile" let cookieSyncDomainsKey = "openclaw.cookieSyncDomains" +func isTalkRealtimeRelayEnabled(defaults: UserDefaults = AppDefaults.standard) -> Bool { + defaults.object(forKey: talkRealtimeRelayEnabledKey) as? Bool ?? false +} + func isComputerControlEnabled(defaults: UserDefaults = AppDefaults.standard) -> Bool { // object(forKey:) preserves an explicit false; bool(forKey:) would conflate it with an unset default. defaults.object(forKey: computerControlEnabledKey) as? Bool ?? true diff --git a/apps/macos/Sources/OpenClaw/TalkModeGatewayConfig.swift b/apps/macos/Sources/OpenClaw/TalkModeGatewayConfig.swift index 6a5e397b6cf3..dfb6e12d0412 100644 --- a/apps/macos/Sources/OpenClaw/TalkModeGatewayConfig.swift +++ b/apps/macos/Sources/OpenClaw/TalkModeGatewayConfig.swift @@ -16,6 +16,18 @@ struct TalkModeGatewayConfigState { let referenceAudioPath: String? let referenceText: String? let seamColorHex: String? + let realtimeProvider: String? + let realtimeModelId: String? + let realtimeSpeakerVoice: String? + let realtimeMode: String? + let realtimeTransport: String? + let realtimeBrain: String? + + var hasGatewayRealtimeRelayTuple: Bool { + self.realtimeMode == "realtime" && + self.realtimeTransport == "gateway-relay" && + self.realtimeBrain == "agent-consult" + } } enum TalkModeGatewayConfigParser { @@ -62,6 +74,22 @@ enum TalkModeGatewayConfigParser { .trimmingCharacters(in: .whitespacesAndNewlines) let referenceText = activeConfig?["referenceText"]?.stringValue? .trimmingCharacters(in: .whitespacesAndNewlines) + let realtime = talk?["realtime"]?.dictionaryValue + let realtimeProviders = realtime?["providers"]?.dictionaryValue + let realtimeProvider = Self.firstString(realtime, keys: ["provider"]) + ?? Self.singleRealtimeProviderId(realtimeProviders) + let realtimeProviderConfig = Self.realtimeProviderConfig( + providers: realtimeProviders, + provider: realtimeProvider) + let realtimeModelId = Self.firstString(realtime, keys: ["model"]) + ?? Self.firstString(realtimeProviderConfig, keys: ["model"]) + let realtimeSpeakerVoice = Self.firstString( + realtime, + keys: ["speakerVoice", "voice"]) + ?? Self.firstString(realtimeProviderConfig, keys: ["speakerVoice", "voice"]) + let realtimeMode = Self.firstString(realtime, keys: ["mode"])?.lowercased() + let realtimeTransport = Self.firstString(realtime, keys: ["transport"])?.lowercased() + let realtimeBrain = Self.firstString(realtime, keys: ["brain"])?.lowercased() let resolvedVoice: String? = if activeProvider == defaultProvider { (voice?.trimmingCharacters(in: .whitespacesAndNewlines).isEmpty == false ? voice : nil) ?? (envVoice?.isEmpty == false ? envVoice : nil) ?? @@ -90,7 +118,13 @@ enum TalkModeGatewayConfigParser { apiKey: resolvedApiKey, referenceAudioPath: referenceAudioPath?.isEmpty == false ? referenceAudioPath : nil, referenceText: referenceText?.isEmpty == false ? referenceText : nil, - seamColorHex: rawSeam.isEmpty ? nil : rawSeam) + seamColorHex: rawSeam.isEmpty ? nil : rawSeam, + realtimeProvider: realtimeProvider, + realtimeModelId: realtimeModelId, + realtimeSpeakerVoice: realtimeSpeakerVoice, + realtimeMode: realtimeMode, + realtimeTransport: realtimeTransport, + realtimeBrain: realtimeBrain) } static func fallback( @@ -119,6 +153,50 @@ enum TalkModeGatewayConfigParser { apiKey: resolvedApiKey, referenceAudioPath: nil, referenceText: nil, - seamColorHex: nil) + seamColorHex: nil, + realtimeProvider: nil, + realtimeModelId: nil, + realtimeSpeakerVoice: nil, + realtimeMode: nil, + realtimeTransport: nil, + realtimeBrain: nil) + } + + private static func firstString( + _ config: [String: AnyCodable]?, + keys: [String]) -> String? + { + guard let config else { return nil } + for key in keys { + let value = config[key]?.stringValue?.trimmingCharacters(in: .whitespacesAndNewlines) + if value?.isEmpty == false { return value } + } + return nil + } + + private static func singleRealtimeProviderId(_ providers: [String: AnyCodable]?) -> String? { + guard let providers, providers.count == 1 else { return nil } + let provider = providers.keys.first?.trimmingCharacters(in: .whitespacesAndNewlines) + return provider?.isEmpty == false ? provider : nil + } + + private static func realtimeProviderConfig( + providers: [String: AnyCodable]?, + provider: String?) -> [String: AnyCodable]? + { + guard let providers else { return nil } + if let provider { + if let exact = providers[provider]?.dictionaryValue { + return exact + } + return providers.first { key, _ in + key.trimmingCharacters(in: .whitespacesAndNewlines) + .caseInsensitiveCompare(provider) == .orderedSame + }?.value.dictionaryValue + } + if providers.count == 1 { + return providers.values.first?.dictionaryValue + } + return nil } } diff --git a/apps/macos/Sources/OpenClaw/TalkModeRuntime+Realtime.swift b/apps/macos/Sources/OpenClaw/TalkModeRuntime+Realtime.swift new file mode 100644 index 000000000000..b6fdbb43581c --- /dev/null +++ b/apps/macos/Sources/OpenClaw/TalkModeRuntime+Realtime.swift @@ -0,0 +1,764 @@ +import Foundation +import OpenClawChatUI +import OpenClawKit +import OSLog + +private enum RealtimeRelayConfigurationError: LocalizedError { + case incompatible + + var errorDescription: String? { + "Gateway realtime Talk configuration changed before startup" + } +} + +extension TalkModeRuntime { + #if DEBUG + typealias VoiceWakePermissionProvider = @Sendable () async -> Bool + #endif + + private enum ScheduledRealtimeRecoveryState: Equatable { + case cancelled + case waitingForStartToFinish + case ready + } + + private static let realtimeStableSessionSeconds: TimeInterval = 30 + private static let realtimeRestartDelaysNanoseconds: [UInt64] = [500_000_000, 2_000_000_000] + + static func realtimeRestartAttempt( + previousRapidRestarts: Int, + activeDuration: TimeInterval) -> Int + { + activeDuration >= self.realtimeStableSessionSeconds ? 1 : previousRapidRestarts + 1 + } + + static func realtimeRestartDelayNanoseconds(attempt: Int) -> UInt64? { + guard attempt > 0, attempt <= self.realtimeRestartDelaysNanoseconds.count else { return nil } + return self.realtimeRestartDelaysNanoseconds[attempt - 1] + } + + func stop( + reconfigurationGeneration expectedReconfigurationGeneration: UInt64?, + lifecycleGeneration expectedLifecycleGeneration: Int?) async + { + guard self.ownsReconfiguration( + expectedReconfigurationGeneration, + lifecycleGeneration: expectedLifecycleGeneration) + else { return } + self.pendingRealtimeRelayStartLifecycleGeneration = nil + self.resetRealtimeRecoveryState() + self.realtimeRelayGeneration &+= 1 + self.realtimeRelayStartGeneration = nil + let realtimeSession = self.detachResourcesForRealtimeStop() + + if let realtimeSession { + await MainActor.run { realtimeSession.stop() } + } + guard self.ownsReconfiguration( + expectedReconfigurationGeneration, + lifecycleGeneration: expectedLifecycleGeneration) + else { return } + await stopSpeaking( + reason: .manual, + reconfigurationGeneration: expectedReconfigurationGeneration, + lifecycleGeneration: expectedLifecycleGeneration) + guard self.ownsReconfiguration( + expectedReconfigurationGeneration, + lifecycleGeneration: expectedLifecycleGeneration) + else { return } + let projectionGeneration = self.realtimeRelayGeneration + _ = await MainActor.run { + self.realtimeRelayDeliveryGate.deliver(ifActive: projectionGeneration) { + TalkModeController.shared.updateLevel(0) + TalkModeController.shared.updatePartialTranscript("") + TalkModeController.shared.updatePhase(.idle) + } + } + } + + func ownsReconfiguration( + _ expectedReconfigurationGeneration: UInt64?, + lifecycleGeneration expectedLifecycleGeneration: Int?) -> Bool + { + (expectedReconfigurationGeneration.map { $0 == self.realtimeReconfigurationGeneration } ?? true) && + (expectedLifecycleGeneration.map { $0 == self.lifecycleGeneration } ?? true) + } + + func beginRealtimeReconfiguration() -> (generation: UInt64, lifecycleGeneration: Int) { + self.lifecycleGeneration &+= 1 + self.realtimeReconfigurationGeneration &+= 1 + return (self.realtimeReconfigurationGeneration, self.lifecycleGeneration) + } + + func realtimeRelayPreferenceDidChange() async { + guard isEnabled else { return } + let reconfiguration = self.beginRealtimeReconfiguration() + await self.stop( + reconfigurationGeneration: reconfiguration.generation, + lifecycleGeneration: reconfiguration.lifecycleGeneration) + guard self.isCurrent(reconfiguration.lifecycleGeneration) else { return } + await self.start() + } + + func start() async { + let gen = lifecycleGeneration + #if DEBUG + guard self.voiceWakeSupportedProvider() else { return } + let hasVoiceWakePermission = await self.voiceWakePermissionProvider() + #else + guard voiceWakeSupported else { return } + let hasVoiceWakePermission = await PermissionManager.ensureVoiceWakePermissions(interactive: true) + #endif + guard hasVoiceWakePermission else { + logger.error("talk runtime not starting: permissions missing") + return + } + macOSRealtimeRelayOptIn = await MainActor.run { + AppStateStore.shared.talkRealtimeRelayEnabled + } + let bypassRealtime = bypassRealtimeOnNextStart + bypassRealtimeOnNextStart = false + if !macOSRealtimeRelayOptIn || bypassRealtime { + await reloadConfig() + } + guard isCurrent(gen) else { return } + if isPaused { + phase = .idle + await MainActor.run { + TalkModeController.shared.updateLevel(0) + TalkModeController.shared.updatePhase(.idle) + } + return + } + var nativeFallbackStatus: String? + if self.macOSRealtimeRelayOptIn, !bypassRealtime { + let fallbackRecognitionGeneration = recognitionGeneration + let fallbackRealtimeRelayGeneration = realtimeRelayGeneration &+ 1 + do { + try await self.startRealtimeRelay(generation: gen) + return + } catch is CancellationError { + if self.consumePendingRealtimeRelayStart() { + await self.start() + } + return + } catch { + pendingRealtimeRelayStartLifecycleGeneration = nil + guard isCurrent(gen), !isPaused, + recognitionGeneration == fallbackRecognitionGeneration, + realtimeRelayGeneration == fallbackRealtimeRelayGeneration, + realtimeRelayStartGeneration == nil, + realtimeSession == nil + else { return } + logger.error( + "talk realtime unavailable; using native fallback: " + + "\(error.localizedDescription, privacy: .public)") + nativeFallbackStatus = String(localized: "Realtime unavailable — using native speech") + } + } + await self.startNativeFallback(generation: gen, status: nativeFallbackStatus) + } + + func consumePendingRealtimeRelayStart() -> Bool { + guard let generation = pendingRealtimeRelayStartLifecycleGeneration else { return false } + pendingRealtimeRelayStartLifecycleGeneration = nil + return generation == lifecycleGeneration && isEnabled && !isPaused && + realtimeSession == nil && realtimeRelayStartGeneration == nil && + self.macOSRealtimeRelayOptIn + } + + private func startNativeFallback(generation: Int, status: String? = nil) async { + let relayGeneration = realtimeRelayGeneration + guard await startRecognition(lifecycleGeneration: generation), + await self.commitNativeFallback( + lifecycleGeneration: generation, + recognitionGeneration: recognitionGeneration, + relayGeneration: relayGeneration, + status: status) + else { return } + startAudioInputObserver() + startSilenceMonitor() + } + + func commitNativeFallback( + lifecycleGeneration: Int, + recognitionGeneration: Int, + relayGeneration: UInt64, + status: String?) async -> Bool + { + let ownsFallback = { + self.canCommitRecognitionStart( + lifecycleGeneration: lifecycleGeneration, + recognitionAttempt: recognitionGeneration) && + self.realtimeRelayGeneration == relayGeneration && + self.realtimeRelayStartGeneration == nil && self.realtimeSession == nil + } + guard ownsFallback() else { return false } + phase = .listening + return await self.projectRealtimeRelay(relayGeneration, nil) { + if let status { + TalkModeController.shared.updatePartialTranscript(status) + } + TalkModeController.shared.updatePhase(.listening) + } + } + + func inputDeviceSelectionDidChange() async { + if let realtimeSession { + guard isEnabled, !isPaused else { return } + let relayGeneration = realtimeRelayGeneration + do { + try await MainActor.run { + try realtimeSession.setInputPaused(true) + try realtimeSession.setInputPaused(false) + } + } catch { + logger.error( + "talk realtime input restart failed: \(error.localizedDescription, privacy: .public)") + await self.handleRealtimeInputRestartFailure( + error.localizedDescription, + relayGeneration: relayGeneration) + } + return + } + guard isEnabled, !isPaused, phase == .listening else { return } + logger.info("talk input selection changed; restarting capture") + let lifecycleGeneration = self.lifecycleGeneration + _ = await startRecognition(lifecycleGeneration: lifecycleGeneration) + } + + func shouldAttemptRealtimeRelay() -> Bool { + guard macOSRealtimeRelayOptIn else { + logger.debug("talk macOS realtime relay disabled locally; using native speech") + return false + } + guard hasGatewayRealtimeRelayTuple else { + logger.warning( + "talk macOS realtime relay opted in but Gateway tuple is incompatible: " + + "mode=\(realtimeMode ?? "missing", privacy: .public) " + + "transport=\(realtimeTransport ?? "missing", privacy: .public) " + + "brain=\(realtimeBrain ?? "missing", privacy: .public); using native fallback") + return false + } + return Self.shouldUseRealtimeRelay( + localOptIn: macOSRealtimeRelayOptIn, + hasGatewayRealtimeRelayTuple: hasGatewayRealtimeRelayTuple) + } + + static func shouldUseRealtimeRelay( + localOptIn: Bool, + hasGatewayRealtimeRelayTuple: Bool) -> Bool + { + localOptIn && hasGatewayRealtimeRelayTuple + } + + func startRealtimeRelay(generation: Int) async throws { + let relayGeneration = try beginRealtimeRelayStart() + defer { + if self.realtimeRelayStartGeneration == relayGeneration { + self.realtimeRelayStartGeneration = nil + } + } + let bootstrap: GatewayConnection.RealtimeTalkBootstrap + do { + bootstrap = try await self.realtimeTalkBootstrapProvider() + } catch { + try await self.applyRealtimeTalkConfig( + self.fallbackTalkConfig(), + lifecycleGeneration: generation, + relayGeneration: relayGeneration) + throw error + } + guard isCurrent(generation), !isPaused, + realtimeRelayGeneration == relayGeneration + else { throw CancellationError() } + let config = self.parseTalkConfig(bootstrap.configSnapshot) + try await self.applyRealtimeTalkConfig( + config, + lifecycleGeneration: generation, + relayGeneration: relayGeneration) + guard self.shouldAttemptRealtimeRelay() else { + throw RealtimeRelayConfigurationError.incompatible + } + let session = try await makeRealtimeRelaySession( + bootstrap: bootstrap, + lifecycleGeneration: generation, + relayGeneration: relayGeneration) + try await ownAndStartRealtimeSession( + session, + lifecycleGeneration: generation, + relayGeneration: relayGeneration, + start: { session in try await session.start() }) + realtimeSessionReadyAt = Date() + phase = .listening + _ = await self.projectRealtimeRelay(relayGeneration, session) { + TalkModeController.shared.updatePartialTranscript("") + TalkModeController.shared.updatePhase(.listening) + } + logger.info( + "talk realtime ready provider=\(realtimeProvider ?? "default", privacy: .public) " + + "model=\(realtimeModelId ?? "default", privacy: .public)") + } + + func applyRealtimeTalkConfig( + _ config: TalkModeGatewayConfigState, + lifecycleGeneration: Int, + relayGeneration: UInt64) async throws + { + let locale = await MainActor.run { () -> String? in + guard self.realtimeRelayDeliveryGate.deliver(ifActive: relayGeneration, { + AppStateStore.shared.seamColorHex = config.seamColorHex + }) else { return nil } + return AppStateStore.shared.voiceWakeLocaleID + } + #if DEBUG + if let checkpoint = self.realtimeConfigApplicationCheckpoint { + await checkpoint() + } + #endif + guard let locale, + self.isCurrent(lifecycleGeneration), + !self.isPaused, + self.realtimeRelayGeneration == relayGeneration, + self.realtimeRelayStartGeneration == relayGeneration + else { throw CancellationError() } + self.commitTalkConfig(config, locale: locale) + } + + func fallbackTalkConfig() -> TalkModeGatewayConfigState { + let env = ProcessInfo.processInfo.environment + return TalkModeGatewayConfigParser.fallback( + defaultModelIdFallback: Self.defaultModelIdFallback, + defaultSilenceTimeoutMs: Self.defaultSilenceTimeoutMs, + envVoice: env["ELEVENLABS_VOICE_ID"], + sagVoice: env["SAG_VOICE_ID"], + envApiKey: env["ELEVENLABS_API_KEY"]) + } + + private func beginRealtimeRelayStart() throws -> UInt64 { + guard realtimeSession == nil, realtimeRelayStartGeneration == nil else { + throw CancellationError() + } + realtimeRelayGeneration &+= 1 + let relayGeneration = realtimeRelayGeneration + realtimeRelayStartGeneration = relayGeneration + return relayGeneration + } + + private func makeRealtimeRelaySession( + bootstrap: GatewayConnection.RealtimeTalkBootstrap, + lifecycleGeneration: Int, + relayGeneration: UInt64) async throws -> RealtimeTalkRelaySession + { + guard isCurrent(lifecycleGeneration), !isPaused, + realtimeRelayGeneration == relayGeneration + else { throw CancellationError() } + let activeSessionKey = await MainActor.run { + WebChatManager.shared.activeSessionKey + } + let sessionKey: String = if let activeSessionKey { + activeSessionKey + } else { + bootstrap.sessionKey + } + let options = RealtimeTalkRelaySession.Options( + sessionKey: sessionKey, + provider: realtimeProvider, + model: realtimeModelId, + voice: realtimeSpeakerVoice) + return await MainActor.run { + RealtimeTalkRelaySession( + transport: bootstrap.transport, + options: options, + audioCapture: MacRealtimeTalkAudioCapture(), + pcmPlayer: RealtimePCMStreamingAudioPlayer(), + onStatus: { [weak self] status in + Task { await self?.handleRealtimeStatus(status, relayGeneration: relayGeneration) } + }, + onIssue: { [weak self] issue in + Task { await self?.handleRealtimeIssue(issue, relayGeneration: relayGeneration) } + }, + onTermination: { [weak self] termination in + Task { + await self?.handleRealtimeTermination( + termination, + relayGeneration: relayGeneration) + } + }, + onSpeakingChanged: { [weak self] speaking in + Task { + await self?.handleRealtimeSpeakingChanged( + speaking, + relayGeneration: relayGeneration) + } + }, + onInputLevel: { [weak self] level in + Task { await self?.handleRealtimeInputLevel(level, relayGeneration: relayGeneration) } + }, + onOutputLevel: { [weak self] level in + Task { await self?.handleRealtimeOutputLevel(level, relayGeneration: relayGeneration) } + }, + onTranscript: { [weak self] transcript in + Task { + await self?.handleRealtimeTranscript( + transcript, + relayGeneration: relayGeneration) + } + }) + } + } + + private func ownAndStartRealtimeSession( + _ session: RealtimeTalkRelaySession, + lifecycleGeneration: Int, + relayGeneration: UInt64, + start: @MainActor @Sendable (RealtimeTalkRelaySession) async throws -> Void) async throws + { + // Construction crosses executors. Claim ownership only after every lifecycle and + // attempt fact is revalidated, then publish before start can suspend. + guard isCurrent(lifecycleGeneration), !isPaused, + realtimeRelayGeneration == relayGeneration, + realtimeRelayStartGeneration == relayGeneration, + realtimeSession == nil + else { + await MainActor.run { session.stop() } + throw CancellationError() + } + realtimeSession = session + do { + try await start(session) + } catch { + await MainActor.run { session.stop() } + if realtimeSession === session { + realtimeSession = nil + } + guard isCurrent(lifecycleGeneration), !isPaused, + realtimeRelayGeneration == relayGeneration, + realtimeRelayStartGeneration == relayGeneration + else { + throw CancellationError() + } + throw error + } + guard isCurrent(lifecycleGeneration), !isPaused, + realtimeRelayGeneration == relayGeneration, + realtimeSession === session + else { + await MainActor.run { session.stop() } + if realtimeSession === session { + realtimeSession = nil + } + throw CancellationError() + } + } + + private func handleRealtimeStatus(_ status: String, relayGeneration: UInt64) { + guard let session = realtimeSession, + ownsRealtimeRelay(relayGeneration, session) + else { return } + logger.debug("talk realtime status=\(status, privacy: .public)") + } + + private func handleRealtimeIssue(_ issue: RealtimeTalkRelayIssue, relayGeneration: UInt64) async { + guard let session = realtimeSession, + ownsRealtimeRelay(relayGeneration, session) + else { return } + logger.error( + "talk realtime issue code=\(issue.code, privacy: .public) " + + "message=\(issue.message, privacy: .public)") + _ = await self.projectRealtimeRelay(relayGeneration, session) { TalkModeController.shared + .updatePartialTranscript(issue.message) + } + } + + func handleRealtimeInputRestartFailure( + _ message: String, + relayGeneration: UInt64) async + { + let issue = RealtimeTalkRelayIssue( + code: "audio_input_unavailable", + message: message, + provider: realtimeProvider, + model: realtimeModelId, + transport: "gateway-relay", + phase: "audio-input") + await handleRealtimeIssue(issue, relayGeneration: relayGeneration) + await handleRealtimeTermination( + .audioInputFailed(message: issue.message), + relayGeneration: relayGeneration) + } + + func setRealtimeInputPaused( + _ paused: Bool, + session: RealtimeTalkRelaySession, + relayGeneration: UInt64) async -> Bool + { + do { + try await MainActor.run { + try session.setInputPaused(paused) + } + return true + } catch { + logger.error( + "talk realtime pause transition failed: \(error.localizedDescription, privacy: .public)") + await self.handleRealtimeInputRestartFailure( + error.localizedDescription, + relayGeneration: relayGeneration) + return false + } + } + + func handleRealtimeTermination( + _ termination: RealtimeTalkRelayTermination, + relayGeneration: UInt64) async + { + guard let session = realtimeSession, + ownsRealtimeRelay(relayGeneration, session) + else { return } + logger.warning( + "talk realtime terminated=\(String(describing: termination), privacy: .public)") + let activeDuration = realtimeSessionReadyAt.map { Date().timeIntervalSince($0) } ?? 0 + realtimeRelayGeneration &+= 1 + let terminalGeneration = realtimeRelayGeneration + // Session-owned terminations close before signalling; runtime-initiated ones do not. + // stop() is idempotent, so closing here keeps a dead relay and its event subscription + // from outliving their owner while recovery starts a replacement session. + await MainActor.run { session.stop() } + guard self.ownsRealtimeRelay(terminalGeneration, session) else { return } + realtimeSession = nil + realtimeSessionReadyAt = nil + phase = .idle + let shouldRecover = isEnabled && !isPaused + let attempt = Self.realtimeRestartAttempt( + previousRapidRestarts: rapidRealtimeRestartCount, + activeDuration: activeDuration) + let delay = Self.realtimeRestartDelayNanoseconds(attempt: attempt) + guard await self.projectRealtimeRelay(terminalGeneration, nil, { + TalkModeController.shared.updateLevel(0) + TalkModeController.shared.updateSpeakingLevel(nil) + TalkModeController.shared.updatePhase(.idle) + if shouldRecover { + TalkModeController.shared.updatePartialTranscript(delay == nil + ? String(localized: "Realtime disconnected repeatedly — using native speech") + : String(localized: "Realtime disconnected — reconnecting…")) + } + }) else { return } + let lifecycleGeneration = self.lifecycleGeneration + let restartGeneration = realtimeRestartGeneration &+ 1 + guard shouldRecover, self.ownsRealtimeRelay(terminalGeneration, nil), + isEnabled, !isPaused + else { return } + rapidRealtimeRestartCount = attempt + realtimeRestartGeneration = restartGeneration + bypassRealtimeOnNextStart = delay == nil + self.scheduleRealtimeRecovery( + after: delay, + lifecycleGeneration: lifecycleGeneration, + restartGeneration: restartGeneration) + } + + func handleRealtimeSpeakingChanged(_ speaking: Bool, relayGeneration: UInt64) async { + guard let session = realtimeSession, + ownsRealtimeRelay(relayGeneration, session), + isEnabled, + !self.isPaused + else { return } + if speaking { + phase = .speaking + _ = await self.projectRealtimeRelay(relayGeneration, session) { + TalkModeController.shared.updatePhase(.speaking) + } + } else if !isPaused { + phase = .listening + _ = await self.projectRealtimeRelay(relayGeneration, session) { + TalkModeController.shared.updatePhase(.listening) + } + } + } + + func handleRealtimeInputLevel(_ level: Double, relayGeneration: UInt64) async { + guard let session = realtimeSession, + ownsRealtimeRelay(relayGeneration, session), + isEnabled, + !self.isPaused + else { return } + _ = await self.projectRealtimeRelay(relayGeneration, session) { + TalkModeController.shared.updateLevel(level) + } + } + + func handleRealtimeOutputLevel(_ level: Double?, relayGeneration: UInt64) async { + guard let session = realtimeSession, + ownsRealtimeRelay(relayGeneration, session), + isEnabled, + !self.isPaused + else { return } + _ = await self.projectRealtimeRelay(relayGeneration, session) { + TalkModeController.shared.updateSpeakingLevel(level) + } + } + + func handleRealtimeTranscript( + _ transcript: RealtimeTalkTranscript, + relayGeneration: UInt64) async + { + guard let session = realtimeSession, + ownsRealtimeRelay(relayGeneration, session), + isEnabled, + !self.isPaused + else { return } + let text = transcript.text.trimmingCharacters(in: .whitespacesAndNewlines) + guard !text.isEmpty else { return } + guard transcript.role == "user" else { return } + if transcript.isFinal { + phase = .thinking + _ = await self.projectRealtimeRelay(relayGeneration, session) { + TalkModeController.shared.commitTranscript(text) + TalkModeController.shared.updatePhase(.thinking) + } + } else { + _ = await self.projectRealtimeRelay(relayGeneration, session) { + TalkModeController.shared.updatePartialTranscript(text) + } + } + } + + func ownsRealtimeRelay(_ generation: UInt64, _ session: RealtimeTalkRelaySession?) -> Bool { + realtimeRelayGeneration == generation && realtimeSession === session + } + + func projectRealtimeRelay( + _ generation: UInt64, + _ session: RealtimeTalkRelaySession?, + _ body: @escaping @MainActor @Sendable () -> Void) async -> Bool + { + guard self.ownsRealtimeRelay(generation, session) else { return false } + return await MainActor.run { + self.realtimeRelayDeliveryGate.deliver(ifActive: generation, body) + } + } + + func resetRealtimeRecoveryState() { + self.cancelScheduledRealtimeRecovery() + realtimeSessionReadyAt = nil + rapidRealtimeRestartCount = 0 + bypassRealtimeOnNextStart = false + } + + func cancelScheduledRealtimeRecovery() { + realtimeRestartGeneration &+= 1 + realtimeRestartTask?.cancel() + realtimeRestartTask = nil + } + + private func scheduleRealtimeRecovery( + after delayNanoseconds: UInt64?, + lifecycleGeneration: Int, + restartGeneration: UInt64) + { + realtimeRestartTask?.cancel() + realtimeRestartTask = Task { [weak self] in + if let delayNanoseconds { + do { + try await Task.sleep(nanoseconds: delayNanoseconds) + } catch { + return + } + } + while let self { + switch await self.scheduledRealtimeRecoveryState( + lifecycleGeneration: lifecycleGeneration, + restartGeneration: restartGeneration) + { + case .cancelled: + return + case .waitingForStartToFinish: + do { + try await Task.sleep(nanoseconds: 50_000_000) + } catch { + return + } + case .ready: + await self.performScheduledRealtimeRecovery( + lifecycleGeneration: lifecycleGeneration, + restartGeneration: restartGeneration) + return + } + } + } + } + + private func scheduledRealtimeRecoveryState( + lifecycleGeneration: Int, + restartGeneration: UInt64) -> ScheduledRealtimeRecoveryState + { + guard self.lifecycleGeneration == lifecycleGeneration, + realtimeRestartGeneration == restartGeneration, + isEnabled, + !isPaused, + realtimeSession == nil + else { return .cancelled } + return realtimeRelayStartGeneration == nil ? .ready : .waitingForStartToFinish + } + + private func performScheduledRealtimeRecovery( + lifecycleGeneration: Int, + restartGeneration: UInt64) async + { + guard self.scheduledRealtimeRecoveryState( + lifecycleGeneration: lifecycleGeneration, + restartGeneration: restartGeneration) == .ready + else { return } + realtimeRestartTask = nil + await self.start() + } +} + +#if DEBUG +extension TalkModeRuntime { + func _test_beginRecognitionAttempt(lifecycleGeneration: Int) -> Int? { + self.beginRecognitionAttempt(lifecycleGeneration: lifecycleGeneration) + } + + func _test_setRecognitionCleanupProbe(_ probe: (@Sendable () -> Void)?) { + self.recognitionCleanupProbe = probe + } + + func _test_setVoiceWakeReadiness(supported: Bool, permissionGranted: Bool) { + self.voiceWakeSupportedProvider = { supported } + self.voiceWakePermissionProvider = { permissionGranted } + } + + func _test_setRealtimeConfigApplicationCheckpoint( + _ checkpoint: (@Sendable () async -> Void)?) + { + self.realtimeConfigApplicationCheckpoint = checkpoint + } + + func _test_enableRealtimeRelaySelection() { + (macOSRealtimeRelayOptIn, hasGatewayRealtimeRelayTuple) = (true, true) + } + + func _test_prepareEnabledLifecycle() -> Int { + isEnabled = true + isPaused = false + lifecycleGeneration &+= 1 + return lifecycleGeneration + } + + func _test_prepareEnabledRealtimeSessionForClose( + _ session: RealtimeTalkRelaySession) -> UInt64 + { + self.cancelScheduledRealtimeRecovery() + isEnabled = true + isPaused = false + lifecycleGeneration &+= 1 + realtimeRelayGeneration &+= 1 + realtimeSession = session + realtimeSessionReadyAt = nil + rapidRealtimeRestartCount = 0 + bypassRealtimeOnNextStart = false + return realtimeRelayGeneration + } +} +#endif diff --git a/apps/macos/Sources/OpenClaw/TalkModeRuntime.swift b/apps/macos/Sources/OpenClaw/TalkModeRuntime.swift index 0cc3846754c7..cd22661ad61c 100644 --- a/apps/macos/Sources/OpenClaw/TalkModeRuntime.swift +++ b/apps/macos/Sources/OpenClaw/TalkModeRuntime.swift @@ -9,6 +9,8 @@ import Speech actor TalkModeRuntime { static let shared = TalkModeRuntime() + typealias RealtimeTalkBootstrapProvider = + @Sendable () async throws -> GatewayConnection.RealtimeTalkBootstrap enum PlaybackPlan: Equatable { case elevenLabsThenSystemVoice(apiKey: String, voiceId: String) @@ -22,13 +24,13 @@ actor TalkModeRuntime { case fallback } - private let logger = Logger(subsystem: "ai.openclaw", category: "talk.runtime") - private let ttsLogger = Logger(subsystem: "ai.openclaw", category: "talk.tts") - private static let defaultModelIdFallback = "eleven_v3" - private static let defaultTalkProvider = "elevenlabs" - private static let mlxTalkProvider = "mlx" - private static let systemTalkProvider = "system" - private static let defaultSilenceTimeoutMs = TalkDefaults.silenceTimeoutMs + let logger = Logger(subsystem: "ai.openclaw", category: "talk.runtime") + let ttsLogger = Logger(subsystem: "ai.openclaw", category: "talk.tts") + static let defaultModelIdFallback = "eleven_v3" + static let defaultTalkProvider = "elevenlabs" + static let mlxTalkProvider = "mlx" + static let systemTalkProvider = "system" + static let defaultSilenceTimeoutMs = TalkDefaults.silenceTimeoutMs private final class RMSMeter: @unchecked Sendable { private let lock = NSLock() @@ -54,16 +56,16 @@ actor TalkModeRuntime { private var activeInputResolution: AudioInputDeviceResolution? private var recognitionRequest: SFSpeechAudioBufferRecognitionRequest? private var recognitionTask: SFSpeechRecognitionTask? - private var recognitionGeneration: Int = 0 + var recognitionGeneration: Int = 0 private var rmsTask: Task? private let rmsMeter = RMSMeter() private var captureTask: Task? private var silenceTask: Task? - private var phase: TalkModePhase = .idle - private var isEnabled = false - private var isPaused = false - private var lifecycleGeneration: Int = 0 + var phase: TalkModePhase = .idle + var isEnabled = false + var isPaused = false + var lifecycleGeneration: Int = 0 private var lastHeard: Date? private var noiseFloorRMS: Double = 1e-4 @@ -79,6 +81,38 @@ actor TalkModeRuntime { private var defaultOutputFormat: String? private var interruptOnSpeech: Bool = true private var activeTalkProvider = TalkModeRuntime.defaultTalkProvider + var realtimeProvider: String? + var realtimeModelId: String? + var realtimeSpeakerVoice: String? + var realtimeMode: String? + var realtimeTransport: String? + var realtimeBrain: String? + var hasGatewayRealtimeRelayTuple = false + var macOSRealtimeRelayOptIn = false + var realtimeSession: RealtimeTalkRelaySession? + var realtimeSessionReadyAt: Date? + var rapidRealtimeRestartCount = 0 + var bypassRealtimeOnNextStart = false + let realtimeRelayDeliveryGate = TalkGenerationDeliveryGate() + var realtimeRelayGeneration: UInt64 = 0 { + didSet { _ = self.realtimeRelayDeliveryGate.activate() } + } + + var realtimeRelayStartGeneration: UInt64? + var pendingRealtimeRelayStartLifecycleGeneration: Int? + var realtimeRestartGeneration: UInt64 = 0 + var realtimeRestartTask: Task? + var realtimeReconfigurationGeneration: UInt64 = 0 + let realtimeTalkBootstrapProvider: RealtimeTalkBootstrapProvider + #if DEBUG + var voiceWakeSupportedProvider: @Sendable () -> Bool = { voiceWakeSupported } + var voiceWakePermissionProvider: VoiceWakePermissionProvider = { + await PermissionManager.ensureVoiceWakePermissions(interactive: true) + } + + var realtimeConfigApplicationCheckpoint: (@Sendable () async -> Void)? + var recognitionCleanupProbe: (@Sendable () -> Void)? + #endif private var speechLocaleID: String? private var lastInterruptedAtSeconds: Double? private var voiceAliases: [String: String] = [:] @@ -93,8 +127,10 @@ actor TalkModeRuntime { private let minSpeechRMS: Double = 1e-3 private let speechBoostFactor: Double = 6.0 - static func configureRecognitionRequest(_ request: SFSpeechAudioBufferRecognitionRequest) { - SpeechRecognitionRequestPolicy.configureInteractiveTranscription(request) + init(realtimeTalkBootstrapProvider: @escaping RealtimeTalkBootstrapProvider = { + try await GatewayConnection.shared.acquireRealtimeTalkBootstrap() + }) { + self.realtimeTalkBootstrapProvider = realtimeTalkBootstrapProvider } // MARK: - Lifecycle @@ -103,8 +139,9 @@ actor TalkModeRuntime { guard enabled != self.isEnabled else { return } self.isEnabled = enabled self.lifecycleGeneration &+= 1 + resetRealtimeRecoveryState() if enabled { - await self.start() + await start() } else { await self.stop() } @@ -117,83 +154,117 @@ actor TalkModeRuntime { guard self.isEnabled else { return } + if paused { + self.pendingRealtimeRelayStartLifecycleGeneration = nil + if self.realtimeRelayStartGeneration != nil { + self.realtimeRelayGeneration &+= 1 + } + } else if self.realtimeRelayStartGeneration != nil, self.macOSRealtimeRelayOptIn { + self.pendingRealtimeRelayStartLifecycleGeneration = self.lifecycleGeneration + return + } + + if paused, realtimeSession == nil { + cancelScheduledRealtimeRecovery() + } + + if let realtimeSession { + let relayGeneration = self.realtimeRelayGeneration + if paused { + await MainActor.run { realtimeSession.setOutputPaused(true) } + guard self.isPaused, + self.realtimeRelayGeneration == relayGeneration, + self.realtimeSession === realtimeSession, + await setRealtimeInputPaused( + true, + session: realtimeSession, + relayGeneration: relayGeneration) + else { return } + self.lastTranscript = "" + self.lastHeard = nil + self.lastSpeechEnergyAt = nil + self.phase = .idle + _ = await projectRealtimeRelay(relayGeneration, realtimeSession) { + TalkModeController.shared.updateLevel(0) + TalkModeController.shared.updateSpeakingLevel(nil) + TalkModeController.shared.updatePartialTranscript("") + TalkModeController.shared.updatePhase(.idle) + } + } else { + guard await setRealtimeInputPaused( + false, + session: realtimeSession, + relayGeneration: relayGeneration), + !self.isPaused, + self.realtimeRelayGeneration == relayGeneration, + self.realtimeSession === realtimeSession + else { return } + await MainActor.run { realtimeSession.setOutputPaused(false) } + guard !self.isPaused, + self.realtimeRelayGeneration == relayGeneration, + self.realtimeSession === realtimeSession + else { + await MainActor.run { realtimeSession.setOutputPaused(true) } + return + } + self.phase = .listening + _ = await projectRealtimeRelay(relayGeneration, realtimeSession) { + TalkModeController.shared.updatePhase(.listening) + } + } + return + } + + if !paused, self.macOSRealtimeRelayOptIn { + await start() + return + } + if paused { self.lastTranscript = "" self.lastHeard = nil self.lastSpeechEnergyAt = nil await MainActor.run { TalkModeController.shared.updatePartialTranscript("") } - await self.stopRecognition() + self.stopRecognition() return } if self.phase == .idle || self.phase == .listening { - await self.startRecognition() + let lifecycleGeneration = self.lifecycleGeneration + guard await self.startRecognition(lifecycleGeneration: lifecycleGeneration), + self.isCurrent(lifecycleGeneration), !self.isPaused else { return } self.phase = .listening await MainActor.run { TalkModeController.shared.updatePhase(.listening) } self.startSilenceMonitor() } } - private func isCurrent(_ generation: Int) -> Bool { + func isCurrent(_ generation: Int) -> Bool { generation == self.lifecycleGeneration && self.isEnabled } - private func start() async { - let gen = self.lifecycleGeneration - guard voiceWakeSupported else { return } - - guard await PermissionManager.ensureVoiceWakePermissions(interactive: true) else { - self.logger.error("talk runtime not starting: permissions missing") - return - } - self.startAudioInputObserver() - await reloadConfig() - guard self.isCurrent(gen) else { return } - if self.isPaused { - self.phase = .idle - await MainActor.run { - TalkModeController.shared.updateLevel(0) - TalkModeController.shared.updatePhase(.idle) - } - return - } - await self.startRecognition() - guard self.isCurrent(gen) else { return } - self.phase = .listening - await MainActor.run { TalkModeController.shared.updatePhase(.listening) } - self.startSilenceMonitor() + func stop() async { + await self.stop(reconfigurationGeneration: nil, lifecycleGeneration: nil) } - private func stop() async { + func detachResourcesForRealtimeStop() -> RealtimeTalkRelaySession? { + let realtimeSession = self.realtimeSession + self.realtimeSession = nil self.audioInputObserver?.stop() self.audioInputObserver = nil self.captureTask?.cancel() self.captureTask = nil self.silenceTask?.cancel() self.silenceTask = nil - - // Stop audio before changing phase (stopSpeaking is gated on .speaking). - await stopSpeaking(reason: .manual) - self.lastTranscript = "" self.lastHeard = nil self.lastSpeechEnergyAt = nil self.phase = .idle - await self.stopRecognition() - await MainActor.run { - TalkModeController.shared.updateLevel(0) - TalkModeController.shared.updatePartialTranscript("") - TalkModeController.shared.updatePhase(.idle) - } + self.stopRecognition() + return realtimeSession } - func inputDeviceSelectionDidChange() async { - guard self.isEnabled, !self.isPaused, self.phase == .listening else { return } - self.logger.info("talk input selection changed; restarting capture") - await self.startRecognition() - } - - private func startAudioInputObserver() { + func startAudioInputObserver() { guard self.audioInputObserver == nil else { return } let observer = AudioInputDeviceObserver() observer.start { @@ -204,9 +275,10 @@ actor TalkModeRuntime { private func audioInputDevicesDidChange() async { guard self.isEnabled, !self.isPaused, self.phase == .listening else { return } + let lifecycleGeneration = self.lifecycleGeneration let availableUIDs = AudioInputDeviceObserver.aliveInputDeviceUIDs() guard let activeInputResolution else { - await self.startRecognition() + _ = await self.startRecognition(lifecycleGeneration: lifecycleGeneration) return } guard activeInputResolution.shouldRestart( @@ -215,7 +287,7 @@ actor TalkModeRuntime { else { return } self.logger.warning("talk active/default input changed; restarting capture") - await self.startRecognition() + _ = await self.startRecognition(lifecycleGeneration: lifecycleGeneration) } // MARK: - Speech recognition @@ -228,12 +300,13 @@ actor TalkModeRuntime { let generation: Int } - private func startRecognition() async { - await self.stopRecognition() - self.recognitionGeneration &+= 1 - let generation = self.recognitionGeneration + func startRecognition(lifecycleGeneration: Int) async -> Bool { + guard let recognitionAttempt = beginRecognitionAttempt( + lifecycleGeneration: lifecycleGeneration) + else { return false } let voiceWakeLocale = await MainActor.run { AppStateStore.shared.voiceWakeLocaleID } + let selectedInputUID = await MainActor.run { AppStateStore.shared.voiceWakeMicID } let supportedLocaleIDs = Set(SFSpeechRecognizer.supportedLocales().map(\.identifier)) let localeID = TalkConfigParsing.resolvedSpeechRecognitionLocaleID( preferredLocaleIDs: [ @@ -242,71 +315,97 @@ actor TalkModeRuntime { Locale.autoupdatingCurrent.identifier, ], supportedLocaleIDs: supportedLocaleIDs) - self.recognizer = localeID + let recognizer = localeID .map { SFSpeechRecognizer(locale: Locale(identifier: $0)) } ?? SFSpeechRecognizer() guard let recognizer, recognizer.isAvailable else { self.logger.error("talk recognizer unavailable") - return + return false } self.logger.debug("talk recognizer locale=\(recognizer.locale.identifier, privacy: .public)") - let request = SFSpeechAudioBufferRecognitionRequest() - Self.configureRecognitionRequest(request) - self.recognitionRequest = request - - let selectedInputUID = await MainActor.run { AppStateStore.shared.voiceWakeMicID } let selection = AudioInputDeviceObserver.resolveSelection(selectedInputUID) + // Preparation above crosses MainActor. Hardware can become owned/running only + // after this lifecycle and recognition attempt are revalidated together. + guard self.canCommitRecognitionStart( + lifecycleGeneration: lifecycleGeneration, + recognitionAttempt: recognitionAttempt) else { return false } + // AVAudioEngine materializes inputNode from the system default before CurrentDevice can bind. // Without a usable default, accessing inputNode can SIGABRT even when another UID is alive. guard selection.resolvedUID != nil, AudioInputDeviceObserver.hasUsableDefaultInputDevice() else { - self.audioEngine = nil - self.activeInputResolution = nil self.logger.error("talk mode: no usable audio input device") - return + return false } - do { - try self.configureAudioEngine( - selection: selection, - request: request, - enableVoiceProcessing: true) - } catch { - self.logger.warning( - "talk processed input setup failed; retrying without voice processing: " + - "\(error.localizedDescription, privacy: .public)") - self.discardAudioEngine() - do { - try self.configureAudioEngine( + let started = TalkRecognitionCaptureLifecycle.start( + isCurrent: { + self.canCommitRecognitionStart( + lifecycleGeneration: lifecycleGeneration, + recognitionAttempt: recognitionAttempt) + }, + prepare: { enableVoiceProcessing in + try self.prepareStartedRecognitionCapture( selection: selection, - request: request, - enableVoiceProcessing: false) - } catch { - self.discardAudioEngine() - self.logger.error( - "talk audio engine start failed: \(error.localizedDescription, privacy: .public)") - return - } - } - + enableVoiceProcessing: enableVoiceProcessing) + }, + discard: { $0.discard() }, + publish: { preparedCapture in + self.recognizer = recognizer + self.recognitionRequest = preparedCapture.request + self.audioEngine = preparedCapture.engine + self.activeInputResolution = preparedCapture.activeInputResolution + self.recognitionTask = recognizer.recognitionTask( + with: preparedCapture.request, + resultHandler: { [weak self, recognitionAttempt] result, error in + guard let self else { return } + let segments = result?.bestTranscription.segments ?? [] + let transcript = result?.bestTranscription.formattedString + let update = RecognitionUpdate( + transcript: transcript, + hasConfidence: segments.contains { $0.confidence > 0.6 }, + isFinal: result?.isFinal ?? false, + errorDescription: error?.localizedDescription, + generation: recognitionAttempt) + Task { await self.handleRecognition(update) } + }) + }, + onFailure: { enableVoiceProcessing, error in + if enableVoiceProcessing { + self.logger.warning( + "talk processed input start failed; retrying without voice processing: " + + "\(error.localizedDescription, privacy: .public)") + } else { + self.logger.error( + "talk audio engine start failed: \(error.localizedDescription, privacy: .public)") + } + }) + guard started else { return false } self.startRMSTicker(meter: self.rmsMeter) - - self.recognitionTask = recognizer.recognitionTask(with: request) { [weak self, generation] result, error in - guard let self else { return } - let segments = result?.bestTranscription.segments ?? [] - let transcript = result?.bestTranscription.formattedString - let update = RecognitionUpdate( - transcript: transcript, - hasConfidence: segments.contains { $0.confidence > 0.6 }, - isFinal: result?.isFinal ?? false, - errorDescription: error?.localizedDescription, - generation: generation) - Task { await self.handleRecognition(update) } - } + return true } - private func stopRecognition() async { + func beginRecognitionAttempt(lifecycleGeneration: Int) -> Int? { + guard self.isCurrent(lifecycleGeneration), !self.isPaused else { return nil } self.recognitionGeneration &+= 1 + let recognitionAttempt = self.recognitionGeneration + self.discardRecognitionResources() + return recognitionAttempt + } + + func canCommitRecognitionStart(lifecycleGeneration: Int, recognitionAttempt: Int) -> Bool { + self.isCurrent(lifecycleGeneration) && !self.isPaused && self.recognitionGeneration == recognitionAttempt + } + + private func stopRecognition() { + self.recognitionGeneration &+= 1 + self.discardRecognitionResources() + } + + func discardRecognitionResources() { + #if DEBUG + self.recognitionCleanupProbe?() + #endif self.recognitionTask?.cancel() self.recognitionTask = nil self.recognitionRequest?.endAudio() @@ -368,7 +467,7 @@ actor TalkModeRuntime { // MARK: - Silence handling - private func startSilenceMonitor() { + func startSilenceMonitor() { self.silenceTask?.cancel() self.silenceTask = Task { [weak self] in await self?.silenceLoop() @@ -417,7 +516,7 @@ actor TalkModeRuntime { if sendChime != .none { await MainActor.run { VoiceWakeChimePlayer.play(sendChime, reason: "talk.send") } } - await self.stopRecognition() + self.stopRecognition() await sendAndSpeak(text) } @@ -451,49 +550,55 @@ actor TalkModeRuntime { return selection } - private func configureAudioEngine( + private func prepareStartedRecognitionCapture( selection: AudioInputDeviceResolution, - request: SFSpeechAudioBufferRecognitionRequest, - enableVoiceProcessing: Bool) throws + enableVoiceProcessing: Bool) + throws -> PreparedRecognitionCapture { + let request = SFSpeechAudioBufferRecognitionRequest() + TalkRecognitionCaptureLifecycle.configure(request) let audioEngine = AVAudioEngine() - self.audioEngine = audioEngine let input = audioEngine.inputNode - if enableVoiceProcessing { - do { + var tapInstalled = false + do { + if enableVoiceProcessing { try input.setVoiceProcessingEnabled(true) - } catch { - // Aggregate devices can reject voice processing; capture still works without it. - self.logger.warning( - "talk voice processing unavailable: \(error.localizedDescription, privacy: .public)") } - } - let activeResolution = self.bindSelectedInputIfNeeded(selection, to: input) - guard activeResolution.resolvedUID != nil else { - throw TalkAudioInputError.unavailable - } + let activeResolution = self.bindSelectedInputIfNeeded(selection, to: input) + guard activeResolution.resolvedUID != nil else { + throw TalkAudioInputError.unavailable + } - let format = input.outputFormat(forBus: 0) - guard format.channelCount > 0, format.sampleRate > 0 else { - throw TalkAudioInputError.invalidFormat + let format = input.outputFormat(forBus: 0) + guard format.channelCount > 0, format.sampleRate > 0 else { + throw TalkAudioInputError.invalidFormat + } + input.removeTap(onBus: 0) + let meter = self.rmsMeter + input.installTap( + onBus: 0, + bufferSize: 2048, + format: format) + { [weak request, meter] buffer, _ in + request?.append(SpeechAudioBufferNormalizer.speechCompatibleBuffer(from: buffer)) + meter.set(TalkAudioLevel.rms(buffer: buffer)) + } + tapInstalled = true + audioEngine.prepare() + try audioEngine.start() + return PreparedRecognitionCapture( + request: request, + engine: audioEngine, + activeInputResolution: activeResolution) + } catch { + request.endAudio() + if tapInstalled { + input.removeTap(onBus: 0) + } + audioEngine.stop() + throw error } - input.removeTap(onBus: 0) - let meter = self.rmsMeter - input.installTap(onBus: 0, bufferSize: 2048, format: format) { [weak request, meter] buffer, _ in - request?.append(SpeechAudioBufferNormalizer.speechCompatibleBuffer(from: buffer)) - meter.set(TalkAudioLevel.rms(buffer: buffer)) - } - audioEngine.prepare() - try audioEngine.start() - self.activeInputResolution = activeResolution - } - - private func discardAudioEngine() { - self.audioEngine?.inputNode.removeTap(onBus: 0) - self.audioEngine?.stop() - self.audioEngine = nil - self.activeInputResolution = nil } private func defaultFallback(for selection: AudioInputDeviceResolution) -> AudioInputDeviceResolution { @@ -504,18 +609,6 @@ actor TalkModeRuntime { } } -private enum TalkAudioInputError: LocalizedError { - case unavailable - case invalidFormat - - var errorDescription: String? { - switch self { - case .unavailable: "Selected input and system default are unavailable" - case .invalidFormat: "Selected audio input has no usable format" - } - } -} - // MARK: - Gateway + TTS extension TalkModeRuntime { @@ -583,8 +676,9 @@ extension TalkModeRuntime { guard let assistantText else { self.logger.warning("talk assistant text missing after timeout") + guard await self.startRecognition(lifecycleGeneration: gen), + self.isCurrent(gen), !self.isPaused else { return } await self.startListening() - await self.startRecognition() return } guard self.isCurrent(gen) else { return } @@ -611,8 +705,10 @@ extension TalkModeRuntime { } return } + let lifecycleGeneration = self.lifecycleGeneration + guard await self.startRecognition(lifecycleGeneration: lifecycleGeneration), + self.isCurrent(lifecycleGeneration), !self.isPaused else { return } await self.startListening() - await self.startRecognition() } private func buildPrompt(transcript: String) -> String { @@ -1113,8 +1209,8 @@ extension TalkModeRuntime { } private func prepareForPlayback(generation: Int) async -> Bool { - await self.startRecognition() - return self.isCurrent(generation) + guard await self.startRecognition(lifecycleGeneration: generation) else { return false } + return self.isCurrent(generation) && !self.isPaused } private func resolveVoiceId(preferred: String?, apiKey: String) async -> String? { @@ -1170,13 +1266,60 @@ extension TalkModeRuntime { return value.allSatisfy { $0.isLetter || $0.isNumber || $0 == "-" || $0 == "_" } } - func stopSpeaking(reason: TalkStopReason) async { + func stopSpeaking( + reason: TalkStopReason, + reconfigurationGeneration expectedReconfigurationGeneration: UInt64? = nil, + lifecycleGeneration expectedLifecycleGeneration: Int? = nil) async + { + guard self.ownsReconfiguration( + expectedReconfigurationGeneration, + lifecycleGeneration: expectedLifecycleGeneration) + else { return } + if let realtimeSession { + let relayReason = switch reason { + case .userTap: "user" + case .speech: "barge-in" + case .manual: "shutdown" + } + let cancelled = await MainActor.run { + realtimeSession.cancelOutput(reason: relayReason) + } + guard self.ownsReconfiguration( + expectedReconfigurationGeneration, + lifecycleGeneration: expectedLifecycleGeneration) + else { return } + if cancelled, reason != .manual, !self.isPaused { + self.phase = .listening + await MainActor.run { TalkModeController.shared.updatePhase(.listening) } + } + return + } let usePCM = self.lastPlaybackWasPCM let remoteInterruptedAt = usePCM ? await stopPCM() : await stopMP3() + guard self.ownsReconfiguration( + expectedReconfigurationGeneration, + lifecycleGeneration: expectedLifecycleGeneration) + else { return } _ = usePCM ? await stopMP3() : await stopPCM() + guard self.ownsReconfiguration( + expectedReconfigurationGeneration, + lifecycleGeneration: expectedLifecycleGeneration) + else { return } let localInterruptedAt = await stopTalkAudio() + guard self.ownsReconfiguration( + expectedReconfigurationGeneration, + lifecycleGeneration: expectedLifecycleGeneration) + else { return } await TalkSystemSpeechSynthesizer.shared.stop() + guard self.ownsReconfiguration( + expectedReconfigurationGeneration, + lifecycleGeneration: expectedLifecycleGeneration) + else { return } await stopMLXVoice() + guard self.ownsReconfiguration( + expectedReconfigurationGeneration, + lifecycleGeneration: expectedLifecycleGeneration) + else { return } guard self.phase == .speaking else { return } let interruptedAt = remoteInterruptedAt ?? localInterruptedAt if reason == .speech, let interruptedAt { @@ -1308,10 +1451,72 @@ extension TalkModeRuntime { await TalkMLXSpeechSynthesizer.shared.cancelCurrent() } + func parseTalkConfig(_ snap: ConfigSnapshot) -> TalkModeGatewayConfigState { + let env = ProcessInfo.processInfo.environment + let envVoice = env["ELEVENLABS_VOICE_ID"]?.trimmingCharacters(in: .whitespacesAndNewlines) + let sagVoice = env["SAG_VOICE_ID"]?.trimmingCharacters(in: .whitespacesAndNewlines) + let envApiKey = env["ELEVENLABS_API_KEY"]?.trimmingCharacters(in: .whitespacesAndNewlines) + + let parsed = TalkModeGatewayConfigParser.parse( + snapshot: snap, + defaultProvider: Self.defaultTalkProvider, + defaultModelIdFallback: Self.defaultModelIdFallback, + defaultSilenceTimeoutMs: Self.defaultSilenceTimeoutMs, + envVoice: envVoice, + sagVoice: sagVoice, + envApiKey: envApiKey) + if parsed.missingResolvedPayload { + self.ttsLogger.info("talk config ignored: normalized payload missing talk.resolved") + } + if parsed.activeProvider == Self.defaultTalkProvider { + self.ttsLogger.info("talk config provider from talk.resolved") + } else if parsed.activeProvider == Self.mlxTalkProvider || + parsed.activeProvider == Self.systemTalkProvider + { + self.ttsLogger.info( + "talk provider \(parsed.activeProvider, privacy: .public) active") + } else { + self.ttsLogger + .info( + """ + talk provider \(parsed.activeProvider, privacy: .public) uses gateway talk.speak \ + with system voice fallback + """) + } + return parsed + } + + private func fetchTalkConfig() async -> TalkModeGatewayConfigState { + do { + let snap: ConfigSnapshot = try await GatewayConnection.shared.requestDecoded( + method: .talkConfig, + params: ["includeSecrets": AnyCodable(true)], + timeoutMs: 8000) + return self.parseTalkConfig(snap) + } catch { + return self.fallbackTalkConfig() + } + } + // MARK: - Config - private func reloadConfig() async { + func reloadConfig() async { let cfg = await fetchTalkConfig() + await self.applyTalkConfig(cfg) + self.macOSRealtimeRelayOptIn = await MainActor.run { + AppStateStore.shared.talkRealtimeRelayEnabled + } + } + + func applyTalkConfig(_ cfg: TalkModeGatewayConfigState) async { + let locale = await MainActor.run { + AppStateStore.shared.seamColorHex = cfg.seamColorHex + return AppStateStore.shared.voiceWakeLocaleID + } + self.commitTalkConfig(cfg, locale: locale) + } + + func commitTalkConfig(_ cfg: TalkModeGatewayConfigState, locale: String) { self.defaultVoiceId = cfg.voiceId self.voiceAliases = cfg.voiceAliases if !self.voiceOverrideActive { @@ -1324,8 +1529,14 @@ extension TalkModeRuntime { self.defaultOutputFormat = cfg.outputFormat self.interruptOnSpeech = cfg.interruptOnSpeech self.activeTalkProvider = cfg.activeProvider + self.realtimeProvider = cfg.realtimeProvider + self.realtimeModelId = cfg.realtimeModelId + self.realtimeSpeakerVoice = cfg.realtimeSpeakerVoice + self.realtimeMode = cfg.realtimeMode + self.realtimeTransport = cfg.realtimeTransport + self.realtimeBrain = cfg.realtimeBrain + self.hasGatewayRealtimeRelayTuple = cfg.hasGatewayRealtimeRelayTuple let configuredSilenceMs = cfg.silenceTimeoutMs - let locale = await MainActor.run { AppStateStore.shared.voiceWakeLocaleID } let isCJKLocale = locale.hasPrefix("ko") || locale.hasPrefix("ja") || locale.hasPrefix("zh") let effectiveSilenceMs = isCJKLocale ? max(configuredSilenceMs, 2000) : configuredSilenceMs if isCJKLocale, configuredSilenceMs < 2000 { @@ -1351,7 +1562,11 @@ extension TalkModeRuntime { "apiKey=\(hasApiKey, privacy: .public) " + "interrupt=\(cfg.interruptOnSpeech, privacy: .public) " + "silenceTimeoutMs=\(cfg.silenceTimeoutMs, privacy: .public) " + - "speechLocale=\(cfg.speechLocaleID ?? "device", privacy: .public)") + "speechLocale=\(cfg.speechLocaleID ?? "device", privacy: .public) " + + "realtimeMode=\(cfg.realtimeMode ?? "off", privacy: .public) " + + "realtimeTransport=\(cfg.realtimeTransport ?? "default", privacy: .public) " + + "realtimeBrain=\(cfg.realtimeBrain ?? "default", privacy: .public) " + + "macOSRealtimeOptIn=\(self.macOSRealtimeRelayOptIn, privacy: .public)") } static func selectTalkProviderConfig( @@ -1364,57 +1579,6 @@ extension TalkModeRuntime { TalkConfigParsing.resolvedSilenceTimeoutMs(talk, fallback: self.defaultSilenceTimeoutMs) } - private func fetchTalkConfig() async -> TalkModeGatewayConfigState { - let env = ProcessInfo.processInfo.environment - let envVoice = env["ELEVENLABS_VOICE_ID"]?.trimmingCharacters(in: .whitespacesAndNewlines) - let sagVoice = env["SAG_VOICE_ID"]?.trimmingCharacters(in: .whitespacesAndNewlines) - let envApiKey = env["ELEVENLABS_API_KEY"]?.trimmingCharacters(in: .whitespacesAndNewlines) - - do { - let snap: ConfigSnapshot = try await GatewayConnection.shared.requestDecoded( - method: .talkConfig, - params: ["includeSecrets": AnyCodable(true)], - timeoutMs: 8000) - let parsed = TalkModeGatewayConfigParser.parse( - snapshot: snap, - defaultProvider: Self.defaultTalkProvider, - defaultModelIdFallback: Self.defaultModelIdFallback, - defaultSilenceTimeoutMs: Self.defaultSilenceTimeoutMs, - envVoice: envVoice, - sagVoice: sagVoice, - envApiKey: envApiKey) - if parsed.missingResolvedPayload { - self.ttsLogger.info("talk config ignored: normalized payload missing talk.resolved") - } - await MainActor.run { - AppStateStore.shared.seamColorHex = parsed.seamColorHex - } - if parsed.activeProvider == Self.defaultTalkProvider { - self.ttsLogger.info("talk config provider from talk.resolved") - } else if parsed.activeProvider == Self.mlxTalkProvider || - parsed.activeProvider == Self.systemTalkProvider - { - self.ttsLogger.info( - "talk provider \(parsed.activeProvider, privacy: .public) active") - } else { - self.ttsLogger - .info( - """ - talk provider \(parsed.activeProvider, privacy: .public) uses gateway talk.speak \ - with system voice fallback - """) - } - return parsed - } catch { - return TalkModeGatewayConfigParser.fallback( - defaultModelIdFallback: Self.defaultModelIdFallback, - defaultSilenceTimeoutMs: Self.defaultSilenceTimeoutMs, - envVoice: envVoice, - sagVoice: sagVoice, - envApiKey: envApiKey) - } - } - // MARK: - Audio level handling private func noteAudioLevel(rms: Double) async { diff --git a/apps/macos/Sources/OpenClaw/TalkRecognitionCaptureLifecycle.swift b/apps/macos/Sources/OpenClaw/TalkRecognitionCaptureLifecycle.swift new file mode 100644 index 000000000000..6a4058cc3b35 --- /dev/null +++ b/apps/macos/Sources/OpenClaw/TalkRecognitionCaptureLifecycle.swift @@ -0,0 +1,57 @@ +@preconcurrency import AVFoundation +import Foundation +import Speech + +struct PreparedRecognitionCapture { + let request: SFSpeechAudioBufferRecognitionRequest + let engine: AVAudioEngine + let activeInputResolution: AudioInputDeviceResolution + + func discard() { + self.request.endAudio() + self.engine.inputNode.removeTap(onBus: 0) + self.engine.stop() + } +} + +enum TalkAudioInputError: LocalizedError { + case unavailable + case invalidFormat + + var errorDescription: String? { + switch self { + case .unavailable: "Selected input and system default are unavailable" + case .invalidFormat: "Selected audio input has no usable format" + } + } +} + +enum TalkRecognitionCaptureLifecycle { + static func configure(_ request: SFSpeechAudioBufferRecognitionRequest) { + SpeechRecognitionRequestPolicy.configureInteractiveTranscription(request) + } + + static func start( + isCurrent: () -> Bool, + prepare: (_ enableVoiceProcessing: Bool) throws -> Capture, + discard: (Capture) -> Void, + publish: (Capture) -> Void, + onFailure: (_ enableVoiceProcessing: Bool, _ error: Error) -> Void) -> Bool + { + for enableVoiceProcessing in [true, false] { + guard isCurrent() else { return false } + do { + let capture = try prepare(enableVoiceProcessing) + guard isCurrent() else { + discard(capture) + return false + } + publish(capture) + return true + } catch { + onFailure(enableVoiceProcessing, error) + } + } + return false + } +} diff --git a/apps/macos/Sources/OpenClaw/VoiceWakeSettings.swift b/apps/macos/Sources/OpenClaw/VoiceWakeSettings.swift index 24823e23250e..195839891861 100644 --- a/apps/macos/Sources/OpenClaw/VoiceWakeSettings.swift +++ b/apps/macos/Sources/OpenClaw/VoiceWakeSettings.swift @@ -182,6 +182,8 @@ struct VoiceWakeSettings: View { subtitle: "Start listening while you hold the key and show the preview overlay.", binding: self.$state.voicePushToTalkEnabled) + self.realtimeRelayToggle + if self.state.voicePushToTalkEnabled, self.state.talkEnabled { SettingsCardRow( title: "Push-to-talk paused", @@ -930,6 +932,18 @@ private struct TriggerPhraseHelpRow: View { } } +extension VoiceWakeSettings { + private var realtimeRelayToggle: some View { + SettingsCardToggleRow( + title: "Use realtime Gateway relay", + subtitle: """ + Use the Gateway's configured realtime voice session on this Mac. \ + Requires realtime, gateway-relay, and agent-consult in Talk settings. + """, + binding: self.$state.talkRealtimeRelayEnabled) + } +} + #if DEBUG struct VoiceWakeSettings_Previews: PreviewProvider { static var previews: some View { diff --git a/apps/macos/Tests/OpenClawIPCTests/TalkModeGatewayConfigTests.swift b/apps/macos/Tests/OpenClawIPCTests/TalkModeGatewayConfigTests.swift index 98dba2610353..fdb06c06d7e7 100644 --- a/apps/macos/Tests/OpenClawIPCTests/TalkModeGatewayConfigTests.swift +++ b/apps/macos/Tests/OpenClawIPCTests/TalkModeGatewayConfigTests.swift @@ -53,4 +53,116 @@ struct TalkModeGatewayConfigTests { #expect(parsed.referenceAudioPath == "/tmp/reference.wav") #expect(parsed.referenceText == "reference transcript") } + + @Test func `realtime config uses top level overrides and normalizes control values`() { + let snapshot = Self.snapshot(talk: [ + "realtime": [ + "provider": " OpenAI ", + "providers": [ + "openai": [ + "model": "provider-model", + "speakerVoice": "alloy", + ], + ], + "model": " gpt-live-1-codex ", + "speakerVoice": " cedar ", + "mode": " Realtime ", + "transport": " Gateway-Relay ", + "brain": " Agent-Consult ", + ], + ]) + + let parsed = Self.parse(snapshot) + + #expect(parsed.realtimeProvider == "OpenAI") + #expect(parsed.realtimeModelId == "gpt-live-1-codex") + #expect(parsed.realtimeSpeakerVoice == "cedar") + #expect(parsed.realtimeMode == "realtime") + #expect(parsed.realtimeTransport == "gateway-relay") + #expect(parsed.realtimeBrain == "agent-consult") + #expect(parsed.hasGatewayRealtimeRelayTuple) + } + + @Test func `realtime config infers its sole provider and reads provider defaults`() { + let snapshot = Self.snapshot(talk: [ + "realtime": [ + "providers": [ + "openai": [ + "model": "gpt-realtime-2.1", + "voice": "marin", + ], + ], + "mode": "realtime", + ], + ]) + + let parsed = Self.parse(snapshot) + + #expect(parsed.realtimeProvider == "openai") + #expect(parsed.realtimeModelId == "gpt-realtime-2.1") + #expect(parsed.realtimeSpeakerVoice == "marin") + #expect(parsed.realtimeMode == "realtime") + #expect(parsed.realtimeTransport == nil) + #expect(parsed.realtimeBrain == nil) + #expect(!parsed.hasGatewayRealtimeRelayTuple) + } + + @Test func `realtime provider config lookup is case insensitive`() { + let snapshot = Self.snapshot(talk: [ + "realtime": [ + "provider": "OPENAI", + "providers": [ + "openai": [ + "model": "gpt-live-1-codex", + "speakerVoice": "cedar", + ], + ], + ], + ]) + + let parsed = Self.parse(snapshot) + + #expect(parsed.realtimeProvider == "OPENAI") + #expect(parsed.realtimeModelId == "gpt-live-1-codex") + #expect(parsed.realtimeSpeakerVoice == "cedar") + } + + @Test func `fallback has no realtime selection`() { + let parsed = TalkModeGatewayConfigParser.fallback( + defaultModelIdFallback: "eleven_v3", + defaultSilenceTimeoutMs: TalkDefaults.silenceTimeoutMs, + envVoice: nil, + sagVoice: nil, + envApiKey: nil) + + #expect(parsed.realtimeProvider == nil) + #expect(parsed.realtimeModelId == nil) + #expect(parsed.realtimeSpeakerVoice == nil) + #expect(parsed.realtimeMode == nil) + #expect(parsed.realtimeTransport == nil) + #expect(parsed.realtimeBrain == nil) + } + + private static func snapshot(talk: [String: Any]) -> ConfigSnapshot { + ConfigSnapshot( + path: nil, + exists: true, + raw: nil, + hash: nil, + parsed: nil, + valid: true, + config: ["talk": AnyCodable(talk)], + issues: nil) + } + + private static func parse(_ snapshot: ConfigSnapshot) -> TalkModeGatewayConfigState { + TalkModeGatewayConfigParser.parse( + snapshot: snapshot, + defaultProvider: "elevenlabs", + defaultModelIdFallback: "eleven_v3", + defaultSilenceTimeoutMs: TalkDefaults.silenceTimeoutMs, + envVoice: nil, + sagVoice: nil, + envApiKey: nil) + } } diff --git a/apps/macos/Tests/OpenClawIPCTests/TalkModeRuntimeSpeechTests.swift b/apps/macos/Tests/OpenClawIPCTests/TalkModeRuntimeSpeechTests.swift index d303e06451ca..6bd85f125ac4 100644 --- a/apps/macos/Tests/OpenClawIPCTests/TalkModeRuntimeSpeechTests.swift +++ b/apps/macos/Tests/OpenClawIPCTests/TalkModeRuntimeSpeechTests.swift @@ -1,13 +1,458 @@ -import OpenClawKit +import Foundation +import OpenClawProtocol import Speech import Testing @testable import OpenClaw +@testable import OpenClawKit +private enum RuntimeTestAudioCaptureError: Error { + case inputUnavailable +} + +private struct RuntimeTestTimeout: Error, CustomStringConvertible { + let operation: String + + var description: String { + "timed out waiting for \(self.operation)" + } +} + +private final class RuntimeTestSignal: @unchecked Sendable { + private struct Waiter { + let id: UUID + let continuation: CheckedContinuation + var deadline: Task? + } + + private let lock = NSLock() + private var values: [Value] = [] + private var waiters: [Waiter] = [] + + func send(_ value: Value) { + let waiter: Waiter? = self.lock.withLock { + guard !self.waiters.isEmpty else { + self.values.append(value) + return nil + } + return self.waiters.removeFirst() + } + self.resume(waiter, with: .success(value)) + } + + func next(_ operation: String) async throws -> Value { + try Task.checkCancellation() + let id = UUID() + return try await withTaskCancellationHandler { + try await withCheckedThrowingContinuation { continuation in + let result: Result? = self.lock.withLock { + if Task.isCancelled { + return .failure(CancellationError()) + } + if !self.values.isEmpty { + return .success(self.values.removeFirst()) + } + self.waiters.append(Waiter(id: id, continuation: continuation)) + return nil + } + if let result { + continuation.resume(with: result) + return + } + let deadline = Task { + do { + try await Task.sleep(for: .seconds(5)) + self.fail(id, with: RuntimeTestTimeout(operation: operation)) + } catch {} + } + let retained = self.lock.withLock { + guard let index = self.waiters.firstIndex(where: { $0.id == id }) else { + return false + } + self.waiters[index].deadline = deadline + return true + } + if !retained { + deadline.cancel() + } + } + } onCancel: { + self.fail(id, with: CancellationError()) + } + } + + private func fail(_ id: UUID, with error: any Error) { + let waiter: Waiter? = self.lock.withLock { + guard let index = self.waiters.firstIndex(where: { $0.id == id }) else { + return nil + } + return self.waiters.remove(at: index) + } + self.resume(waiter, with: .failure(error)) + } + + private func resume(_ waiter: Waiter?, with result: Result) { + guard let waiter else { return } + waiter.deadline?.cancel() + waiter.continuation.resume(with: result) + } +} + +@MainActor +private final class RuntimeTestAudioCapture: RealtimeTalkAudioCapturing { + let suppressesInputDuringOutput = false + var startError: Error? + private(set) var startCount = 0 + + func start( + targetSampleRate: Double, + onAudio: @escaping @Sendable (RealtimeTalkAudioFrame) -> Void, + onFailure: @escaping @MainActor (String) -> Void) throws + { + self.startCount += 1 + if let startError = self.startError { + throw startError + } + } + + func stop() {} +} + +private actor RuntimeTestRelayRequestLog { + private var methods: [String] = [] + private var sessionIds: [String?] = [] + private nonisolated let changed = RuntimeTestSignal() + + func record(method: String, params: [String: AnyCodable]?) { + self.methods.append(method) + self.sessionIds.append(params?["sessionId"]?.stringValue) + self.changed.send(()) + } + + func snapshot() -> (methods: [String], sessionIds: [String?]) { + (self.methods, self.sessionIds) + } + + func waitForCount(_ count: Int) async throws { + while self.methods.count < count { + _ = try await self.changed.next("relay request \(count)") + } + } +} + +/// Relay whose close RPC is observable, so runtime termination paths can be proven to release the +/// server-side session instead of only dropping their local reference. +@MainActor +private func makeRecordingRelaySession( + requests: RuntimeTestRelayRequestLog, + audioCapture: RuntimeTestAudioCapture) -> RealtimeTalkRelaySession +{ + let session = RealtimeTalkRelaySession( + transport: RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in AsyncStream { $0.finish() } }, + request: { method, params, _ in + await requests.record(method: method, params: params) + return Data("{\"ok\":true}".utf8) + }), + options: .init(sessionKey: "main", provider: "openai", model: "gpt-realtime-2", voice: nil), + audioCapture: audioCapture, + pcmPlayer: RuntimeTestPCMPlayer(), + onStatus: { _ in }, + onSpeakingChanged: { _ in }) + session._test_setRelaySessionId("relay-1") + return session +} + +private func waitForRelayClose(_ requests: RuntimeTestRelayRequestLog) async throws -> [String] { + try await requests.waitForCount(1) + return await requests.snapshot().methods +} + +@MainActor +private final class RuntimeTestPCMPlayer: PCMStreamingAudioPlaying { + private(set) var stopCount = 0 + private let onStop: (() -> Void)? + + init(onStop: (() -> Void)? = nil) { + self.onStop = onStop + } + + func play( + stream: AsyncThrowingStream, + sampleRate: Double) async -> StreamingPlaybackResult + { + fatalError("Playback is not used by this test") + } + + func stop() -> Double? { + self.stopCount += 1 + self.onStop?() + return nil + } +} + +private actor RuntimeContinuationBarrier { + private var entered = false + private var released = false + private nonisolated let enteredSignal = RuntimeTestSignal() + private nonisolated let releaseSignal = RuntimeTestSignal() + + func wait() async throws { + self.entered = true + self.enteredSignal.send(()) + guard !self.released else { return } + do { + _ = try await self.releaseSignal.next("barrier release") + } catch { + self.release() + throw error + } + } + + func waitUntilEntered() async throws { + if self.entered { return } + _ = try await self.enteredSignal.next("barrier entry") + } + + func release() { + guard !self.released else { return } + self.released = true + self.releaseSignal.send(()) + } +} + +private func waitForRuntimeBarrier( + _ barrier: RuntimeContinuationBarrier, + cleaningUp attempt: Task) async throws +{ + do { + try await barrier.waitUntilEntered() + } catch { + attempt.cancel() + await barrier.release() + _ = try? await AsyncTimeout.withTimeout( + seconds: 5, + onTimeout: { RuntimeTestTimeout(operation: "cancelled runtime attempt") }, + operation: { await attempt.value }) + throw error + } +} + +private final class RuntimeCommitProbe: @unchecked Sendable { + private let lock = NSLock() + private var recordedValues: [String] = [] + + func record(_ value: String) { + self.lock.withLock { + self.recordedValues.append(value) + } + } + + func values() -> [String] { + self.lock.withLock { self.recordedValues } + } +} + +private final class RuntimeRecognitionCapture { + let name: String + + init(_ name: String) { + self.name = name + } +} + +private enum RuntimeRecognitionStartError: Error { + case failed +} + +private enum RuntimeRelayStartError: Error { + case failed +} + +private struct RuntimeConditionTimeout: Error, CustomStringConvertible { + let operation: String + + var description: String { + "timed out waiting for \(self.operation)" + } +} + +private func waitForRuntimeCondition( + _ operation: String, + condition: @escaping @Sendable () async -> Bool) async throws +{ + try await AsyncTimeout.withTimeout( + seconds: 1, + onTimeout: { RuntimeConditionTimeout(operation: operation) }, + operation: { + while !Task.isCancelled { + if await condition() { + return + } + } + throw CancellationError() + }) +} + +enum RuntimeRelayStartupPauseOutcome: Equatable { + case resume + case remainPaused + case disable +} + +@MainActor +private func makeRuntimeTestRealtimeSession( + player: RuntimeTestPCMPlayer) -> RealtimeTalkRelaySession +{ + RealtimeTalkRelaySession( + transport: RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in AsyncStream { $0.finish() } }, + request: { _, _, _ in Data("{\"ok\":true}".utf8) }), + options: .init(sessionKey: "main", provider: "openai", model: "gpt-realtime-2", voice: nil), + audioCapture: RuntimeTestAudioCapture(), + pcmPlayer: player, + onStatus: { _ in }, + onSpeakingChanged: { _ in }) +} + +private func makeRuntimeTestConfigSnapshot( + sessionKey: String = "main", + realtimeModel: String = "gpt-realtime-2") -> ConfigSnapshot +{ + ConfigSnapshot( + path: nil, + exists: true, + raw: nil, + hash: nil, + parsed: nil, + valid: true, + config: [ + "session": AnyCodable(["mainKey": AnyCodable(sessionKey)]), + "talk": AnyCodable([ + "realtime": AnyCodable([ + "provider": AnyCodable("openai"), + "providers": AnyCodable([ + "openai": AnyCodable([ + "model": AnyCodable(realtimeModel), + ]), + ]), + "mode": AnyCodable("realtime"), + "transport": AnyCodable("gateway-relay"), + "brain": AnyCodable("agent-consult"), + ]), + ]), + ], + issues: nil) +} + +private func makeRuntimeTestBootstrap( + requests: RuntimeTestRelayRequestLog = RuntimeTestRelayRequestLog(), + createBarrier: RuntimeContinuationBarrier? = nil, + probe: RuntimeCommitProbe? = nil, + sessionKey: String = "main", + realtimeModel: String = "gpt-realtime-2") throws -> GatewayConnection.RealtimeTalkBootstrap +{ + let events = AsyncStream.makeStream(bufferingPolicy: .bufferingNewest(8)) + let result = TalkSessionCreateResult( + sessionid: "talk-session", + mode: AnyCodable("realtime"), + transport: AnyCodable("gateway-relay"), + brain: AnyCodable("agent-consult"), + relaysessionid: "relay-1") + let resultData = try JSONEncoder().encode(result) + let transport = RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in events.stream }, + request: { method, params, _ in + await requests.record(method: method, params: params) + if method == "talk.session.create" { + probe?.record("start") + if let createBarrier { + try await createBarrier.wait() + } + events.continuation.yield(EventFrame( + type: "event", + event: "talk.event", + payload: AnyCodable([ + "relaySessionId": "relay-1", + "type": "ready", + ]), + seq: nil, + stateversion: nil)) + return resultData + } + if method == "talk.session.close" { + events.continuation.finish() + } + return Data("{\"ok\":true}".utf8) + }, + isCurrent: { true }) + return GatewayConnection.RealtimeTalkBootstrap( + transport: transport, + configSnapshot: makeRuntimeTestConfigSnapshot( + sessionKey: sessionKey, + realtimeModel: realtimeModel), + sessionKey: sessionKey) +} + +private actor RuntimeTestBootstrapSequence { + private var bootstraps: [GatewayConnection.RealtimeTalkBootstrap] + private let firstBarrier: RuntimeContinuationBarrier? + private var count = 0 + + init( + bootstraps: [GatewayConnection.RealtimeTalkBootstrap], + firstBarrier: RuntimeContinuationBarrier? = nil) + { + self.bootstraps = bootstraps + self.firstBarrier = firstBarrier + } + + func next() async throws -> GatewayConnection.RealtimeTalkBootstrap { + let index = self.count + self.count += 1 + if index == 0, let firstBarrier { + try await firstBarrier.wait() + } + guard self.bootstraps.indices.contains(index) else { + throw RuntimeRelayStartError.failed + } + return self.bootstraps[index] + } + + func requestCount() -> Int { + self.count + } +} + +@Suite(.serialized) struct TalkModeRuntimeSpeechTests { + @Test func `macOS realtime relay requires local opt in and exact Gateway tuple`() { + #expect(!TalkModeRuntime.shouldUseRealtimeRelay( + localOptIn: false, + hasGatewayRealtimeRelayTuple: false)) + #expect(!TalkModeRuntime.shouldUseRealtimeRelay( + localOptIn: false, + hasGatewayRealtimeRelayTuple: true)) + #expect(!TalkModeRuntime.shouldUseRealtimeRelay( + localOptIn: true, + hasGatewayRealtimeRelayTuple: false)) + #expect(TalkModeRuntime.shouldUseRealtimeRelay( + localOptIn: true, + hasGatewayRealtimeRelayTuple: true)) + } + + @Test @MainActor func `macOS realtime relay preference defaults off and reads explicit opt in`() async { + await TestIsolation.withUserDefaultsValues([talkRealtimeRelayEnabledKey: nil]) { + #expect(!AppState(preview: true).talkRealtimeRelayEnabled) + } + await TestIsolation.withUserDefaultsValues([talkRealtimeRelayEnabledKey: true]) { + #expect(AppState(preview: true).talkRealtimeRelayEnabled) + } + } + @Test func `speech request uses dictation defaults`() { let request = SFSpeechAudioBufferRecognitionRequest() - TalkModeRuntime.configureRecognitionRequest(request) + TalkRecognitionCaptureLifecycle.configure(request) #expect(request.shouldReportPartialResults) #expect(request.taskHint == .dictation) @@ -55,6 +500,724 @@ struct TalkModeRuntimeSpeechTests { TalkMLXSpeechSynthesizer.SynthesizeError.modelLoadFailed("missing")) == .fallback) } + @Test func `realtime recovery uses the iOS retry budget`() { + #expect(TalkModeRuntime.realtimeRestartAttempt( + previousRapidRestarts: 1, + activeDuration: 5) == 2) + #expect(TalkModeRuntime.realtimeRestartAttempt( + previousRapidRestarts: 2, + activeDuration: 31) == 1) + #expect(TalkModeRuntime.realtimeRestartDelayNanoseconds(attempt: 1) == 500_000_000) + #expect(TalkModeRuntime.realtimeRestartDelayNanoseconds(attempt: 2) == 2_000_000_000) + #expect(TalkModeRuntime.realtimeRestartDelayNanoseconds(attempt: 3) == nil) + } + + @Test @MainActor func `ready then audio failure clears relay owner and schedules bounded recovery`() async { + let runtime = TalkModeRuntime() + let session = RealtimeTalkRelaySession( + transport: RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in AsyncStream { $0.finish() } }, + request: { _, _, _ in throw CancellationError() }), + options: .init(sessionKey: "main", provider: "openai", model: "gpt-realtime-2", voice: nil), + audioCapture: RuntimeTestAudioCapture(), + pcmPlayer: RuntimeTestPCMPlayer(), + onStatus: { _ in }, + onSpeakingChanged: { _ in }) + let relayGeneration = await runtime._test_prepareEnabledRealtimeSessionForClose(session) + + await runtime.handleRealtimeTermination( + .remoteClose(reason: "stale"), + relayGeneration: relayGeneration &- 1) + #expect(await runtime.realtimeSession != nil) + + await runtime.handleRealtimeTermination( + .audioInputFailed(message: "microphone unavailable"), + relayGeneration: relayGeneration) + + #expect(await runtime.realtimeSession == nil) + #expect(await runtime.rapidRealtimeRestartCount == 1) + #expect(await runtime.realtimeRestartTask != nil) + + await runtime.setEnabled(false) + session.stop() + } + + @Test @MainActor func `selected microphone restart failure closes relay and schedules recovery`() async throws { + let runtime = TalkModeRuntime() + let requests = RuntimeTestRelayRequestLog() + let session = makeRecordingRelaySession( + requests: requests, + audioCapture: RuntimeTestAudioCapture()) + let relayGeneration = await runtime._test_prepareEnabledRealtimeSessionForClose(session) + + await runtime.handleRealtimeInputRestartFailure( + "selected microphone unavailable", + relayGeneration: relayGeneration) + + #expect(await runtime.realtimeSession == nil) + #expect(await runtime.rapidRealtimeRestartCount == 1) + #expect(await runtime.realtimeRestartTask != nil) + + // Ownership must not be dropped while the server relay stays live; recovery would then + // run a second session against the same gateway lease. + let recorded = try await waitForRelayClose(requests) + #expect(recorded == ["talk.session.close"]) + #expect(await requests.snapshot().sessionIds == ["relay-1"]) + + await runtime.setEnabled(false) + session.stop() + } + + @Test func `stale termination and callbacks cannot tear down or project over a successor`() async throws { + let runtime = TalkModeRuntime() + let stopEntered = RuntimeTestSignal() + let releaseStop = DispatchSemaphore(value: 0) + var released = false + defer { + if !released { + releaseStop.signal() + } + } + let sessionA = await MainActor.run { + makeRuntimeTestRealtimeSession(player: RuntimeTestPCMPlayer(onStop: { + stopEntered.send(()) + _ = releaseStop.wait(timeout: .now() + 5) + })) + } + let sessionB = await MainActor.run { makeRuntimeTestRealtimeSession(player: RuntimeTestPCMPlayer()) } + let generationA = await runtime._test_prepareEnabledRealtimeSessionForClose(sessionA) + + let staleTermination = Task { + await runtime.handleRealtimeTermination( + .remoteClose(reason: "replaced"), + relayGeneration: generationA) + } + do { + _ = try await stopEntered.next("stale session stop") + } catch { + released = true + releaseStop.signal() + await staleTermination.value + throw error + } + let invalidatedGeneration = await runtime.realtimeRelayGeneration + #expect(invalidatedGeneration != generationA) + let generationB = await runtime._test_prepareEnabledRealtimeSessionForClose(sessionB) + #expect(await runtime.realtimeSession === sessionB) + released = true + releaseStop.signal() + await staleTermination.value + + await MainActor.run { TalkModeController.shared.updatePartialTranscript("successor") } + await runtime.handleRealtimeSpeakingChanged(false, relayGeneration: generationA) + await runtime.handleRealtimeInputLevel(0.9, relayGeneration: generationA) + await runtime.handleRealtimeOutputLevel(0.8, relayGeneration: generationA) + await runtime.handleRealtimeTranscript( + .init(role: "user", text: "stale", isFinal: false), + relayGeneration: generationA) + + #expect(await runtime.realtimeRelayGeneration == generationB) + #expect(await runtime.realtimeSession === sessionB) + #expect(await runtime.realtimeRestartTask == nil) + #expect(await MainActor.run { TalkModeController.shared.partialTranscript } == "successor") + + await runtime.setEnabled(false) + await MainActor.run { sessionB.stop() } + } + + @Test func `stale preference cleanup preserves successor session recognition and UI`() async throws { + let runtime = TalkModeRuntime() + let mainActorEntered = RuntimeTestSignal() + let mainActorFinished = RuntimeTestSignal() + let releaseMainActor = DispatchSemaphore(value: 0) + var released = false + defer { + if !released { + releaseMainActor.signal() + } + } + let sessionA = await MainActor.run { + makeRuntimeTestRealtimeSession(player: RuntimeTestPCMPlayer()) + } + let sessionB = await MainActor.run { + makeRuntimeTestRealtimeSession(player: RuntimeTestPCMPlayer()) + } + _ = await runtime._test_prepareEnabledRealtimeSessionForClose(sessionA) + let lifecycleA = await runtime.lifecycleGeneration + _ = try #require(await runtime._test_beginRecognitionAttempt( + lifecycleGeneration: lifecycleA)) + let recognitionCleanup = RuntimeCommitProbe() + await runtime._test_setRecognitionCleanupProbe { + recognitionCleanup.record("cleanup") + } + await MainActor.run { + TalkModeController.shared.updatePartialTranscript("successor") + } + DispatchQueue.main.async { + mainActorEntered.send(()) + _ = releaseMainActor.wait(timeout: .now() + 5) + mainActorFinished.send(()) + } + _ = try await mainActorEntered.next("MainActor preference blocker") + + let stalePreference = Task { + await runtime.realtimeRelayPreferenceDidChange() + } + do { + try await waitForRuntimeCondition("preference owner detachment") { + await runtime.realtimeSession == nil + } + } catch { + released = true + releaseMainActor.signal() + _ = try? await mainActorFinished.next("MainActor preference cleanup") + await stalePreference.value + throw error + } + #expect(recognitionCleanup.values() == ["cleanup"]) + + _ = await runtime._test_prepareEnabledRealtimeSessionForClose(sessionB) + let lifecycleB = await runtime.lifecycleGeneration + let recognitionB = try #require(await runtime._test_beginRecognitionAttempt( + lifecycleGeneration: lifecycleB)) + #expect(recognitionCleanup.values() == ["cleanup", "cleanup"]) + + released = true + releaseMainActor.signal() + _ = try await mainActorFinished.next("MainActor preference cleanup") + await stalePreference.value + + #expect(await runtime.realtimeSession === sessionB) + #expect(await runtime.recognitionGeneration == recognitionB) + #expect(recognitionCleanup.values() == ["cleanup", "cleanup"]) + #expect(await MainActor.run { TalkModeController.shared.partialTranscript } == "successor") + + await runtime._test_setRecognitionCleanupProbe(nil) + await runtime.setEnabled(false) + await MainActor.run { sessionB.stop() } + } + + @Test @MainActor func `unpause that cannot restart capture closes relay and schedules recovery`() async throws { + let runtime = TalkModeRuntime() + let requests = RuntimeTestRelayRequestLog() + let audioCapture = RuntimeTestAudioCapture() + let session = makeRecordingRelaySession(requests: requests, audioCapture: audioCapture) + _ = await runtime._test_prepareEnabledRealtimeSessionForClose(session) + + await runtime.setPaused(true) + audioCapture.startError = RuntimeTestAudioCaptureError.inputUnavailable + await runtime.setPaused(false) + + // Talk must never stay enabled with no microphone and no route back: the failed unpause + // has to reach the same bounded recovery / native-speech fallback as any other capture loss. + #expect(await runtime.realtimeSession == nil) + #expect(await runtime.rapidRealtimeRestartCount == 1) + #expect(await runtime.realtimeRestartTask != nil) + + let recorded = try await waitForRelayClose(requests) + #expect(recorded == ["talk.session.close"]) + + await runtime.setEnabled(false) + session.stop() + } + + @Test @MainActor func `paused reenable lets pinned bootstrap refresh realtime selection`() async throws { + try await TestIsolation.withUserDefaultsValues([talkRealtimeRelayEnabledKey: true]) { + let previousRelayPreference = AppStateStore.shared.talkRealtimeRelayEnabled + AppStateStore.shared.talkRealtimeRelayEnabled = true + defer { + AppStateStore.shared.talkRealtimeRelayEnabled = previousRelayPreference + } + let requests = RuntimeTestRelayRequestLog() + let bootstrap = try makeRuntimeTestBootstrap( + requests: requests, + realtimeModel: "fresh-model") + let runtime = TalkModeRuntime(realtimeTalkBootstrapProvider: { bootstrap }) + await runtime._test_setVoiceWakeReadiness(supported: true, permissionGranted: true) + _ = await runtime._test_prepareEnabledLifecycle() + await runtime.setPaused(true) + await runtime.setEnabled(false) + let staleConfig = await runtime.fallbackTalkConfig() + await runtime.applyTalkConfig(staleConfig) + + await runtime.setEnabled(true) + #expect(await runtime.realtimeSession == nil) + #expect(await requests.snapshot().methods.isEmpty) + + await runtime.setPaused(false) + + #expect(await runtime.realtimeSession != nil) + #expect(await runtime.realtimeModelId == "fresh-model") + #expect(await requests.snapshot().methods.first == "talk.session.create") + await runtime.setEnabled(false) + } + } + + @Test @MainActor func `pausing realtime resets visible state and ignores late callbacks`() async { + let runtime = TalkModeRuntime() + let player = RuntimeTestPCMPlayer() + let session = makeRuntimeTestRealtimeSession(player: player) + let relayGeneration = await runtime._test_prepareEnabledRealtimeSessionForClose(session) + TalkModeController.shared.updatePhase(.speaking) + TalkModeController.shared.updateLevel(0.8) + TalkModeController.shared.updatePartialTranscript("stale") + + await runtime.setPaused(true) + await runtime.handleRealtimeSpeakingChanged(true, relayGeneration: relayGeneration) + await runtime.handleRealtimeInputLevel(0.9, relayGeneration: relayGeneration) + await runtime.handleRealtimeOutputLevel(0.8, relayGeneration: relayGeneration) + await runtime.handleRealtimeTranscript( + .init(role: "user", text: "late transcript", isFinal: false), + relayGeneration: relayGeneration) + + #expect(await runtime.phase == .idle) + #expect(TalkModeController.shared.phase == .idle) + #expect(TalkModeController.shared.level == 0) + #expect(TalkModeController.shared.partialTranscript.isEmpty) + #expect(player.stopCount == 0) + + await runtime.setEnabled(false) + session.stop() + } + + @Test @MainActor func `resuming realtime restarts input and reuses the relay`() async throws { + let runtime = TalkModeRuntime() + let audioCapture = RuntimeTestAudioCapture() + let player = RuntimeTestPCMPlayer() + let eventChannel = AsyncStream.makeStream() + let result = TalkSessionCreateResult( + sessionid: "talk-session", + mode: AnyCodable("realtime"), + transport: AnyCodable("gateway-relay"), + brain: AnyCodable("agent-consult"), + relaysessionid: "relay-1") + let resultData = try JSONEncoder().encode(result) + let session = RealtimeTalkRelaySession( + transport: RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in eventChannel.stream }, + request: { method, _, _ in + if method == "talk.session.create" { + eventChannel.continuation.yield(EventFrame( + type: "event", + event: "talk.event", + payload: AnyCodable([ + "relaySessionId": "relay-1", + "type": "ready", + ]), + seq: nil, + stateversion: nil)) + return resultData + } + return Data("{\"ok\":true}".utf8) + }), + options: .init(sessionKey: "main", provider: "openai", model: "gpt-realtime-2", voice: nil), + audioCapture: audioCapture, + pcmPlayer: player, + onStatus: { _ in }, + onSpeakingChanged: { _ in }) + try await session.start() + let relayGeneration = await runtime._test_prepareEnabledRealtimeSessionForClose(session) + + await runtime.setPaused(true) + await runtime.setPaused(false) + + #expect(audioCapture.startCount == 2) + #expect(await runtime.realtimeSession === session) + await runtime.handleRealtimeSpeakingChanged(true, relayGeneration: relayGeneration) + #expect(await runtime.phase == .speaking) + + await runtime.setEnabled(false) + session.stop() + eventChannel.continuation.finish() + } + + @Test @MainActor func `disabling during relay startup stops the published session`() async throws { + let barrier = RuntimeContinuationBarrier() + let requests = RuntimeTestRelayRequestLog() + let probe = RuntimeCommitProbe() + let bootstrap = try makeRuntimeTestBootstrap( + requests: requests, + createBarrier: barrier, + probe: probe) + let runtime = TalkModeRuntime(realtimeTalkBootstrapProvider: { bootstrap }) + let lifecycleGeneration = await runtime._test_prepareEnabledLifecycle() + await runtime._test_enableRealtimeRelaySelection() + let attempt = Task { + do { + try await runtime.startRealtimeRelay(generation: lifecycleGeneration) + return true + } catch { + return false + } + } + + try await waitForRuntimeBarrier(barrier, cleaningUp: attempt) + #expect(await runtime.realtimeSession != nil) + await runtime.setEnabled(false) + await barrier.release() + + #expect(await attempt.value == false) + #expect(await runtime.realtimeSession == nil) + #expect(probe.values() == ["start"]) + #expect(await requests.snapshot().methods.contains("talk.session.create")) + } + + @Test func `failed realtime bootstrap clears prior Gateway selection`() async throws { + let runtime = TalkModeRuntime(realtimeTalkBootstrapProvider: { + throw RuntimeRelayStartError.failed + }) + let staleConfig = await runtime.parseTalkConfig( + makeRuntimeTestConfigSnapshot(realtimeModel: "stale-model")) + await runtime.applyTalkConfig(staleConfig) + let lifecycleGeneration = await runtime._test_prepareEnabledLifecycle() + await runtime._test_enableRealtimeRelaySelection() + + do { + try await runtime.startRealtimeRelay(generation: lifecycleGeneration) + Issue.record("expected bootstrap failure") + } catch {} + + #expect(await runtime.realtimeProvider == nil) + #expect(await runtime.realtimeModelId == nil) + #expect(await !runtime.hasGatewayRealtimeRelayTuple) + await runtime.setEnabled(false) + } + + @Test func `stale realtime config application cannot replace current selection`() async throws { + let checkpoint = RuntimeContinuationBarrier() + let bootstrap = try makeRuntimeTestBootstrap(realtimeModel: "stale-model") + let runtime = TalkModeRuntime(realtimeTalkBootstrapProvider: { bootstrap }) + await runtime._test_setRealtimeConfigApplicationCheckpoint { + try? await checkpoint.wait() + } + let currentConfig = await runtime.parseTalkConfig( + makeRuntimeTestConfigSnapshot(realtimeModel: "current-model")) + await runtime.applyTalkConfig(currentConfig) + let lifecycleGeneration = await runtime._test_prepareEnabledLifecycle() + await runtime._test_enableRealtimeRelaySelection() + + let attempt = Task { + do { + try await runtime.startRealtimeRelay(generation: lifecycleGeneration) + return true + } catch { + return false + } + } + try await waitForRuntimeBarrier(checkpoint, cleaningUp: attempt) + _ = await runtime.beginRealtimeReconfiguration() + await checkpoint.release() + + #expect(await attempt.value == false) + #expect(await runtime.realtimeModelId == "current-model") + await runtime._test_setRealtimeConfigApplicationCheckpoint(nil) + await runtime.setEnabled(false) + } + + @Test(arguments: [ + RuntimeRelayStartupPauseOutcome.resume, + .remainPaused, + .disable, + ]) + @MainActor + func `relay startup pause retries only a matching resume`( + outcome: RuntimeRelayStartupPauseOutcome) async throws + { + let barrier = RuntimeContinuationBarrier() + let requests = RuntimeTestRelayRequestLog() + let probe = RuntimeCommitProbe() + let bootstrap = try makeRuntimeTestBootstrap( + requests: requests, + createBarrier: barrier, + probe: probe) + let runtime = TalkModeRuntime(realtimeTalkBootstrapProvider: { bootstrap }) + let lifecycleGeneration = await runtime._test_prepareEnabledLifecycle() + await runtime._test_enableRealtimeRelaySelection() + let attempt = Task { + do { + try await runtime.startRealtimeRelay(generation: lifecycleGeneration) + return true + } catch { + return false + } + } + + try await waitForRuntimeBarrier(barrier, cleaningUp: attempt) + #expect(await runtime.realtimeSession != nil) + await runtime.setPaused(true) + if outcome != .remainPaused { + await runtime.setPaused(false) + } + if outcome == .disable { + await runtime.setEnabled(false) + } + await barrier.release() + + #expect(await attempt.value == false) + #expect(await runtime.realtimeSession == nil) + if await runtime.consumePendingRealtimeRelayStart() { probe.record("retry") } + if await runtime.consumePendingRealtimeRelayStart() { probe.record("retry") } + #expect(probe.values() == (outcome == .resume ? ["start", "retry"] : ["start"])) + + await runtime.setEnabled(false) + } + + @Test @MainActor func `paused pinned bootstrap retries before stale tuple fallback`() async throws { + try await TestIsolation.withUserDefaultsValues([talkRealtimeRelayEnabledKey: true]) { + let previousRelayPreference = AppStateStore.shared.talkRealtimeRelayEnabled + AppStateStore.shared.talkRealtimeRelayEnabled = true + defer { + AppStateStore.shared.talkRealtimeRelayEnabled = previousRelayPreference + } + let barrier = RuntimeContinuationBarrier() + let requests = RuntimeTestRelayRequestLog() + let firstBootstrap = try makeRuntimeTestBootstrap( + requests: requests, + realtimeModel: "fresh-model") + let retryBootstrap = try makeRuntimeTestBootstrap( + requests: requests, + realtimeModel: "fresh-model") + let sequence = RuntimeTestBootstrapSequence( + bootstraps: [firstBootstrap, retryBootstrap], + firstBarrier: barrier) + let runtime = TalkModeRuntime(realtimeTalkBootstrapProvider: { try await sequence.next() }) + await runtime._test_setVoiceWakeReadiness(supported: true, permissionGranted: true) + let attempt = Task { + await runtime.setEnabled(true) + return true + } + + try await waitForRuntimeBarrier(barrier, cleaningUp: attempt) + #expect(await runtime.hasGatewayRealtimeRelayTuple == false) + await runtime.setPaused(true) + await runtime.setPaused(false) + await barrier.release() + + _ = await attempt.value + #expect(await sequence.requestCount() == 2) + #expect(await requests.snapshot().methods == ["talk.session.create"]) + #expect(await runtime.realtimeSession != nil) + #expect(await runtime.realtimeModelId == "fresh-model") + await runtime.setEnabled(false) + } + } + + @Test func `processed recognition start failure retries a fresh raw capture`() { + let probe = RuntimeCommitProbe() + + let started = TalkRecognitionCaptureLifecycle.start( + isCurrent: { true }, + prepare: { enableVoiceProcessing -> RuntimeRecognitionCapture in + if enableVoiceProcessing { + probe.record("prepare-processed") + probe.record("cleanup-processed") + throw RuntimeRecognitionStartError.failed + } + probe.record("prepare-raw") + return RuntimeRecognitionCapture("raw") + }, + discard: { probe.record("discard-\($0.name)") }, + publish: { probe.record("publish-\($0.name)") }, + onFailure: { enableVoiceProcessing, _ in + probe.record(enableVoiceProcessing ? "failed-processed" : "failed-raw") + }) + + #expect(started) + #expect(probe.values() == [ + "prepare-processed", + "cleanup-processed", + "failed-processed", + "prepare-raw", + "publish-raw", + ]) + } + + @Test func `failed recognition candidates clean up without publishing`() { + let probe = RuntimeCommitProbe() + + let started = TalkRecognitionCaptureLifecycle.start( + isCurrent: { true }, + prepare: { enableVoiceProcessing -> RuntimeRecognitionCapture in + let kind = enableVoiceProcessing ? "processed" : "raw" + probe.record("prepare-\(kind)") + probe.record("cleanup-\(kind)") + throw RuntimeRecognitionStartError.failed + }, + discard: { probe.record("discard-\($0.name)") }, + publish: { probe.record("publish-\($0.name)") }, + onFailure: { enableVoiceProcessing, _ in + probe.record(enableVoiceProcessing ? "failed-processed" : "failed-raw") + }) + + #expect(!started) + #expect(probe.values() == [ + "prepare-processed", + "cleanup-processed", + "failed-processed", + "prepare-raw", + "cleanup-raw", + "failed-raw", + ]) + } + + @Test @MainActor func `stale relay cleanup cannot clear a newer owned session`() async throws { + let barrier = RuntimeContinuationBarrier() + let requestsA = RuntimeTestRelayRequestLog() + let requestsB = RuntimeTestRelayRequestLog() + let bootstrapA = try makeRuntimeTestBootstrap(requests: requestsA) + let bootstrapB = try makeRuntimeTestBootstrap(requests: requestsB) + let sequence = RuntimeTestBootstrapSequence( + bootstraps: [bootstrapA, bootstrapB], + firstBarrier: barrier) + let runtime = TalkModeRuntime(realtimeTalkBootstrapProvider: { + try await sequence.next() + }) + let lifecycleA = await runtime._test_prepareEnabledLifecycle() + await runtime._test_enableRealtimeRelaySelection() + let attemptA = Task { + do { + try await runtime.startRealtimeRelay(generation: lifecycleA) + return true + } catch { + return false + } + } + + try await waitForRuntimeBarrier(barrier, cleaningUp: attemptA) + await runtime.setEnabled(false) + let lifecycleB = await runtime._test_prepareEnabledLifecycle() + await runtime._test_enableRealtimeRelaySelection() + try await runtime.startRealtimeRelay(generation: lifecycleB) + let sessionB = try #require(await runtime.realtimeSession) + await barrier.release() + + #expect(await attemptA.value == false) + #expect(await runtime.realtimeSession === sessionB) + #expect(await requestsA.snapshot().methods.isEmpty) + #expect(await requestsB.snapshot().methods.first == "talk.session.create") + + await runtime.setEnabled(false) + sessionB.stop() + } + + @Test @MainActor func `current relay failure owner can transition to native fallback`() async throws { + let runtime = TalkModeRuntime() + let lifecycleGeneration = await runtime._test_prepareEnabledLifecycle() + let recognitionGeneration = try #require(await runtime._test_beginRecognitionAttempt( + lifecycleGeneration: lifecycleGeneration)) + let relayGeneration = await runtime.realtimeRelayGeneration + + #expect(await runtime.commitNativeFallback( + lifecycleGeneration: lifecycleGeneration, + recognitionGeneration: recognitionGeneration, + relayGeneration: relayGeneration, + status: "native")) + #expect(TalkModeController.shared.partialTranscript == "native") + + await runtime.setEnabled(false) + } + + @Test @MainActor func `stale relay fallback cannot replace successor recognition owner`() async throws { + let runtime = TalkModeRuntime() + let lifecycleGeneration = await runtime._test_prepareEnabledLifecycle() + let staleRecognition = try #require(await runtime._test_beginRecognitionAttempt( + lifecycleGeneration: lifecycleGeneration)) + let relayGeneration = await runtime.realtimeRelayGeneration + let successorRecognition = try #require(await runtime._test_beginRecognitionAttempt( + lifecycleGeneration: lifecycleGeneration)) + TalkModeController.shared.updatePartialTranscript("successor") + + let accepted = await runtime.commitNativeFallback( + lifecycleGeneration: lifecycleGeneration, + recognitionGeneration: staleRecognition, + relayGeneration: relayGeneration, + status: "stale fallback") + + #expect(await runtime.recognitionGeneration == successorRecognition) + #expect(TalkModeController.shared.partialTranscript == "successor") + #expect(!accepted) + + await runtime.setEnabled(false) + } + + @Test func `blocked fallback projection cannot overwrite a successor`() async throws { + let runtime = TalkModeRuntime() + let lifecycleGeneration = await runtime._test_prepareEnabledLifecycle() + let recognitionGeneration = try #require(await runtime._test_beginRecognitionAttempt( + lifecycleGeneration: lifecycleGeneration)) + let relayGeneration = await runtime.realtimeRelayGeneration + await MainActor.run { + TalkModeController.shared.updatePartialTranscript("successor") + } + let blocked = AsyncStream.makeStream() + let releaseMainActor = DispatchSemaphore(value: 0) + var didReleaseMainActor = false + defer { + if !didReleaseMainActor { + releaseMainActor.signal() + } + } + DispatchQueue.main.async { + blocked.continuation.yield() + releaseMainActor.wait() + } + for await _ in blocked.stream.prefix(1) {} + + let staleCommit = Task { + await runtime.commitNativeFallback( + lifecycleGeneration: lifecycleGeneration, + recognitionGeneration: recognitionGeneration, + relayGeneration: relayGeneration, + status: "stale fallback") + } + try await waitForRuntimeCondition("fallback to enter projection") { + await runtime.phase == .listening + } + let successor = Task { await runtime.setEnabled(false) } + try await waitForRuntimeCondition("successor generation") { + await runtime.realtimeRelayGeneration != relayGeneration + } + didReleaseMainActor = true + releaseMainActor.signal() + + #expect(await staleCommit.value == false) + await successor.value + #expect(await MainActor.run { TalkModeController.shared.partialTranscript }.isEmpty) + } + + @Test func `stale recognition attempt preserves current owner`() async { + let runtime = TalkModeRuntime() + let lifecycleGeneration = await runtime._test_prepareEnabledLifecycle() + let currentRecognition = await runtime._test_beginRecognitionAttempt( + lifecycleGeneration: lifecycleGeneration) + + let staleRecognition = await runtime._test_beginRecognitionAttempt( + lifecycleGeneration: lifecycleGeneration &- 1) + + #expect(currentRecognition != nil) + #expect(staleRecognition == nil) + #expect(await runtime.recognitionGeneration == currentRecognition) + await runtime.setEnabled(false) + } + + @Test func `cancelled recognition attempt discards capture before publication`() { + let probe = RuntimeCommitProbe() + var isCurrent = true + + let started = TalkRecognitionCaptureLifecycle.start( + isCurrent: { isCurrent }, + prepare: { _ in + isCurrent = false + return RuntimeRecognitionCapture("processed") + }, + discard: { probe.record("discard-\($0.name)") }, + publish: { probe.record("publish-\($0.name)") }, + onFailure: { _, _ in probe.record("failed") }) + + #expect(!started) + #expect(probe.values() == ["discard-processed"]) + } + @Test func `talk speak params carry resolved voice and directive overrides`() { let params = TalkModeRuntime.makeTalkSpeakParams( text: "hello", diff --git a/docs/nodes/talk.md b/docs/nodes/talk.md index 48f10a45bdc6..6d88aa503578 100644 --- a/docs/nodes/talk.md +++ b/docs/nodes/talk.md @@ -58,6 +58,75 @@ stopping Talk releases the camera and microphone tracks. - Replies are written to WebChat (same as typing). - **Interrupt on speech** (default on): if the user talks while the assistant is speaking, playback stops and the interruption timestamp is noted for the next prompt. +## Realtime Talk over the Gateway relay (macOS) + +macOS defaults to the native path above: Apple Speech recognition, Gateway chat, and `talk.speak` +playback. It switches to a streamed realtime session only when `talk.realtime` selects all three +of these together: + +| Key | Required value | +| ----------- | --------------- | +| `mode` | `realtime` | +| `transport` | `gateway-relay` | +| `brain` | `agent-consult` | + +Any other combination — including a partially set one — keeps the native path. + +```json5 +{ + talk: { + realtime: { + provider: "openai", + providers: { + openai: { + model: "gpt-realtime-2.1", + speakerVoice: "cedar", + }, + }, + mode: "realtime", + transport: "gateway-relay", + brain: "agent-consult", + }, + }, +} +``` + +The Mac must also opt in locally with **Settings > Voice Wake > Use realtime Gateway relay**. +This preference defaults off and stays on that Mac; Gateway config alone never activates the +streamed path. Keep `transport: "webrtc"` for browser or iOS client-owned sessions; macOS uses +the relay only when the config explicitly selects `gateway-relay`. + +The Gateway must also advertise `gateway-relay` and `agent-consult` for the selected provider in +`talk.catalog`. Realtime requires macOS 26 or newer, matching Voice Wake; on older versions the +Talk and Voice Wake controls are unavailable. + +### When realtime cannot start + +Talk never silently sits idle. If the relay fails to start — no Gateway route, rejected +credentials, or an unsupported model — the failure is logged, the overlay shows the reason, and +Talk falls back to the native speech path for that session. + +Once a session is running, a dropped relay reconnects on a bounded retry schedule (roughly 0.5 s +then 2 s). If those attempts are exhausted, the overlay reports +`Realtime disconnected repeatedly — using native speech` and the next start bypasses realtime. +Losing the microphone mid-session closes the relay and takes the same route. + +Relay output cancellation is turn-scoped. Clients copy the current `turnId` from the +`talk.event` audio envelope. Matching ids return `applied`, stale ids return `stale`, and +sessions without an active turn return `idle`. Older clients that omit `turnId` still cancel +the current turn: + +```json +{ + "method": "talk.session.cancelOutput", + "params": { + "sessionId": "relay-session-id", + "turnId": "turn-7", + "reason": "barge-in" + } +} +``` + ## Voice directives in replies The assistant can prefix a reply with a single JSON line to control voice: