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
This commit is contained in:
Peter Steinberger
2026-07-09 22:21:00 +01:00
committed by GitHub
parent f8efe14639
commit 4aa8bf12e3
13 changed files with 314 additions and 522 deletions
+2 -2
View File
@@ -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",
@@ -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..<channelCount {
let values = channelData[channel]
for index in 0..<frameCount {
let sample = values[index]
sumSquares += sample * sample
samples += 1
}
}
guard samples > 0 else { return 0 }
return sqrt(sumSquares / Float(samples))
}
private nonisolated static func safeLogMessage(_ value: String) -> String {
let singleLine = value
.replacingOccurrences(of: "\n", with: " ")
@@ -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<StreamingPlaybackResult, Never>?
private var watchdog: Task<Void, Never>?
func setContinuation(_ continuation: CheckedContinuation<StreamingPlaybackResult, Never>) {
self.lock.lock()
defer { self.lock.unlock() }
self.continuation = continuation
}
func setWatchdog(_ task: Task<Void, Never>?) {
self.lock.lock()
let old = self.watchdog
self.watchdog = task
self.lock.unlock()
old?.cancel()
}
func finish(_ result: StreamingPlaybackResult) {
let continuation: CheckedContinuation<StreamingPlaybackResult, Never>?
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
+55 -101
View File
@@ -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<Data, Error>,
sampleRate: Double) -> AsyncThrowingStream<Data, Error>
{
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<Data, Error>,
sampleRate: Double?,
makeMP3Stream: (String?) -> AsyncThrowingStream<Data, Error>) 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..<n {
let v = data[i]
sum += v * v
}
rms = sqrt(sum / Float(n))
}
}
let resolvedRms = rms ?? 0
let resolvedRms = Float(TalkAudioLevel.rms(buffer: buffer))
self.lock.lock()
self.lastRms = resolvedRms
if resolvedRms > self.maxRmsWindow { self.maxRmsWindow = resolvedRms }
@@ -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..<frameCount {
let s = channel[i]
sum += s * s
}
let rms = sqrt(sum / Float(frameCount) + 1e-12)
let db = 20 * log10(Double(rms))
return max(0, min(1, (db + 50) / 50))
}
}
struct MicLevelBar: View {
@@ -1,172 +0,0 @@
import AVFoundation
import Foundation
import OpenClawKit
import OSLog
@MainActor
final class TalkAudioPlayer: NSObject, @preconcurrency AVAudioPlayerDelegate {
static let shared = TalkAudioPlayer()
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
}
private final class Playback: @unchecked Sendable {
private let lock = NSLock()
private var finished = false
private var continuation: CheckedContinuation<TalkPlaybackResult, Never>?
private var watchdog: Task<Void, Never>?
func setContinuation(_ continuation: CheckedContinuation<TalkPlaybackResult, Never>) {
self.lock.lock()
defer { self.lock.unlock() }
self.continuation = continuation
}
func setWatchdog(_ task: Task<Void, Never>?) {
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<TalkPlaybackResult, Never>?
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?
}
@@ -92,24 +92,7 @@ final class TalkModeController {
_ stream: AsyncThrowingStream<Data, Error>,
sampleRate: Double) -> AsyncThrowingStream<Data, Error>
{
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() {
@@ -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..<frameCount {
let sample = Double(channelData[i])
sum += sample * sample
}
return sqrt(sum / Double(frameCount))
}
private func shouldInterrupt(transcript: String, hasConfidence: Bool) async -> Bool {
let trimmed = transcript.trimmingCharacters(in: .whitespacesAndNewlines)
guard trimmed.count >= 3 else { return false }
@@ -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..<frameCount {
let sample = Double(channelData[i])
sum += sample * sample
}
return sqrt(sum / Double(frameCount))
}
private func restartRecognizer() {
// Restart the recognizer so we listen for the next trigger with a clean buffer.
let current = self.currentConfig
@@ -1,15 +1,16 @@
import Foundation
import OpenClawKit
import Testing
@testable import OpenClaw
@Suite(.serialized) struct TalkAudioPlayerTests {
@Suite(.serialized) struct TalkBufferedAudioPlayerTests {
@MainActor
@Test func `play does not hang when playback ends or fails`() async throws {
let wav = makeWav16Mono(sampleRate: 8000, samples: 80)
defer { _ = TalkAudioPlayer.shared.stop() }
defer { _ = TalkBufferedAudioPlayer.shared.stop() }
_ = try await withTimeout(seconds: 10.0) {
await TalkAudioPlayer.shared.play(data: wav)
await TalkBufferedAudioPlayer.shared.play(data: wav)
}
#expect(true)
@@ -18,14 +19,14 @@ import Testing
@MainActor
@Test func `play does not hang when play is called twice`() async throws {
let wav = makeWav16Mono(sampleRate: 8000, samples: 800)
defer { _ = TalkAudioPlayer.shared.stop() }
defer { _ = TalkBufferedAudioPlayer.shared.stop() }
let first = Task { @MainActor in
await TalkAudioPlayer.shared.play(data: wav)
await TalkBufferedAudioPlayer.shared.play(data: wav)
}
await Task.yield()
_ = await TalkAudioPlayer.shared.play(data: wav)
_ = await TalkBufferedAudioPlayer.shared.play(data: wav)
_ = try await withTimeout(seconds: 10.0) {
await first.value
@@ -0,0 +1,171 @@
// canImport mirrors ElevenLabsKitShim: StreamingPlaybackResult only exists on
// platforms ElevenLabsKit builds for (iOS/macOS); the watch build compiles this out.
#if Talk && canImport(ElevenLabsKit)
import AVFoundation
import Foundation
import OSLog
/// Plays one complete TTS clip (MP3/WAV/FLAC container) at a time via
/// AVAudioPlayer, with live level metering and a watchdog so a stalled or
/// silently failing playback can never hang the talk loop. Shared by the iOS
/// and macOS talk runtimes.
@MainActor
public final class TalkBufferedAudioPlayer: NSObject, @preconcurrency AVAudioPlayerDelegate {
public static let shared = TalkBufferedAudioPlayer()
override public init() {
super.init()
}
private final class Playback: @unchecked Sendable {
private let lock = NSLock()
private var finished = false
private var continuation: CheckedContinuation<StreamingPlaybackResult, Never>?
private var watchdog: Task<Void, Never>?
func setContinuation(_ continuation: CheckedContinuation<StreamingPlaybackResult, Never>) {
self.lock.lock()
defer { self.lock.unlock() }
self.continuation = continuation
}
func setWatchdog(_ task: Task<Void, Never>?) {
self.lock.lock()
let old = self.watchdog
self.watchdog = task
self.lock.unlock()
old?.cancel()
}
func finish(_ result: StreamingPlaybackResult) {
let continuation: CheckedContinuation<StreamingPlaybackResult, Never>?
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
@@ -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..<channelCount {
let samples = channelData[channel]
for index in 0..<frameCount {
let sample = Double(samples[index])
sum += sample * sample
}
}
return (sum / Double(frameCount * channelCount)).squareRoot()
}
/// RMS of little-endian PCM16 mono bytes; 0 for empty or odd-length data.
public static func pcm16RMS(_ data: Data) -> 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<Data, Error>,
sampleRate: Double) -> AsyncThrowingStream<Data, Error>
{
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()
@@ -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)
}
}