From 4aa8bf12e3f5dc12f2f9845741f7088ddd91e079 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Thu, 9 Jul 2026 22:21:00 +0100 Subject: [PATCH] refactor(voice): dedupe talk audio primitives into OpenClawKit (#103099) * refactor(voice): dedupe talk audio primitives into OpenClawKit One buffered TTS clip player (watchdog + metering) shared by iOS and macOS replaces the near-identical pair; macOS TalkPlaybackResult was a field-for-field duplicate of StreamingPlaybackResult and is gone. One AVAudioPCMBuffer RMS helper replaces five per-file copies; the metered PCM passthrough moves into PCMPlaybackEnvelope; TalkModeManager's duplicated ElevenLabs PCM-then-mp3-retry blocks collapse into one helper. Net -378 lines, no behavior change. * fix(ios): drop stray blank line and resync i18n inventory * fix(ios): drop stray blank line after helper extraction --- apps/.i18n/native-source.json | 4 +- .../Voice/RealtimeTalkRelaySession.swift | 22 +-- .../Voice/TalkGatewaySpeechClient.swift | 158 +--------------- apps/ios/Sources/Voice/TalkModeManager.swift | 156 ++++++---------- .../Sources/OpenClaw/MicLevelMonitor.swift | 17 +- .../Sources/OpenClaw/TalkAudioPlayer.swift | 172 ------------------ .../Sources/OpenClaw/TalkModeController.swift | 19 +- .../Sources/OpenClaw/TalkModeRuntime.swift | 24 +-- .../Sources/OpenClaw/VoiceWakeRuntime.swift | 15 +- ...ift => TalkBufferedAudioPlayerTests.swift} | 13 +- .../OpenClawKit/TalkBufferedAudioPlayer.swift | 171 +++++++++++++++++ .../OpenClawKit/TalkPlaybackLevelMeters.swift | 44 +++++ .../OpenClawKitTests/TalkWaveformTests.swift | 21 +++ 13 files changed, 314 insertions(+), 522 deletions(-) delete mode 100644 apps/macos/Sources/OpenClaw/TalkAudioPlayer.swift rename apps/macos/Tests/OpenClawIPCTests/{TalkAudioPlayerTests.swift => TalkBufferedAudioPlayerTests.swift} (87%) create mode 100644 apps/shared/OpenClawKit/Sources/OpenClawKit/TalkBufferedAudioPlayer.swift diff --git a/apps/.i18n/native-source.json b/apps/.i18n/native-source.json index 476a1aaf51b4..17772a0b83e1 100644 --- a/apps/.i18n/native-source.json +++ b/apps/.i18n/native-source.json @@ -13491,7 +13491,7 @@ }, { "kind": "conditional-branch", - "line": 3065, + "line": 3032, "path": "apps/ios/Sources/Voice/TalkModeManager.swift", "source": "Native", "surface": "apple", @@ -13499,7 +13499,7 @@ }, { "kind": "conditional-branch", - "line": 3065, + "line": 3032, "path": "apps/ios/Sources/Voice/TalkModeManager.swift", "source": "Native WebRTC", "surface": "apple", diff --git a/apps/ios/Sources/Voice/RealtimeTalkRelaySession.swift b/apps/ios/Sources/Voice/RealtimeTalkRelaySession.swift index ab91c87aa01e..3d2f400bdeaa 100644 --- a/apps/ios/Sources/Voice/RealtimeTalkRelaySession.swift +++ b/apps/ios/Sources/Voice/RealtimeTalkRelaySession.swift @@ -18,7 +18,7 @@ private func makeRealtimeAudioTapBlock( targetSampleRate: targetSampleRate) guard !encoded.isEmpty else { return } let timestampMs = (ProcessInfo.processInfo.systemUptime * 1000).rounded() - let rms = RealtimeTalkRelaySession.rmsLevel(buffer: buffer) + let rms = Float(TalkAudioLevel.rms(buffer: buffer)) onAudio(encoded, timestampMs, rms) } } @@ -916,26 +916,6 @@ final class RealtimeTalkRelaySession { return data } - fileprivate nonisolated static func rmsLevel(buffer: AVAudioPCMBuffer) -> Float { - guard let channelData = buffer.floatChannelData, - buffer.frameLength > 0 - else { return 0 } - let frameCount = Int(buffer.frameLength) - let channelCount = max(1, Int(buffer.format.channelCount)) - var sumSquares: Float = 0 - var samples = 0 - for channel in 0.. 0 else { return 0 } - return sqrt(sumSquares / Float(samples)) - } - private nonisolated static func safeLogMessage(_ value: String) -> String { let singleLine = value .replacingOccurrences(of: "\n", with: " ") diff --git a/apps/ios/Sources/Voice/TalkGatewaySpeechClient.swift b/apps/ios/Sources/Voice/TalkGatewaySpeechClient.swift index 57dc92b4898e..f4b37a380ed9 100644 --- a/apps/ios/Sources/Voice/TalkGatewaySpeechClient.swift +++ b/apps/ios/Sources/Voice/TalkGatewaySpeechClient.swift @@ -127,161 +127,9 @@ extension TalkBufferedAudioPlaying { func setLevelHandler(_: (@MainActor (Double?) -> Void)?) {} } -@MainActor -final class TalkBufferedAudioPlayer: NSObject, TalkBufferedAudioPlaying, @preconcurrency AVAudioPlayerDelegate { - static let shared = TalkBufferedAudioPlayer() - - private final class Playback: @unchecked Sendable { - private let lock = NSLock() - private var finished = false - private var continuation: CheckedContinuation? - private var watchdog: Task? - - func setContinuation(_ continuation: CheckedContinuation) { - self.lock.lock() - defer { self.lock.unlock() } - self.continuation = continuation - } - - func setWatchdog(_ task: Task?) { - self.lock.lock() - let old = self.watchdog - self.watchdog = task - self.lock.unlock() - old?.cancel() - } - - func finish(_ result: StreamingPlaybackResult) { - let continuation: CheckedContinuation? - self.lock.lock() - if self.finished { - continuation = nil - } else { - self.finished = true - continuation = self.continuation - self.continuation = nil - } - self.lock.unlock() - continuation?.resume(returning: result) - } - } - - private let logger = Logger(subsystem: "ai.openclaw", category: "talk.tts") - private var player: AVAudioPlayer? - private var playback: Playback? - private var levelHandler: (@MainActor (Double?) -> Void)? - private var levelMeter: AudioPlayerLevelMeter? - - func setLevelHandler(_ handler: (@MainActor (Double?) -> Void)?) { - self.levelHandler = handler - } - - func play(data: Data) async -> StreamingPlaybackResult { - self.stopInternal() - - let playback = Playback() - self.playback = playback - return await withCheckedContinuation { continuation in - playback.setContinuation(continuation) - do { - let player = try AVAudioPlayer(data: data) - self.player = player - player.delegate = self - player.prepareToPlay() - if let levelHandler { - let meter = AudioPlayerLevelMeter(onLevel: levelHandler) - meter.attach(player) - self.levelMeter = meter - } - self.armWatchdog(playback: playback) - if !player.play() { - self.logger.error("talk buffered audio player refused to play") - self.finish(playback: playback, result: .init(finished: false, interruptedAt: nil)) - } - } catch { - self.logger.error("talk buffered audio player failed: \(error.localizedDescription, privacy: .public)") - self.finish(playback: playback, result: .init(finished: false, interruptedAt: nil)) - } - } - } - - func stop() -> Double? { - guard let player else { return nil } - let interruptedAt = player.currentTime - self.finish( - playback: self.playback, - result: .init(finished: false, interruptedAt: interruptedAt)) - return interruptedAt - } - - func audioPlayerDidFinishPlaying(_ player: AVAudioPlayer, successfully flag: Bool) { - self.finish( - playback: self.activePlayback(for: player), - result: .init(finished: flag, interruptedAt: nil)) - } - - func audioPlayerDecodeErrorDidOccur(_ player: AVAudioPlayer, error: (any Error)?) { - let message = error?.localizedDescription ?? "unknown decode error" - self.logger.error("talk buffered audio decode failed: \(message, privacy: .public)") - self.finish( - playback: self.activePlayback(for: player), - result: .init(finished: false, interruptedAt: nil)) - } - - private func activePlayback(for player: AVAudioPlayer) -> Playback? { - // AVAudioPlayer can deliver callbacks after stop/replacement. Keep a stale - // player from completing the current reply's continuation. - guard self.player === player else { return nil } - return self.playback - } - - private func stopInternal() { - if let player, let playback { - self.finish( - playback: playback, - result: .init(finished: false, interruptedAt: player.currentTime)) - return - } - self.player?.stop() - self.player = nil - } - - private func finish(playback: Playback?, result: StreamingPlaybackResult) { - guard let playback else { return } - playback.setWatchdog(nil) - playback.finish(result) - - guard self.playback === playback else { return } - self.playback = nil - self.levelMeter?.detach() - self.levelMeter = nil - self.player?.stop() - self.player = nil - } - - private func armWatchdog(playback: Playback) { - playback.setWatchdog(Task { @MainActor [weak self] in - guard let self else { return } - try? await Task.sleep(nanoseconds: 650_000_000) - guard !Task.isCancelled, self.playback === playback else { return } - guard self.player?.isPlaying == true else { - self.finish( - playback: playback, - result: .init(finished: false, interruptedAt: nil)) - return - } - - let duration = self.player?.duration ?? 0 - let timeoutSeconds = min(max(2.0, duration + 2.0), 5 * 60.0) - try? await Task.sleep(nanoseconds: UInt64(timeoutSeconds * 1_000_000_000)) - guard !Task.isCancelled, self.playback === playback else { return } - self.logger.error("talk buffered audio player watchdog completed unresolved playback") - self.finish( - playback: playback, - result: .init(finished: false, interruptedAt: nil)) - }) - } -} +/// Playback lives in OpenClawKit's TalkBufferedAudioPlayer, shared with the +/// macOS talk runtime; this file keeps only the iOS test seam conformance. +extension TalkBufferedAudioPlayer: TalkBufferedAudioPlaying {} private enum TalkGatewaySpeechError: LocalizedError { case invalidRequest diff --git a/apps/ios/Sources/Voice/TalkModeManager.swift b/apps/ios/Sources/Voice/TalkModeManager.swift index 887701febdab..45e4130da5bb 100644 --- a/apps/ios/Sources/Voice/TalkModeManager.swift +++ b/apps/ios/Sources/Voice/TalkModeManager.swift @@ -1786,37 +1786,18 @@ final class TalkModeManager: NSObject { self.startSpeechInterruptionRecognitionIfNeeded() self.statusText = "Speaking…" - let sampleRate = TalkTTSValidation.pcmSampleRate(from: outputFormat) - let result: StreamingPlaybackResult - if let sampleRate { - let streamFailure = StreamFailureBox() - let stream = Self.monitorStreamFailures(rawStream, failureBox: streamFailure) - self.lastPlaybackWasPCM = true - var playback = await pcmPlayer.play( - stream: self.meteredPCMStream(stream, sampleRate: sampleRate), - sampleRate: sampleRate) - self.pcmPlaybackEnvelope.cancel() - if !playback.finished, playback.interruptedAt == nil { - let mp3Format = ElevenLabsTTSClient.validatedOutputFormat("mp3_44100_128") - self.logger.warning("pcm playback failed; retrying mp3") - if Self.isPCMFormatRejectedByAPI(streamFailure.value) { - self.pcmFormatUnavailable = true - } - self.lastPlaybackWasPCM = false - let mp3Stream = client.streamSynthesize( - voiceId: voiceId, - request: self.makeElevenLabsTTSRequest( - text: cleaned, - directive: directive, - modelId: modelId, - outputFormat: mp3Format, - language: language)) - playback = await self.mp3Player.play(stream: mp3Stream) - } - result = playback - } else { - self.lastPlaybackWasPCM = false - result = await self.mp3Player.play(stream: rawStream) + let result = await self.playElevenLabsStream( + rawStream, + sampleRate: TalkTTSValidation.pcmSampleRate(from: outputFormat)) + { mp3Format in + client.streamSynthesize( + voiceId: voiceId, + request: self.makeElevenLabsTTSRequest( + text: cleaned, + directive: directive, + modelId: modelId, + outputFormat: mp3Format, + language: language)) } let duration = Date().timeIntervalSince(started) self.logger @@ -1892,7 +1873,7 @@ final class TalkModeManager: NSObject { self.lastPlaybackWasPCM = true let stream = Self.makeBufferedAudioStream(chunks: [audio.data]) result = await self.pcmPlayer.play( - stream: self.meteredPCMStream(stream, sampleRate: sampleRate), + stream: self.pcmPlaybackEnvelope.metering(stream, sampleRate: sampleRate), sampleRate: sampleRate) self.pcmPlaybackEnvelope.cancel() case .buffered: @@ -2452,32 +2433,6 @@ final class TalkModeManager: NSObject { } } - /// Passes PCM chunks through to the player while feeding the playback - /// envelope so `playbackLevel` follows the audible speech. - private func meteredPCMStream( - _ stream: AsyncThrowingStream, - sampleRate: Double) -> AsyncThrowingStream - { - let envelope = self.pcmPlaybackEnvelope - envelope.begin(sampleRate: sampleRate) - return AsyncThrowingStream { continuation in - let task = Task { @MainActor in - do { - for try await chunk in stream { - envelope.append(chunk) - continuation.yield(chunk) - } - continuation.finish() - } catch { - continuation.finish(throwing: error) - } - } - continuation.onTermination = { _ in - task.cancel() - } - } - } - private func speakIncrementalSegment( _ text: String, context preferredContext: IncrementalSpeechContext? = nil, @@ -2515,40 +2470,52 @@ final class TalkModeManager: NSObject { client.streamSynthesize(voiceId: voiceId, request: request) } let playbackFormat = prefetchedAudio?.outputFormat ?? context.outputFormat - let sampleRate = TalkTTSValidation.pcmSampleRate(from: playbackFormat) - let result: StreamingPlaybackResult - if let sampleRate { - let streamFailure = StreamFailureBox() - let stream = Self.monitorStreamFailures(rawStream, failureBox: streamFailure) - self.lastPlaybackWasPCM = true - var playback = await pcmPlayer.play( - stream: self.meteredPCMStream(stream, sampleRate: sampleRate), - sampleRate: sampleRate) - self.pcmPlaybackEnvelope.cancel() - if !playback.finished, playback.interruptedAt == nil { - self.logger.warning("pcm playback failed; retrying mp3") - if Self.isPCMFormatRejectedByAPI(streamFailure.value) { - self.pcmFormatUnavailable = true - } - self.lastPlaybackWasPCM = false - let mp3Format = ElevenLabsTTSClient.validatedOutputFormat("mp3_44100_128") - let mp3Stream = client.streamSynthesize( - voiceId: voiceId, - request: self.makeIncrementalTTSRequest( - text: text, - context: context, - outputFormat: mp3Format)) - playback = await self.mp3Player.play(stream: mp3Stream) - } - result = playback - } else { - self.lastPlaybackWasPCM = false - result = await self.mp3Player.play(stream: rawStream) + let result = await self.playElevenLabsStream( + rawStream, + sampleRate: TalkTTSValidation.pcmSampleRate(from: playbackFormat)) + { mp3Format in + client.streamSynthesize( + voiceId: voiceId, + request: self.makeIncrementalTTSRequest( + text: text, + context: context, + outputFormat: mp3Format)) } if !result.finished, let interruptedAt = result.interruptedAt { self.lastInterruptedAtSeconds = interruptedAt } } + + /// Plays an ElevenLabs stream: metered PCM when the output format is raw + /// PCM, retried once as mp3 when PCM playback fails outright (some plans + /// and formats reject PCM); plain mp3 streaming otherwise. + private func playElevenLabsStream( + _ rawStream: AsyncThrowingStream, + sampleRate: Double?, + makeMP3Stream: (String?) -> AsyncThrowingStream) async -> StreamingPlaybackResult + { + guard let sampleRate else { + self.lastPlaybackWasPCM = false + return await self.mp3Player.play(stream: rawStream) + } + let streamFailure = StreamFailureBox() + let stream = Self.monitorStreamFailures(rawStream, failureBox: streamFailure) + self.lastPlaybackWasPCM = true + var playback = await pcmPlayer.play( + stream: self.pcmPlaybackEnvelope.metering(stream, sampleRate: sampleRate), + sampleRate: sampleRate) + self.pcmPlaybackEnvelope.cancel() + if !playback.finished, playback.interruptedAt == nil { + self.logger.warning("pcm playback failed; retrying mp3") + if Self.isPCMFormatRejectedByAPI(streamFailure.value) { + self.pcmFormatUnavailable = true + } + self.lastPlaybackWasPCM = false + let mp3Format = ElevenLabsTTSClient.validatedOutputFormat("mp3_44100_128") + playback = await self.mp3Player.play(stream: makeMP3Stream(mp3Format)) + } + return playback + } } private struct IncrementalSpeechBuffer { @@ -3286,20 +3253,7 @@ private final class AudioTapDiagnostics: @unchecked Sendable { let ch = buffer.format.channelCount let frames = buffer.frameLength - var rms: Float? - if let data = buffer.floatChannelData?.pointee { - let n = Int(frames) - if n > 0 { - var sum: Float = 0 - for i in 0.. self.maxRmsWindow { self.maxRmsWindow = resolvedRms } diff --git a/apps/macos/Sources/OpenClaw/MicLevelMonitor.swift b/apps/macos/Sources/OpenClaw/MicLevelMonitor.swift index fc44ac614396..93ce0262764c 100644 --- a/apps/macos/Sources/OpenClaw/MicLevelMonitor.swift +++ b/apps/macos/Sources/OpenClaw/MicLevelMonitor.swift @@ -1,4 +1,5 @@ import AVFoundation +import OpenClawKit import OSLog import SwiftUI @@ -41,7 +42,7 @@ actor MicLevelMonitor { input.removeTap(onBus: 0) input.installTap(onBus: 0, bufferSize: 512, format: format) { [weak self] buffer, _ in guard let self else { return } - let level = Self.normalizedLevel(from: buffer) + let level = TalkAudioLevel.normalized(rms: TalkAudioLevel.rms(buffer: buffer)) Task { await self.push(level: level) } } engine.prepare() @@ -71,20 +72,6 @@ actor MicLevelMonitor { self.lastPublishedLevel = value Task { @MainActor in update(value) } } - - private static func normalizedLevel(from buffer: AVAudioPCMBuffer) -> Double { - guard let channel = buffer.floatChannelData?[0] else { return 0 } - let frameCount = Int(buffer.frameLength) - guard frameCount > 0 else { return 0 } - var sum: Float = 0 - for i in 0.. Void)? - private var levelMeter: AudioPlayerLevelMeter? - - func setLevelHandler(_ handler: (@MainActor (Double?) -> Void)?) { - self.levelHandler = handler - } - - private final class Playback: @unchecked Sendable { - private let lock = NSLock() - private var finished = false - private var continuation: CheckedContinuation? - private var watchdog: Task? - - func setContinuation(_ continuation: CheckedContinuation) { - self.lock.lock() - defer { self.lock.unlock() } - self.continuation = continuation - } - - func setWatchdog(_ task: Task?) { - self.lock.lock() - let old = self.watchdog - self.watchdog = task - self.lock.unlock() - old?.cancel() - } - - func cancelWatchdog() { - self.setWatchdog(nil) - } - - func finish(_ result: TalkPlaybackResult) { - let continuation: CheckedContinuation? - self.lock.lock() - if self.finished { - continuation = nil - } else { - self.finished = true - continuation = self.continuation - self.continuation = nil - } - self.lock.unlock() - continuation?.resume(returning: result) - } - } - - func play(data: Data) async -> TalkPlaybackResult { - self.stopInternal() - - let playback = Playback() - self.playback = playback - - return await withCheckedContinuation { continuation in - playback.setContinuation(continuation) - do { - let player = try AVAudioPlayer(data: data) - self.player = player - - player.delegate = self - player.prepareToPlay() - if let levelHandler { - let meter = AudioPlayerLevelMeter(onLevel: levelHandler) - meter.attach(player) - self.levelMeter = meter - } - - self.armWatchdog(playback: playback) - - let ok = player.play() - if !ok { - self.logger.error("talk audio player refused to play") - self.finish(playback: playback, result: TalkPlaybackResult(finished: false, interruptedAt: nil)) - } - } catch { - self.logger.error("talk audio player failed: \(error.localizedDescription, privacy: .public)") - self.finish(playback: playback, result: TalkPlaybackResult(finished: false, interruptedAt: nil)) - } - } - } - - func stop() -> Double? { - guard let player else { return nil } - let time = player.currentTime - self.stopInternal(interruptedAt: time) - return time - } - - func audioPlayerDidFinishPlaying(_: AVAudioPlayer, successfully flag: Bool) { - self.stopInternal(finished: flag) - } - - private func stopInternal(finished: Bool = false, interruptedAt: Double? = nil) { - guard let playback else { return } - let result = TalkPlaybackResult(finished: finished, interruptedAt: interruptedAt) - self.finish(playback: playback, result: result) - } - - private func finish(playback: Playback, result: TalkPlaybackResult) { - playback.cancelWatchdog() - playback.finish(result) - - guard self.playback === playback else { return } - self.playback = nil - self.levelMeter?.detach() - self.levelMeter = nil - self.player?.stop() - self.player = nil - } - - private func stopInternal() { - if let playback = self.playback { - let interruptedAt = self.player?.currentTime - self.finish( - playback: playback, - result: TalkPlaybackResult(finished: false, interruptedAt: interruptedAt)) - return - } - self.player?.stop() - self.player = nil - } - - private func armWatchdog(playback: Playback) { - playback.setWatchdog(Task { @MainActor [weak self] in - guard let self else { return } - - do { - try await Task.sleep(nanoseconds: 650_000_000) - } catch { - return - } - if Task.isCancelled { return } - - guard self.playback === playback else { return } - if self.player?.isPlaying != true { - self.logger.error("talk audio player did not start playing") - self.finish(playback: playback, result: TalkPlaybackResult(finished: false, interruptedAt: nil)) - return - } - - let duration = self.player?.duration ?? 0 - let timeoutSeconds = min(max(2.0, duration + 2.0), 5 * 60.0) - do { - try await Task.sleep(nanoseconds: UInt64(timeoutSeconds * 1_000_000_000)) - } catch { - return - } - if Task.isCancelled { return } - - guard self.playback === playback else { return } - guard self.player?.isPlaying == true else { return } - self.logger.error("talk audio player watchdog fired") - self.finish(playback: playback, result: TalkPlaybackResult(finished: false, interruptedAt: nil)) - }) - } -} - -struct TalkPlaybackResult { - let finished: Bool - let interruptedAt: Double? -} diff --git a/apps/macos/Sources/OpenClaw/TalkModeController.swift b/apps/macos/Sources/OpenClaw/TalkModeController.swift index 24bfc467e072..acefc7f03add 100644 --- a/apps/macos/Sources/OpenClaw/TalkModeController.swift +++ b/apps/macos/Sources/OpenClaw/TalkModeController.swift @@ -92,24 +92,7 @@ final class TalkModeController { _ stream: AsyncThrowingStream, sampleRate: Double) -> AsyncThrowingStream { - let envelope = self.playbackEnvelope - envelope.begin(sampleRate: sampleRate) - return AsyncThrowingStream { continuation in - let task = Task { @MainActor in - do { - for try await chunk in stream { - envelope.append(chunk) - continuation.yield(chunk) - } - continuation.finish() - } catch { - continuation.finish(throwing: error) - } - } - continuation.onTermination = { _ in - task.cancel() - } - } + self.playbackEnvelope.metering(stream, sampleRate: sampleRate) } func endSpeechMetering() { diff --git a/apps/macos/Sources/OpenClaw/TalkModeRuntime.swift b/apps/macos/Sources/OpenClaw/TalkModeRuntime.swift index d841433625c3..4456f9f5e461 100644 --- a/apps/macos/Sources/OpenClaw/TalkModeRuntime.swift +++ b/apps/macos/Sources/OpenClaw/TalkModeRuntime.swift @@ -227,9 +227,7 @@ actor TalkModeRuntime { let meter = self.rmsMeter input.installTap(onBus: 0, bufferSize: 2048, format: format) { [weak request, meter] buffer, _ in request?.append(SpeechAudioBufferNormalizer.speechCompatibleBuffer(from: buffer)) - if let rms = Self.rmsLevel(buffer: buffer) { - meter.set(rms) - } + meter.set(TalkAudioLevel.rms(buffer: buffer)) } audioEngine.prepare() @@ -1102,16 +1100,16 @@ extension TalkModeRuntime { } @MainActor - private func playTalkAudio(data: Data) async -> TalkPlaybackResult { - TalkAudioPlayer.shared.setLevelHandler { level in + private func playTalkAudio(data: Data) async -> StreamingPlaybackResult { + TalkBufferedAudioPlayer.shared.setLevelHandler { level in TalkModeController.shared.updateSpeakingLevel(level) } - return await TalkAudioPlayer.shared.play(data: data) + return await TalkBufferedAudioPlayer.shared.play(data: data) } @MainActor private func stopTalkAudio() -> Double? { - TalkAudioPlayer.shared.stop() + TalkBufferedAudioPlayer.shared.stop() } private func synthesizeMLXVoice( @@ -1255,18 +1253,6 @@ extension TalkModeRuntime { } } - private static func rmsLevel(buffer: AVAudioPCMBuffer) -> Double? { - guard let channelData = buffer.floatChannelData?.pointee else { return nil } - let frameCount = Int(buffer.frameLength) - guard frameCount > 0 else { return nil } - var sum: Double = 0 - for i in 0.. Bool { let trimmed = transcript.trimmingCharacters(in: .whitespacesAndNewlines) guard trimmed.count >= 3 else { return false } diff --git a/apps/macos/Sources/OpenClaw/VoiceWakeRuntime.swift b/apps/macos/Sources/OpenClaw/VoiceWakeRuntime.swift index ea52819ad6d0..a4dac2cfcc1c 100644 --- a/apps/macos/Sources/OpenClaw/VoiceWakeRuntime.swift +++ b/apps/macos/Sources/OpenClaw/VoiceWakeRuntime.swift @@ -1,5 +1,6 @@ import AVFoundation import Foundation +import OpenClawKit import OSLog import Speech import SwabbleKit @@ -188,7 +189,7 @@ actor VoiceWakeRuntime { input.removeTap(onBus: 0) input.installTap(onBus: 0, bufferSize: 2048, format: format) { [weak self, weak request] buffer, _ in request?.append(SpeechAudioBufferNormalizer.speechCompatibleBuffer(from: buffer)) - guard let rms = Self.rmsLevel(buffer: buffer) else { return } + let rms = TalkAudioLevel.rms(buffer: buffer) Task.detached { [weak self] in await self?.noteAudioLevel(rms: rms) await self?.noteAudioTap(rms: rms) @@ -726,18 +727,6 @@ actor VoiceWakeRuntime { } } - private static func rmsLevel(buffer: AVAudioPCMBuffer) -> Double? { - guard let channelData = buffer.floatChannelData?.pointee else { return nil } - let frameCount = Int(buffer.frameLength) - guard frameCount > 0 else { return nil } - var sum: Double = 0 - for i in 0..? + private var watchdog: Task? + + func setContinuation(_ continuation: CheckedContinuation) { + self.lock.lock() + defer { self.lock.unlock() } + self.continuation = continuation + } + + func setWatchdog(_ task: Task?) { + self.lock.lock() + let old = self.watchdog + self.watchdog = task + self.lock.unlock() + old?.cancel() + } + + func finish(_ result: StreamingPlaybackResult) { + let continuation: CheckedContinuation? + self.lock.lock() + if self.finished { + continuation = nil + } else { + self.finished = true + continuation = self.continuation + self.continuation = nil + } + self.lock.unlock() + continuation?.resume(returning: result) + } + } + + private let logger = Logger(subsystem: "ai.openclaw", category: "talk.tts") + private var player: AVAudioPlayer? + private var playback: Playback? + private var levelHandler: (@MainActor (Double?) -> Void)? + private var levelMeter: AudioPlayerLevelMeter? + + public func setLevelHandler(_ handler: (@MainActor (Double?) -> Void)?) { + self.levelHandler = handler + } + + public func play(data: Data) async -> StreamingPlaybackResult { + self.stopInternal() + + let playback = Playback() + self.playback = playback + return await withCheckedContinuation { continuation in + playback.setContinuation(continuation) + do { + let player = try AVAudioPlayer(data: data) + self.player = player + player.delegate = self + player.prepareToPlay() + if let levelHandler { + let meter = AudioPlayerLevelMeter(onLevel: levelHandler) + meter.attach(player) + self.levelMeter = meter + } + self.armWatchdog(playback: playback) + if !player.play() { + self.logger.error("talk buffered audio player refused to play") + self.finish(playback: playback, result: .init(finished: false, interruptedAt: nil)) + } + } catch { + self.logger.error("talk buffered audio player failed: \(error.localizedDescription, privacy: .public)") + self.finish(playback: playback, result: .init(finished: false, interruptedAt: nil)) + } + } + } + + public func stop() -> Double? { + guard let player else { return nil } + let interruptedAt = player.currentTime + self.finish( + playback: self.playback, + result: .init(finished: false, interruptedAt: interruptedAt)) + return interruptedAt + } + + public func audioPlayerDidFinishPlaying(_ player: AVAudioPlayer, successfully flag: Bool) { + self.finish( + playback: self.activePlayback(for: player), + result: .init(finished: flag, interruptedAt: nil)) + } + + public func audioPlayerDecodeErrorDidOccur(_ player: AVAudioPlayer, error: (any Error)?) { + let message = error?.localizedDescription ?? "unknown decode error" + self.logger.error("talk buffered audio decode failed: \(message, privacy: .public)") + self.finish( + playback: self.activePlayback(for: player), + result: .init(finished: false, interruptedAt: nil)) + } + + private func activePlayback(for player: AVAudioPlayer) -> Playback? { + // AVAudioPlayer can deliver callbacks after stop/replacement. Keep a stale + // player from completing the current reply's continuation. + guard self.player === player else { return nil } + return self.playback + } + + private func stopInternal() { + if let player, let playback { + self.finish( + playback: playback, + result: .init(finished: false, interruptedAt: player.currentTime)) + return + } + self.player?.stop() + self.player = nil + } + + private func finish(playback: Playback?, result: StreamingPlaybackResult) { + guard let playback else { return } + playback.setWatchdog(nil) + playback.finish(result) + + guard self.playback === playback else { return } + self.playback = nil + self.levelMeter?.detach() + self.levelMeter = nil + self.player?.stop() + self.player = nil + } + + private func armWatchdog(playback: Playback) { + playback.setWatchdog(Task { @MainActor [weak self] in + guard let self else { return } + try? await Task.sleep(nanoseconds: 650_000_000) + guard !Task.isCancelled, self.playback === playback else { return } + guard self.player?.isPlaying == true else { + self.finish( + playback: playback, + result: .init(finished: false, interruptedAt: nil)) + return + } + + let duration = self.player?.duration ?? 0 + let timeoutSeconds = min(max(2.0, duration + 2.0), 5 * 60.0) + try? await Task.sleep(nanoseconds: UInt64(timeoutSeconds * 1_000_000_000)) + guard !Task.isCancelled, self.playback === playback else { return } + self.logger.error("talk buffered audio player watchdog completed unresolved playback") + self.finish( + playback: playback, + result: .init(finished: false, interruptedAt: nil)) + }) + } +} +#endif diff --git a/apps/shared/OpenClawKit/Sources/OpenClawKit/TalkPlaybackLevelMeters.swift b/apps/shared/OpenClawKit/Sources/OpenClawKit/TalkPlaybackLevelMeters.swift index aba4b0cf4503..7e7a67a7eb55 100644 --- a/apps/shared/OpenClawKit/Sources/OpenClawKit/TalkPlaybackLevelMeters.swift +++ b/apps/shared/OpenClawKit/Sources/OpenClawKit/TalkPlaybackLevelMeters.swift @@ -13,6 +13,24 @@ public enum TalkAudioLevel { max(0, min(1, (decibels + 50) / 50)) } + /// Average RMS across all channels of a float PCM buffer; 0 for degenerate + /// buffers (Core Audio taps never deliver them in practice). Callable from + /// realtime audio tap threads. + public static func rms(buffer: AVAudioPCMBuffer) -> Double { + guard let channelData = buffer.floatChannelData, buffer.frameLength > 0 else { return 0 } + let frameCount = Int(buffer.frameLength) + let channelCount = max(1, Int(buffer.format.channelCount)) + var sum: Double = 0 + for channel in 0.. Double { let sampleCount = data.count / 2 @@ -89,6 +107,32 @@ public final class PCMPlaybackEnvelope { self.scheduleEnd = start } + /// Passes PCM chunks through to a player while metering them into this + /// envelope, so the published level follows the audible speech. Callers + /// `cancel()` once playback returns. + public func metering( + _ stream: AsyncThrowingStream, + sampleRate: Double) -> AsyncThrowingStream + { + self.begin(sampleRate: sampleRate) + return AsyncThrowingStream { continuation in + let task = Task { @MainActor [weak self] in + do { + for try await chunk in stream { + self?.append(chunk) + continuation.yield(chunk) + } + continuation.finish() + } catch { + continuation.finish(throwing: error) + } + } + continuation.onTermination = { _ in + task.cancel() + } + } + } + /// Stops immediately (interruption/teardown) and clears the published level. public func cancel() { self.publishTask?.cancel() diff --git a/apps/shared/OpenClawKit/Tests/OpenClawKitTests/TalkWaveformTests.swift b/apps/shared/OpenClawKit/Tests/OpenClawKitTests/TalkWaveformTests.swift index f3e6d7aa2811..4f674a11e6bf 100644 --- a/apps/shared/OpenClawKit/Tests/OpenClawKitTests/TalkWaveformTests.swift +++ b/apps/shared/OpenClawKit/Tests/OpenClawKitTests/TalkWaveformTests.swift @@ -1,3 +1,4 @@ +import AVFoundation import Foundation import OpenClawChatUI import OpenClawKit @@ -87,4 +88,24 @@ struct TalkAudioLevelTests { #expect(TalkAudioLevel.pcm16RMS(Data()) == 0) #expect(TalkAudioLevel.pcm16RMS(Data([0x7F])) == 0) } + + @Test + func `buffer RMS averages float samples across channels`() throws { + let format = try #require(AVAudioFormat( + commonFormat: .pcmFormatFloat32, + sampleRate: 16000, + channels: 2, + interleaved: false)) + let buffer = try #require(AVAudioPCMBuffer(pcmFormat: format, frameCapacity: 128)) + buffer.frameLength = 128 + let channels = try #require(buffer.floatChannelData) + for index in 0..<128 { + channels[0][index] = 0.5 + channels[1][index] = -0.5 + } + #expect(abs(TalkAudioLevel.rms(buffer: buffer) - 0.5) < 1e-6) + + buffer.frameLength = 0 + #expect(TalkAudioLevel.rms(buffer: buffer) == 0) + } }