fix(ios): play realtime voice replies reliably with xAI (#163146)

Keep native realtime voice replies complete under UI load by moving playback and output routing off the main actor, retaining generation/cancellation fences, and accepting bounded provider bursts. Install acknowledgment identity before replaying startup output.

Related: #163122. Transcript and confirmation work remain separate.

Verified exact-head CI run 36984563290, including native Swift suites and iOS smoke. Contributor supplied physical-iPhone before/after playback evidence; Bluetooth route changes are not certified.

Landing session: https://team.openclaw.ai/chat/roboclaw/dashboard/060f4613-1c9e-4402-8400-66113d2e4463

Co-authored-by: Tak Hoffman <781889+Takhoffman@users.noreply.github.com>
Co-authored-by: Takhoffman <781889+Takhoffman@users.noreply.github.com>
Co-authored-by: Marvinthebored <peter@lindsey.jp>
This commit is contained in:
Marvinthebored 2026-10-02 23:01:14 +08:00 • committed by GitHub
parent 7555e3f881
commit 646916d1c2
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
9 changed files with 1616 additions and 483 deletions

View file

@ -3467,8 +3467,8 @@
{"id":"native.apple.716cff6f0c9637d1","source":"Realtime Voice, Gateway relay ready","surface":"apple","sites":[{"kind":"ui-localized-call","path":"apps/ios/Sources/Voice/TalkModeManager.swift"}]},
{"id":"native.apple.3459811c080ed880","source":"Realtime audio failed: %@","surface":"apple","sites":[{"kind":"ui-localized-call","path":"apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift"}]},
{"id":"native.apple.4f81777177325153","source":"Realtime audio input fell behind. Reconnecting…","surface":"apple","sites":[{"kind":"ui-localized-call","path":"apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift"}]},
{"id":"native.apple.1d3037a7bc3ec7cb","source":"Realtime audio playback failed. Reconnecting…","surface":"apple","sites":[{"kind":"ui-localized-call","path":"apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift"}]},
{"id":"native.apple.90f64f78a77015d7","source":"Realtime audio playback fell behind. Reconnecting…","surface":"apple","sites":[{"kind":"ui-localized-call","path":"apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift"}]},
{"id":"native.apple.1d3037a7bc3ec7cb","source":"Realtime audio playback failed. Reconnecting…","surface":"apple","sites":[{"kind":"ui-localized-call","path":"apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkOutput.swift"}]},
{"id":"native.apple.90f64f78a77015d7","source":"Realtime audio playback fell behind. Reconnecting…","surface":"apple","sites":[{"kind":"ui-localized-call","path":"apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkOutput.swift"}]},
{"id":"native.apple.f40c2764425389ef","source":"Realtime closed before it became ready.","surface":"apple","sites":[{"kind":"ui-localized-call","path":"apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift"}]},
{"id":"native.apple.59d22abf27122e23","source":"Realtime connection ended before it became ready.","surface":"apple","sites":[{"kind":"ui-localized-call","path":"apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift"}]},
{"id":"native.apple.b236c658fc54cbe5","source":"Realtime did not become ready in time.","surface":"apple","sites":[{"kind":"ui-localized-call","path":"apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift"}]},

View file

@ -2,14 +2,34 @@
import AVFAudio
import Foundation
@MainActor
public final class RealtimePCMStreamingAudioPlayer: PCMStreamingAudioPlaying {
/// Synchronous registration lets the relay retire a generation before any delayed task runs.
protocol RealtimePCMPlayback: Sendable {
func beginPlayback(stream: AsyncThrowingStream<Data, Error>, sampleRate: Double)
-> Task<StreamingPlaybackResult, Never>
func stop() -> Double?
}
/// The lock guards playback/generation state; backend closures run on backendQueue without it.
/// Completion callbacks are queued, including callbacks invoked synchronously by node.stop.
public final nonisolated class RealtimePCMStreamingAudioPlayer: PCMStreamingAudioPlaying, RealtimePCMPlayback,
@unchecked Sendable
{
private let lock = NSLock()
/// Backend operations never hold the state lock or block a relay control. Ordering on this
/// queue ensures a retired generation's stop precedes preparation of its replacement.
private let backendQueue = DispatchQueue(label: "ai.openclaw.realtime-playback", qos: .userInitiated)
static let frameDurationSeconds = 0.020
static let maxScheduledBuffers = 3
/// Bound scheduled audio to the relay's 60 s reply limit; completions refill off-main.
static let maxScheduledBuffers = 3000
typealias Completion = @Sendable () -> Void
private let preparePlayback: (Double) throws -> Void
private let scheduleFrame: (Data, Double, @escaping Completion) throws -> Void
private let startPlayback: () -> Void
/// Frames queued before the node starts: 300 ms of cushion so the first seconds of a reply,
/// which arrive while the UI is busy, don't underrun.
static let prebufferFrames = 15
private var playbackStarted = false
private let stopPlayback: () -> Void
private let playbackTime: () -> Double?
@ -17,7 +37,7 @@ public final class RealtimePCMStreamingAudioPlayer: PCMStreamingAudioPlaying {
private var nextBufferID: UInt64 = 0
private var scheduledBufferIDs: Set<UInt64> = []
private var slotWaiters: [CheckedContinuation<Bool, Never>] = []
private var playbackContinuation: CheckedContinuation<StreamingPlaybackResult, Never>?
private var playbackContinuation: AsyncStream<StreamingPlaybackResult>.Continuation?
private var inputTask: Task<Void, Never>?
private var inputFinished = false
@ -43,7 +63,13 @@ public final class RealtimePCMStreamingAudioPlayer: PCMStreamingAudioPlaying {
engine.connect(node, to: engine.mainMixerNode, format: nextFormat)
engine.prepare()
try engine.start()
node.play()
// Preserve the device-startup cushion before the first audible frames.
let padFrames = AVAudioFrameCount(sampleRate * 0.3)
if let pad = AVAudioPCMBuffer(pcmFormat: nextFormat, frameCapacity: padFrames) {
pad.frameLength = padFrames
pad.int16ChannelData?[0].update(repeating: 0, count: Int(padFrames))
node.scheduleBuffer(pad)
}
},
scheduleFrame: { data, _, completion in
guard let format else {
@ -65,6 +91,7 @@ public final class RealtimePCMStreamingAudioPlayer: PCMStreamingAudioPlaying {
completionCallbackType: .dataPlayedBack)
{ _ in completion() }
},
startPlayback: { node.play() },
stopPlayback: {
node.stop()
engine.stop()
@ -80,11 +107,13 @@ public final class RealtimePCMStreamingAudioPlayer: PCMStreamingAudioPlaying {
init(
preparePlayback: @escaping (Double) throws -> Void,
scheduleFrame: @escaping (Data, Double, @escaping Completion) throws -> Void,
startPlayback: @escaping () -> Void = {},
stopPlayback: @escaping () -> Void,
playbackTime: @escaping () -> Double?)
{
self.preparePlayback = preparePlayback
self.scheduleFrame = scheduleFrame
self.startPlayback = startPlayback
self.stopPlayback = stopPlayback
self.playbackTime = playbackTime
}
@ -93,37 +122,64 @@ public final class RealtimePCMStreamingAudioPlayer: PCMStreamingAudioPlaying {
stream: AsyncThrowingStream<Data, Error>,
sampleRate: Double) async -> StreamingPlaybackResult
{
_ = self.stop()
guard sampleRate > 0 else {
return StreamingPlaybackResult(finished: false, interruptedAt: nil)
}
self.generation &+= 1
let generation = self.generation
do {
try self.preparePlayback(sampleRate)
} catch {
return StreamingPlaybackResult(finished: false, interruptedAt: nil)
}
return await withCheckedContinuation { continuation in
self.playbackContinuation = continuation
self.inputTask = Task { @MainActor [weak self] in
await self?.consume(stream: stream, sampleRate: sampleRate, generation: generation)
await self.beginPlayback(stream: stream, sampleRate: sampleRate).value
}
func beginPlayback(
stream: AsyncThrowingStream<Data, Error>,
sampleRate: Double) -> Task<StreamingPlaybackResult, Never>
{
self.lock.withLock {
self.finish(StreamingPlaybackResult(finished: false, interruptedAt: nil), cancelInput: true)
let results = AsyncStream<StreamingPlaybackResult>.makeStream(bufferingPolicy: .bufferingOldest(1))
let resultTask = Task.detached(priority: .high) {
for await result in results.stream {
return result
}
return StreamingPlaybackResult(finished: false, interruptedAt: nil)
}
guard sampleRate > 0 else {
results.continuation.finish()
return resultTask
}
self.generation &+= 1
let generation = self.generation
self.playbackStarted = false
self.playbackContinuation = results.continuation
// Registration and queue submission are atomic with stop; engine work is not.
self.backendQueue.async { [weak self] in
guard let self, self.isCurrent(generation) else { return }
do {
try self.preparePlayback(sampleRate)
self.lock.withLock {
guard self.generation == generation else { return }
self.inputTask = Task.detached(priority: .high) { [weak self] in
await self?.consume(stream: stream, sampleRate: sampleRate, generation: generation)
}
}
} catch {
self.lock.withLock { self.finish(generation: generation, finished: false) }
}
}
return resultTask
}
}
public func stop() -> Double? {
// AVAudioPlayerNode's render-time query is thread-safe and does not wait for preparation.
let interruptedAt = self.playbackTime()
self.finish(StreamingPlaybackResult(
finished: false,
interruptedAt: interruptedAt), cancelInput: true)
self.lock.withLock {
self.finish(StreamingPlaybackResult(finished: false, interruptedAt: interruptedAt), cancelInput: true)
}
return interruptedAt
}
private func isCurrent(_ generation: UInt64) -> Bool {
self.lock.withLock { self.generation == generation }
}
private func consume(
stream: AsyncThrowingStream<Data, Error>,
sampleRate: Double,
generation: UInt64) async
stream: AsyncThrowingStream<Data, Error>, sampleRate: Double, generation: UInt64) async
{
let frameBytes = max(
MemoryLayout<Int16>.size,
@ -136,63 +192,95 @@ public final class RealtimePCMStreamingAudioPlayer: PCMStreamingAudioPlaying {
while pending.count >= frameBytes {
let frame = Data(pending.prefix(frameBytes))
pending.removeFirst(frameBytes)
guard await self.schedule(
frame: frame,
sampleRate: sampleRate,
generation: generation)
guard await self.schedule(frame: frame, sampleRate: sampleRate, generation: generation)
else { return }
}
}
if !pending.isEmpty {
pending.append(Data(repeating: 0, count: frameBytes - pending.count))
guard await self.schedule(
frame: pending,
sampleRate: sampleRate,
generation: generation)
guard await self.schedule(frame: pending, sampleRate: sampleRate, generation: generation)
else { return }
}
guard self.generation == generation else { return }
self.inputFinished = true
self.finishIfDrained(generation: generation)
self.backendQueue.async { [weak self] in
guard let self, self.isCurrent(generation) else { return }
self.startPlaybackOnce(generation: generation)
self.lock.withLock {
guard self.generation == generation else { return }
self.inputFinished = true
self.finishIfDrained(generation: generation)
}
}
} catch {
self.finish(generation: generation, finished: false)
let interruptedAt = self.playbackTime()
self.lock.withLock {
guard self.generation == generation else { return }
self.finish(StreamingPlaybackResult(finished: false, interruptedAt: interruptedAt))
}
}
}
private func schedule(frame: Data, sampleRate: Double, generation: UInt64) async -> Bool {
while self.generation == generation,
self.scheduledBufferIDs.count >= Self.maxScheduledBuffers
{
let admitted = await withCheckedContinuation { continuation in
self.slotWaiters.append(continuation)
let admitted = await withCheckedContinuation { continuation in
self.lock.withLock {
guard self.generation == generation, !Task.isCancelled else {
continuation.resume(returning: false)
return
}
if self.scheduledBufferIDs.count >= Self.maxScheduledBuffers {
self.slotWaiters.append(continuation)
} else { continuation.resume(returning: true) }
}
guard admitted else { return false }
}
guard self.generation == generation, !Task.isCancelled else { return false }
self.nextBufferID &+= 1
let bufferID = self.nextBufferID
self.scheduledBufferIDs.insert(bufferID)
do {
try self.scheduleFrame(frame, sampleRate) { [weak self] in
Task { @MainActor in
self?.completed(bufferID: bufferID, generation: generation)
guard admitted else { return false }
return await withCheckedContinuation { continuation in
self.backendQueue.async { [weak self] in
guard let self else { continuation.resume(returning: false)
return
}
let bufferID: UInt64? = self.lock.withLock {
guard self.generation == generation else { return nil }
self.nextBufferID &+= 1
self.scheduledBufferIDs.insert(self.nextBufferID)
return self.nextBufferID
}
guard let bufferID else { continuation.resume(returning: false)
return
}
do {
try self.scheduleFrame(frame, sampleRate) { [weak self] in
guard let self else { return }
self.backendQueue.async { [weak self] in
self?.lock.withLock { self?.completed(bufferID: bufferID, generation: generation) }
}
}
let shouldStart = self.lock.withLock {
self.generation == generation && self.scheduledBufferIDs.count >= Self.prebufferFrames
}
if shouldStart { self.startPlaybackOnce(generation: generation) }
continuation.resume(returning: self.isCurrent(generation))
} catch {
self.lock.withLock { self.finish(generation: generation, finished: false) }
continuation.resume(returning: false)
}
}
}
}
/// Called only on backendQueue. A concurrently retired generation is stopped next on that queue.
private func startPlaybackOnce(generation: UInt64) {
let start = self.lock.withLock {
guard self.generation == generation, !self.playbackStarted else { return false }
self.playbackStarted = true
return true
} catch {
self.scheduledBufferIDs.remove(bufferID)
self.finish(generation: generation, finished: false)
return false
}
if start {
self.startPlayback()
}
}
private func completed(bufferID: UInt64, generation: UInt64) {
guard self.generation == generation,
self.scheduledBufferIDs.remove(bufferID) != nil
else { return }
if !self.slotWaiters.isEmpty {
self.slotWaiters.removeFirst().resume(returning: true)
}
guard self.generation == generation, self.scheduledBufferIDs.remove(bufferID) != nil else { return }
if !self.slotWaiters.isEmpty { self.slotWaiters.removeFirst().resume(returning: true) }
self.finishIfDrained(generation: generation)
}
@ -203,11 +291,10 @@ public final class RealtimePCMStreamingAudioPlayer: PCMStreamingAudioPlaying {
private func finish(generation: UInt64, finished: Bool) {
guard self.generation == generation else { return }
self.finish(StreamingPlaybackResult(
finished: finished,
interruptedAt: finished ? nil : self.playbackTime()))
self.finish(StreamingPlaybackResult(finished: finished, interruptedAt: nil))
}
/// State-only retirement. Even synchronous backend stop callbacks cannot re-enter this lock.
private func finish(_ result: StreamingPlaybackResult, cancelInput: Bool = false) {
self.generation &+= 1
if cancelInput { self.inputTask?.cancel() }
@ -221,8 +308,20 @@ public final class RealtimePCMStreamingAudioPlayer: PCMStreamingAudioPlaying {
self.inputFinished = false
let continuation = self.playbackContinuation
self.playbackContinuation = nil
self.stopPlayback()
continuation?.resume(returning: result)
if continuation != nil {
self.backendQueue.async { [self] in self.stopPlayback() }
}
continuation?.yield(result)
continuation?.finish()
}
#if DEBUG
// periphery:ignore - tests observe completion processing, not just callback submission.
func _test_waitForBackendOperations() async {
await withCheckedContinuation { continuation in
self.backendQueue.async { continuation.resume() }
}
}
#endif
}
#endif

View file

@ -0,0 +1,463 @@
#if Talk && canImport(ElevenLabsKit) && (os(iOS) || os(macOS))
import Foundation
import OpenClawProtocol
import OSLog
/// All mutable output state and legacy player access are serialized by `lock`.
/// Audio events and synchronous relay controls use the same critical section.
final class RealtimeTalkOutput: @unchecked Sendable {
private let lock = NSLock()
private let player: (any RealtimePCMPlayback)?
private let legacyPlayer: PCMStreamingAudioPlaying
private let transport: RealtimeTalkRelayTransport
private let notification: AsyncStream<Void>.Continuation
private let logger = Logger(subsystem: "ai.openclawfoundation.app", category: "RealtimeTalkRelay")
private var effects: [Effect] = []
private var pendingStartupEvents = 0
var startupRoutingReady = false
private var routingGeneration: UInt64 = 0
var relaySessionId: String?
var isClosed = false
var outputSampleRateHz = 24000.0
var outputTask: Task<Void, Never>?
var outputContinuation: AsyncThrowingStream<Data, Error>.Continuation?
var pendingOutputAudio = Data()
var outputSessionId = 0
var pendingPlaybackMarks: [String] = []
var isOutputPaused = false
var reportedSpeaking = false
var isOutputPlaying = false
var outputIdentity: OutputIdentity?
var suppressedOutputIdentity: OutputIdentity?
var awaitingOutputClear = false
var cancelledOutputTurnId: String?
var terminalOutputCancellationReason: String?
var cancellationInFlight = false
var outputStartedAtMs: Double?
var outputAudioChunkCount = 0
var outputAudioByteCount = 0
var envelope = OutputEnvelope()
enum Effect: Sendable {
case speaking(Bool), beginLevels, cancelLevels, stopLegacyPlayer, cancellationCleared
case failure(String)
}
@MainActor
init(
player: PCMStreamingAudioPlaying,
transport: RealtimeTalkRelayTransport,
notification: AsyncStream<Void>.Continuation)
{
self.player = player as? any RealtimePCMPlayback
self.legacyPlayer = player
self.transport = transport
self.notification = notification
}
deinit {
self.outputContinuation?.finish()
self.outputTask?.cancel()
self.notification.finish()
}
func withLock<T>(_ body: (RealtimeTalkOutput) throws -> T) rethrows -> T {
try self.lock.withLock { try body(self) }
}
private func emit(_ effect: Effect) {
let notify = self.effects.isEmpty
self.effects.append(effect)
if notify {
self.notification.yield(())
}
}
func takeEffects() -> [Effect] {
let effects = self.effects
self.effects.removeAll(keepingCapacity: true)
return effects
}
/// Do not let a newly routed frame overtake events buffered before relay creation. Once
/// startup has drained, non-audio/UI events never put audio behind the main consumer.
func resetRouting(lifecycleGeneration: UInt64) {
self.routingGeneration = lifecycleGeneration
self.pendingStartupEvents = 0
self.startupRoutingReady = false
}
func route(_ event: EventFrame, lifecycleGeneration: UInt64) -> (handled: Bool, startup: Bool) {
self.withLock { output in
guard !output.isClosed, output.routingGeneration == lifecycleGeneration else { return (true, false) }
if !output.startupRoutingReady || output.pendingStartupEvents > 0 {
output.pendingStartupEvents += 1
return (false, true)
}
return (output.handleAudioEvent(event), false)
}
}
func mainEventHandled(startup: Bool, lifecycleGeneration: UInt64) {
if startup, self.routingGeneration == lifecycleGeneration {
self.pendingStartupEvents -= 1
}
}
/// Caller holds the output lock (route or the main startup consumer).
func handleAudioEvent(_ event: EventFrame) -> Bool {
guard !self.isClosed, let relaySessionId = self.relaySessionId,
event.event == "talk.event", let payload = event.payload?.dictionaryValue,
payload["relaySessionId"]?.stringValue == relaySessionId else { return false }
switch payload["type"]?.stringValue {
case "audio": self.handleOutputAudio(payload)
case "audioDone": self.handleOutputAudioDone(payload)
case "clear": self.handleOutputClear(payload)
case "mark": self.handlePlaybackMark(payload)
default: return false
}
return true
}
struct OutputIdentity: Equatable, Sendable {
let turnId: String?
init(_ payload: [String: AnyCodable]) {
self.turnId = payload["talkEvent"]?.dictionaryValue?["turnId"]?.stringValue?.trimmedNonEmpty
}
}
func retireCancellation() {
self.suppressedOutputIdentity = nil
self.awaitingOutputClear = false
self.cancellationInFlight = false
}
func handleOutputClear(_ payload: [String: AnyCodable]) {
let clearIdentity = OutputIdentity(payload)
// Provider clears retire playback; only turn.cancelled acknowledges turn cancellation.
let clearsSuppressed = self.awaitingOutputClear &&
payload["talkEvent"]?.dictionaryValue?["type"]?.stringValue == "turn.cancelled" &&
self.suppressedOutputIdentity == clearIdentity
if clearsSuppressed {
self.awaitingOutputClear = false
if !self.cancellationInFlight {
self.suppressedOutputIdentity = nil
self.emit(.cancellationCleared)
}
}
let currentMatches = clearIdentity.turnId == nil || self.outputIdentity == clearIdentity
guard clearsSuppressed || currentMatches else { return }
let marks = self.takePendingPlaybackMarks()
// Cancellation already published the stopped state. A later clear with no
// active output only retires the fence; it must not emit a duplicate callback.
if self.isOutputPlaying || self.outputIdentity != nil {
self.stopOutputPlayback()
}
self.acknowledgePlaybackMarks(marks)
}
func recordOutputAudioChunk(byteCount: Int) {
self.outputAudioChunkCount += 1
self.outputAudioByteCount += byteCount
guard self.outputAudioChunkCount == 1 || self.outputAudioChunkCount % 20 == 0 else { return }
let chunks = self.outputAudioChunkCount
let bytes = self.outputAudioByteCount
self.logger.debug("talk realtime audio: chunks=\(chunks) bytes=\(bytes)")
}
func markOutputAudioStarted(nowMs: Double) {
if !self.isOutputPlaying {
self.outputStartedAtMs = nowMs
}
self.isOutputPlaying = true
}
func finishOutputPlaybackStream() {
guard let continuation = outputContinuation else { return }
if !self.pendingOutputAudio.isEmpty {
let trailingFrame = self.pendingOutputAudio
self.pendingOutputAudio.removeAll(keepingCapacity: true)
guard self.yieldOutputAudioFrame(trailingFrame) else { return }
}
continuation.finish()
self.outputContinuation = nil
}
func markOutputPlaybackFinished() {
// Only drained playback completes output; elapsed time cannot prove the
// device finished queued audio. Publish the terminal transition once.
guard self.isOutputPlaying else { return }
self.isOutputPlaying = false
self.outputIdentity = nil
self.outputStartedAtMs = nil
self.envelope.cancel()
self.emit(.cancelLevels)
self.reportSpeaking(false)
self.acknowledgePlaybackMarks(self.takePendingPlaybackMarks())
}
func takePendingPlaybackMarks() -> [String] {
let marks = self.pendingPlaybackMarks
self.pendingPlaybackMarks.removeAll()
return marks
}
func handlePlaybackMark(_ payload: [String: AnyCodable]) {
guard let markName = payload["markName"]?.stringValue?.trimmedNonEmpty else { return }
if self.isOutputPlaying {
self.pendingPlaybackMarks.append(markName)
} else {
self.acknowledgePlaybackMarks([markName])
}
}
func acknowledgePlaybackMarks(_ marks: [String]) {
guard !marks.isEmpty,
let relaySessionId
else { return }
for markName in marks {
Task { [transport, logger] in
let payload: [String: AnyCodable] = [
"sessionId": AnyCodable(relaySessionId),
"markName": AnyCodable(markName),
]
do {
_ = try await transport.request("talk.session.acknowledgeMark", payload, 8000)
} catch {
let message = String(error.localizedDescription.prefix(180))
logger.warning(
"talk realtime: mark acknowledgement failed=\(message, privacy: .public)")
}
}
}
}
func stopOutputPlayback() {
self.outputSessionId += 1
self.outputContinuation?.finish()
self.outputContinuation = nil
self.pendingOutputAudio.removeAll(keepingCapacity: true)
self.outputTask?.cancel()
self.outputTask = nil
if let player {
_ = player.stop()
} else {
self.emit(.stopLegacyPlayer)
}
self.isOutputPlaying = false
self.outputIdentity = nil
self.outputStartedAtMs = nil
self.envelope.cancel()
self.emit(.cancelLevels)
self.reportSpeaking(false)
}
func handleOutputAudio(_ payload: [String: AnyCodable]) {
guard !self.isOutputPaused else { return }
let incomingIdentity = OutputIdentity(payload)
guard let incomingTurnId = incomingIdentity.turnId else {
self.handleOutputPlaybackOverflow()
return
}
guard !self.awaitingOutputClear else { return }
guard incomingTurnId != self.cancelledOutputTurnId else { return }
guard let base64 = payload["audioBase64"]?.stringValue else { return }
guard let data = Data(base64Encoded: base64) else {
self.handleOutputPlaybackOverflow()
return
}
self.terminalOutputCancellationReason = nil
if let currentIdentity = outputIdentity,
currentIdentity != incomingIdentity
{
let marks = self.takePendingPlaybackMarks()
self.stopOutputPlayback()
self.acknowledgePlaybackMarks(marks)
} else if self.outputContinuation == nil, self.outputTask != nil {
self.stopOutputPlayback()
}
self.outputIdentity = incomingIdentity
self.recordOutputAudioChunk(byteCount: data.count)
self.markOutputAudioStarted(nowMs: ProcessInfo.processInfo.systemUptime * 1000)
self.reportSpeaking(true)
self.ensureOutputPlaybackStarted()
self.bufferOutputAudio(data)
}
func reportSpeaking(_ speaking: Bool) {
// Only the per-chunk `true` floods; `false` comes from idempotent teardown paths.
guard !(speaking && self.reportedSpeaking) else { return }
self.reportedSpeaking = speaking
self.emit(.speaking(speaking))
}
func bufferOutputAudio(_ data: Data) {
let frameByteCount = max(2, Int((outputSampleRateHz * 0.02).rounded()) * 2)
var offset = data.startIndex
if !self.pendingOutputAudio.isEmpty {
let fillCount = min(frameByteCount - self.pendingOutputAudio.count, data.count)
let fillEnd = data.index(offset, offsetBy: fillCount)
self.pendingOutputAudio.append(data[offset..<fillEnd])
offset = fillEnd
if self.pendingOutputAudio.count == frameByteCount {
let frame = self.pendingOutputAudio
self.pendingOutputAudio.removeAll(keepingCapacity: true)
guard self.yieldOutputAudioFrame(frame) else { return }
}
}
while data.distance(from: offset, to: data.endIndex) >= frameByteCount {
let frameEnd = data.index(offset, offsetBy: frameByteCount)
let frame = Data(data[offset..<frameEnd])
offset = frameEnd
guard self.yieldOutputAudioFrame(frame) else { return }
}
if offset < data.endIndex {
self.pendingOutputAudio.append(data[offset...])
}
}
func yieldOutputAudioFrame(_ data: Data) -> Bool {
guard let continuation = outputContinuation else { return false }
switch continuation.yield(data) {
case .enqueued:
if self.envelope.append(data) { self.emit(.beginLevels) }
return true
case .dropped:
self.handleOutputPlaybackOverflow()
return false
case .terminated:
return false
@unknown default:
self.handleOutputPlaybackOverflow()
return false
}
}
func handleOutputAudioDone(_ payload: [String: AnyCodable]) {
let incomingIdentity = OutputIdentity(payload)
if incomingIdentity.turnId != nil,
let outputIdentity,
outputIdentity != incomingIdentity
{
return
}
self.finishOutputPlaybackStream()
}
private func ensureOutputPlaybackStarted() {
guard self.outputContinuation == nil, self.outputTask == nil else { return }
self.outputSessionId += 1
let sessionId = self.outputSessionId
let sampleRate = self.outputSampleRateHz
self.envelope.begin(sampleRate: sampleRate)
self.emit(.beginLevels)
let (stream, continuation) = AsyncThrowingStream<Data, Error>.makeStream(
bufferingPolicy: .bufferingOldest(RealtimeTalkRelaySession.maxBufferedOutputChunks))
self.outputContinuation = continuation
// The real backend registers its generation under the output lock, so a stop
// cannot race a delayed task that starts an already-retired reply. Actor-bound legacy
// players remain supported for existing clients/fakes, but are not the realtime backend.
let playback: Task<StreamingPlaybackResult, Never> = if let player {
player.beginPlayback(stream: stream, sampleRate: sampleRate)
} else {
Task { @MainActor [legacyPlayer] in
guard !Task.isCancelled else {
return StreamingPlaybackResult(finished: false, interruptedAt: nil)
}
return await legacyPlayer.play(stream: stream, sampleRate: sampleRate)
}
}
self.outputTask = Task.detached(priority: .high) { [weak self] in
let result = await withTaskCancellationHandler {
await playback.value
} onCancel: { playback.cancel() }
self?.withLock { output in
guard output.outputSessionId == sessionId, !output.isClosed else { return }
output.outputTask = nil
output.outputContinuation = nil
if result.finished {
output.markOutputPlaybackFinished()
} else {
output.handleOutputPlaybackFailure(
String(localized: "Realtime audio playback failed. Reconnecting…"))
}
}
}
}
private func handleOutputPlaybackOverflow() {
self.handleOutputPlaybackFailure(
String(localized: "Realtime audio playback fell behind. Reconnecting…"))
}
private func handleOutputPlaybackFailure(_ message: String) {
guard !self.isClosed else { return }
// Fence immediately; main may be busy, but no more audio can enter after failure.
self.isClosed = true
self.stopOutputPlayback()
self.emit(.failure(message))
}
/// Same PCM-time envelope as PCMPlaybackEnvelope, sampled by the UI at 30 Hz. Metering
/// and its timeline are output-lock-owned; publishing levels never gates scheduling.
struct OutputEnvelope {
private struct Segment { let start: Double
let end: Double
let level: Double
}
private var segments: [Segment] = []
private var bytesPerSecond = 0.0
private var startedAt: Double?
private var scheduleEnd = 0.0
mutating func begin(sampleRate: Double) {
self.cancel()
self.bytesPerSecond = max(1, sampleRate * 2)
}
mutating func append(_ data: Data) -> Bool {
guard self.bytesPerSecond > 1, !data.isEmpty else { return false }
let restarting = self.startedAt == nil
let now = ProcessInfo.processInfo.systemUptime
if self.startedAt == nil {
self.startedAt = now
}
var start = max(now - (startedAt ?? now), self.scheduleEnd)
let windowBytes = max(2, Int(bytesPerSecond * 0.05) & ~1)
var offset = data.startIndex
while offset < data.endIndex {
let end = min(offset + windowBytes, data.endIndex)
let window = Data(data[offset..<end])
let duration = Double(window.count) / self.bytesPerSecond
self.segments.append(Segment(
start: start,
end: start + duration,
level: TalkAudioLevel.normalized(rms: TalkAudioLevel.pcm16RMS(window))))
start += duration
offset = end
}
self.scheduleEnd = start
return restarting
}
mutating func level() -> Double? {
guard let startedAt else { return 0 }
let elapsed = ProcessInfo.processInfo.systemUptime - startedAt
guard elapsed <= self.scheduleEnd + 0.5 else {
self.cancel()
return nil
}
self.segments.removeAll { $0.end < elapsed }
return self.segments.first { elapsed >= $0.start && elapsed < $0.end }?.level ?? 0
}
mutating func cancel() {
self.segments.removeAll(keepingCapacity: true)
self.startedAt = nil
self.scheduleEnd = 0
}
}
}
#endif

View file

@ -145,7 +145,6 @@ private actor RealtimeAudioSender {
private let request: @Sendable (String, [String: AnyCodable]?, Double) async throws -> Data
private var relaySessionId: String?
private var pendingSends = 0
private let maxPendingSends = 4
init(
relaySessionId: String,
@ -161,7 +160,7 @@ private actor RealtimeAudioSender {
func send(_ data: Data, timestampMs: Double) async -> RealtimeAudioSendOutcome {
guard !Task.isCancelled, let relaySessionId else { return .inactive }
guard self.pendingSends < self.maxPendingSends else { return .saturated }
guard self.pendingSends < RealtimeTalkRelaySession.maxPendingAudioSends else { return .saturated }
self.pendingSends += 1
defer { self.pendingSends -= 1 }
// The Gateway carries this straight into the provider's media timeline, and OpenAI rejects
@ -251,14 +250,22 @@ public final class RealtimeTalkRelaySession {
private nonisolated static let bargeInCooldownMs: Double = 900
private nonisolated static let minOutputBeforeBargeInMs: Double = 250
private nonisolated static let startupReadyTimeoutSeconds = 12
/// At the protocol's 20 ms cadence this bounds queued relay audio to 640 ms / 30,720 bytes.
/// Overflow terminates the session so recovery replaces a lagging playback path.
private nonisolated static let maxBufferedOutputChunks = 32
/// Providers may deliver a whole reply faster than realtime (xAI sends it in one burst), so the
/// bound must hold a full reply: 60 s of 20 ms frames (~2.9 MB at 24 kHz). Overflow still
/// terminates the session so recovery replaces a stalled playback path.
/// In-flight `talk.session.appendAudio` requests before input counts as stalled. Each mic
/// callback (~43 ms) is one request; Gateway round trips over Wi-Fi/Tailscale spike to ~1 s, so
/// 48 tolerates ~2 s of latency instead of the ~170 ms that 4 allowed.
nonisolated static let maxPendingAudioSends = 48
nonisolated static let maxBufferedOutputChunks = 3000
private let transport: RealtimeTalkRelayTransport
private let audioCapture: any RealtimeTalkAudioCapturing
private let options: Options
private let pcmPlayer: PCMStreamingAudioPlaying
private let output: RealtimeTalkOutput
private var outputEffectsTask: Task<Void, Never>?
private var outputLevelTask: Task<Void, Never>?
private let logger = Logger(subsystem: "ai.openclawfoundation.app", category: "RealtimeTalkRelay")
private let onStatus: (String) -> Void
private let onIssue: (RealtimeTalkRelayIssue) -> Void
@ -267,9 +274,6 @@ public final class RealtimeTalkRelaySession {
private let onInputLevel: (Double) -> Void
private let onOutputLevel: (Double?) -> Void
private let onTranscript: (RealtimeTalkTranscript) -> Void
/// Playback-time-aligned envelope of the assistant PCM the relay schedules;
/// drives the speaking waveform with real audio instead of a synthetic pulse.
private var outputEnvelope: PCMPlaybackEnvelope?
private var relaySessionId: String?
private var serverClose: (sessionId: String, task: Task<Void, Error>)?
@ -279,33 +283,17 @@ public final class RealtimeTalkRelaySession {
private var startupWaiter: CheckedContinuation<StartupWaitResult, Never>?
private var pendingPreRelayEvents: [EventFrame] = []
private var inputSampleRateHz = Double(RealtimeTalkRelaySession.defaultSampleRateHz)
private var outputSampleRateHz = Double(RealtimeTalkRelaySession.defaultSampleRateHz)
private var supportsBargeIn: Bool? = true
private var eventTask: Task<Void, Never>?
private var toolCallTasks: [UUID: Task<Void, Never>] = [:]
private var audioSendTasks: [UUID: Task<Void, Never>] = [:]
private var outputTask: Task<Void, Never>?
private var outputContinuation: AsyncThrowingStream<Data, Error>.Continuation?
/// Provider deltas may span any number of frames; retain only the partial tail so the
/// AsyncStream's 32 slots always contain bounded 20 ms PCM chunks.
private var pendingOutputAudio = Data()
private var outputSessionId = 0
private var pendingPlaybackMarks: [String] = []
private var audioSender: RealtimeAudioSender?
private var isInputPaused = false
private var isOutputPaused = false
private var audioCaptureGeneration: UInt64 = 0
private var isClosed = false
private var lifecycleGeneration: UInt64 = 0
private var outputCancellationGeneration: UInt64 = 0
private var isOutputPlaying = false
private var outputIdentity: OutputIdentity?
private var suppressedOutputIdentity: OutputIdentity?
private var awaitingOutputClear = false
private var cancelledOutputTurnId: String?
private var terminalOutputCancellationReason: String?
private var outputCancellationTask: Task<Void, Never>?
private var outputStartedAtMs: Double?
private var lastBargeInAtMs: Double = 0
private var micLogFrameCount = 0
private var micLogByteCount = 0
@ -315,8 +303,6 @@ public final class RealtimeTalkRelaySession {
private var suppressedEchoByteCount = 0
private var suppressedEchoMaxRms: Float = 0
private var lastSuppressedEchoLogAtMs: Double = 0
private var outputAudioChunkCount = 0
private var outputAudioByteCount = 0
public var voiceSessionId: String? {
self.relaySessionId
@ -343,6 +329,11 @@ public final class RealtimeTalkRelaySession {
self.audioCapture = audioCapture
self.options = options
self.pcmPlayer = pcmPlayer
let notifications = AsyncStream<Void>.makeStream()
self.output = RealtimeTalkOutput(
player: pcmPlayer,
transport: transport,
notification: notifications.continuation)
self.onStatus = onStatus
self.onIssue = onIssue
self.onTermination = onTermination
@ -350,12 +341,29 @@ public final class RealtimeTalkRelaySession {
self.onInputLevel = onInputLevel
self.onOutputLevel = onOutputLevel
self.onTranscript = onTranscript
self.outputEffectsTask = Task { @MainActor [weak self] in
for await _ in notifications.stream {
guard let self else { return }
self.drainOutputEffects()
}
}
}
deinit {
self.eventTask?.cancel()
self.outputEffectsTask?.cancel()
self.outputLevelTask?.cancel()
}
public func start() async throws {
self.lifecycleGeneration &+= 1
let lifecycleGeneration = self.lifecycleGeneration
self.isClosed = false
self.output.withLock {
$0.isClosed = false
$0.relaySessionId = nil
$0.resetRouting(lifecycleGeneration: lifecycleGeneration)
}
self.hasReceivedReady = false
self.hasReceivedFailure = false
self.supportsBargeIn = nil
@ -363,17 +371,17 @@ public final class RealtimeTalkRelaySession {
self.startupWaiter = nil
self.pendingPreRelayEvents.removeAll()
self.onStatus("Connecting realtime…")
let eventStream = await self.transport.subscribeServerEvents(200)
switch await self.lifecycleStatus(lifecycleGeneration) {
let eventStream = await transport.subscribeServerEvents(200)
switch await lifecycleStatus(lifecycleGeneration) {
case .current: break
case .cancelledLocally: return
case .routeLost: throw Self.gatewayRouteLostError()
}
self.startEventPump(stream: eventStream, lifecycleGeneration: lifecycleGeneration)
do {
let result = try await self.createRelaySession()
let result = try await createRelaySession()
let createdRelaySessionId = result.relaysessionid?.trimmedNonEmpty
let statusAfterCreate = await self.lifecycleStatus(lifecycleGeneration)
let statusAfterCreate = await lifecycleStatus(lifecycleGeneration)
if statusAfterCreate != .current {
if let relaySessionId = createdRelaySessionId {
try? await self.beginServerClose(relaySessionId: relaySessionId).value
@ -396,8 +404,10 @@ public final class RealtimeTalkRelaySession {
])
}
self.relaySessionId = relaySessionId
let supportsBargeIn = await self.resolveSupportsBargeIn(result)
switch await self.lifecycleStatus(lifecycleGeneration) {
// Acknowledgments need identity during startup; routing stays gated until replay finishes.
self.output.withLock { $0.relaySessionId = relaySessionId }
let supportsBargeIn = await resolveSupportsBargeIn(result)
switch await lifecycleStatus(lifecycleGeneration) {
case .current: break
case .cancelledLocally: return
case .routeLost: throw Self.gatewayRouteLostError()
@ -407,15 +417,16 @@ public final class RealtimeTalkRelaySession {
relaySessionId: relaySessionId,
request: self.transport.request)
self.configureAudioContract(result.audio)
try self.startMicrophonePump(lifecycleGeneration: lifecycleGeneration)
try startMicrophonePump(lifecycleGeneration: lifecycleGeneration)
self.onStatus("Waiting for realtime…")
await self.drainPendingPreRelayEvents(lifecycleGeneration: lifecycleGeneration)
switch await self.lifecycleStatus(lifecycleGeneration) {
await drainPendingPreRelayEvents(lifecycleGeneration: lifecycleGeneration)
self.output.withLock { $0.startupRoutingReady = true }
switch await lifecycleStatus(lifecycleGeneration) {
case .current: break
case .cancelledLocally: return
case .routeLost: throw Self.gatewayRouteLostError()
}
switch await self.waitForStartupResult(
switch await waitForStartupResult(
timeoutSeconds: Self.startupReadyTimeoutSeconds,
lifecycleGeneration: lifecycleGeneration)
{
@ -427,7 +438,9 @@ public final class RealtimeTalkRelaySession {
} catch {
// A lost route must still surface: swallowing here would discard both the original
// failure and the route loss, leaving the runtime with nothing to fall back from.
if await self.lifecycleStatus(lifecycleGeneration) == .cancelledLocally { return }
if await lifecycleStatus(lifecycleGeneration) == .cancelledLocally {
return
}
let createdRelaySessionId = self.relaySessionId
self.close(sendClose: false)
if let createdRelaySessionId {
@ -450,28 +463,37 @@ public final class RealtimeTalkRelaySession {
private func close(sendClose: Bool) {
guard !self.isClosed else { return }
self.isClosed = true
self.outputCancellationGeneration &+= 1
self.outputCancellationTask?.cancel()
self.outputCancellationTask = nil
self.output.withLock { output in
let wasClosed = output.isClosed
output.isClosed = true
output.relaySessionId = nil
output.pendingPlaybackMarks.removeAll()
output.cancelledOutputTurnId = nil
output.terminalOutputCancellationReason = nil
output.isOutputPaused = false
output.retireCancellation()
if !wasClosed { output.stopOutputPlayback() }
output.reportSpeaking(false)
}
self.lifecycleGeneration &+= 1
self.finishStartupWait(.cancelled)
self.stopMicrophonePump()
finishStartupWait(.cancelled)
stopMicrophonePump()
self.eventTask?.cancel()
self.eventTask = nil
for task in self.toolCallTasks.values {
task.cancel()
}
self.pendingPlaybackMarks.removeAll()
let audioSender = self.audioSender
self.audioSender = nil
Task { await audioSender?.close() }
self.retireOutputCancellation()
self.cancelledOutputTurnId = nil
self.terminalOutputCancellationReason = nil
self.isOutputPaused = false
self.stopOutputPlayback()
if sendClose, let relaySessionId = self.relaySessionId {
self.drainOutputEffects()
if sendClose, let relaySessionId {
self.beginServerClose(relaySessionId: relaySessionId)
}
self.relaySessionId = nil
self.onSpeakingChanged(false)
relaySessionId = nil
}
/// Deliberately not a `CancellationError`: the runtime treats those as caller-initiated and
@ -521,10 +543,13 @@ public final class RealtimeTalkRelaySession {
}
public func setOutputPaused(_ paused: Bool) {
guard self.isOutputPaused != paused else { return }
self.isOutputPaused = paused
if paused, self.isOutputPlaying {
self.cancelOutput(reason: "pause")
let cancel = self.output.withLock { output in
guard output.isOutputPaused != paused else { return false }
output.isOutputPaused = paused
return paused && output.isOutputPlaying
}
if cancel {
cancelOutput(reason: "pause")
}
}
@ -564,8 +589,9 @@ public final class RealtimeTalkRelaySession {
}
self.inputSampleRateHz = audio["inputSampleRateHz"]?.doubleValue
?? Double(Self.defaultSampleRateHz)
self.outputSampleRateHz = audio["outputSampleRateHz"]?.doubleValue
?? Double(Self.defaultSampleRateHz)
self.output.withLock {
$0.outputSampleRateHz = audio["outputSampleRateHz"]?.doubleValue ?? Double(Self.defaultSampleRateHz)
}
}
private func resolveSupportsBargeIn(_ session: TalkSessionCreateResult) async -> Bool? {
@ -591,16 +617,68 @@ public final class RealtimeTalkRelaySession {
return entry?["supportsBargeIn"]?.boolValue ?? true
}
private func drainOutputEffects() {
for effect in self.output.withLock({ $0.takeEffects() }) {
switch effect {
case let .speaking(speaking): self.onSpeakingChanged(speaking)
case .stopLegacyPlayer: _ = self.pcmPlayer.stop()
case .beginLevels:
self.outputLevelTask?.cancel()
self.outputLevelTask = Task { @MainActor [weak self] in
while !Task.isCancelled {
guard let self else { return }
let level = self.output.withLock { $0.envelope.level() }
guard let level else {
self.onOutputLevel(nil)
return
}
self.onOutputLevel(level)
try? await Task.sleep(for: .milliseconds(33))
}
}
case .cancelLevels:
self.outputLevelTask?.cancel()
self.outputLevelTask = nil
self.onOutputLevel(nil)
case .cancellationCleared:
if self.outputCancellationTask == nil, !self.output.withLock({ $0.awaitingOutputClear }) {
self.retireOutputCancellation()
}
case let .failure(message): handleOutputPlaybackFailure(message)
}
}
}
private func startEventPump(stream: AsyncStream<EventFrame>, lifecycleGeneration: UInt64) {
self.eventTask?.cancel()
self.eventTask = Task { [weak self] in
for await event in stream {
if Task.isCancelled { return }
await self?.handleGatewayEvent(event, lifecycleGeneration: lifecycleGeneration)
let mainEvents = AsyncStream<(event: EventFrame, startup: Bool)>.makeStream()
let consumer = Task { @MainActor [weak self] in
for await delivery in mainEvents.stream {
guard !Task.isCancelled else { return }
await self?.handleGatewayEvent(delivery.event, lifecycleGeneration: lifecycleGeneration)
self?.output.withLock {
$0.mainEventHandled(startup: delivery.startup, lifecycleGeneration: lifecycleGeneration)
}
}
guard !Task.isCancelled else { return }
await self?.handleEventStreamEnded(lifecycleGeneration: lifecycleGeneration)
}
self.eventTask = Task.detached(priority: .high) { [output = self.output] in
await withTaskCancellationHandler {
for await event in stream {
guard !Task.isCancelled else { break }
let route = output.route(event, lifecycleGeneration: lifecycleGeneration)
if !route.handled {
mainEvents.continuation.yield((event, route.startup))
}
}
mainEvents.continuation.finish()
await consumer.value
} onCancel: {
mainEvents.continuation.finish()
consumer.cancel()
}
}
}
}
@ -647,13 +725,17 @@ extension RealtimeTalkRelaySession {
self.finishStartupWait(.ready)
self.onStatus("Listening (Realtime)")
case "audio":
self.handleOutputAudio(payload)
self.output.withLock { $0.handleOutputAudio(payload) }
self.drainOutputEffects()
case "audioDone":
self.handleOutputAudioDone(payload)
self.output.withLock { $0.handleOutputAudioDone(payload) }
self.drainOutputEffects()
case "clear":
self.handleOutputClear(payload)
self.output.withLock { $0.handleOutputClear(payload) }
self.drainOutputEffects()
case "mark":
self.handlePlaybackMark(payload)
self.output.withLock { $0.handlePlaybackMark(payload) }
self.drainOutputEffects()
case "transcript":
self.handleTranscriptEvent(payload)
case "toolCall":
@ -676,8 +758,10 @@ extension RealtimeTalkRelaySession {
let termination: RealtimeTalkRelayTermination = if talkEvent?["type"]?.stringValue == "session.closed",
talkEvent?["payload"]?.dictionaryValue?["reason"]?
.stringValue == "output-cancelled",
let cancellationReason = self
.terminalOutputCancellationReason
let cancellationReason = output
.withLock({
$0.terminalOutputCancellationReason
})
{
.outputCancelled(reason: cancellationReason)
} else {
@ -700,27 +784,6 @@ extension RealtimeTalkRelaySession {
}
}
private func handleOutputClear(_ payload: [String: AnyCodable]) {
let clearIdentity = OutputIdentity(payload)
// Provider clears retire playback; only turn.cancelled acknowledges turn cancellation.
let clearsSuppressed = self.awaitingOutputClear &&
payload["talkEvent"]?.dictionaryValue?["type"]?.stringValue == "turn.cancelled" &&
self.suppressedOutputIdentity == clearIdentity
if clearsSuppressed {
self.awaitingOutputClear = false
if self.outputCancellationTask == nil { self.retireOutputCancellation() }
}
let currentMatches = clearIdentity.turnId == nil || self.outputIdentity == clearIdentity
guard clearsSuppressed || currentMatches else { return }
let marks = self.takePendingPlaybackMarks()
// Cancellation already published the stopped state. A later clear with no
// active output only retires the fence; it must not emit a duplicate callback.
if self.isOutputPlaying || self.outputIdentity != nil {
self.stopOutputPlayback()
}
self.acknowledgePlaybackMarks(marks)
}
private func waitForStartupResult(
timeoutSeconds: Int,
lifecycleGeneration: UInt64) async -> StartupWaitResult
@ -786,30 +849,14 @@ extension RealtimeTalkRelaySession {
phase: payload["phase"]?.stringValue ?? phase)
}
private func recordOutputAudioChunk(byteCount: Int) {
self.outputAudioChunkCount += 1
self.outputAudioByteCount += byteCount
guard self.outputAudioChunkCount == 1 || self.outputAudioChunkCount % 20 == 0 else { return }
self.logger.debug(
"talk realtime audio: chunks=\(self.outputAudioChunkCount) bytes=\(self.outputAudioByteCount)")
}
private func markOutputAudioStarted(nowMs: Double) {
if !self.isOutputPlaying {
self.outputStartedAtMs = nowMs
}
self.isOutputPlaying = true
}
private func handleInputLevelDuringOutput(_ rms: Float, timestampMs: Double) {
guard self.isOutputPlaying else { return }
guard rms >= Self.bargeInRmsThreshold else { return }
if let outputStartedAtMs,
timestampMs - outputStartedAtMs < Self.minOutputBeforeBargeInMs
{
return
let shouldCancel = self.output.withLock { output in
guard output.isOutputPlaying, rms >= Self.bargeInRmsThreshold else { return false }
if let outputStartedAtMs = output.outputStartedAtMs,
timestampMs - outputStartedAtMs < Self.minOutputBeforeBargeInMs { return false }
return timestampMs - self.lastBargeInAtMs >= Self.bargeInCooldownMs
}
guard timestampMs - self.lastBargeInAtMs >= Self.bargeInCooldownMs else { return }
guard shouldCancel else { return }
self.lastBargeInAtMs = timestampMs
self.cancelOutput(reason: "barge-in")
}
@ -1018,113 +1065,6 @@ extension RealtimeTalkRelaySession {
guard await self.isCurrentLifecycle(lifecycleGeneration) else { throw CancellationError() }
}
private func ensureOutputPlaybackStarted() {
guard self.outputContinuation == nil, self.outputTask == nil else { return }
self.outputSessionId += 1
let sessionId = self.outputSessionId
let envelope = self.outputEnvelope ?? PCMPlaybackEnvelope { [weak self] level in
self?.onOutputLevel(level)
}
envelope.begin(sampleRate: self.outputSampleRateHz)
self.outputEnvelope = envelope
let stream = AsyncThrowingStream<Data, Error>(
bufferingPolicy: .bufferingOldest(Self.maxBufferedOutputChunks))
{ continuation in self.outputContinuation = continuation }
self.outputTask = Task { [weak self] in
guard let self else { return }
guard self.outputSessionId == sessionId, !self.isClosed, !Task.isCancelled else { return }
let result = await self.pcmPlayer.play(stream: stream, sampleRate: self.outputSampleRateHz)
await MainActor.run {
guard self.outputSessionId == sessionId else { return }
self.outputTask = nil
self.outputContinuation = nil
if !result.finished {
if let interruptedAt = result.interruptedAt {
self.logger.info("realtime output interrupted at \(interruptedAt, privacy: .public)s")
}
self.handleOutputPlaybackFailure(
String(localized: "Realtime audio playback failed. Reconnecting…"))
return
}
self.markOutputPlaybackFinished()
}
}
}
private func finishOutputPlaybackStream() {
guard let continuation = self.outputContinuation else { return }
if !self.pendingOutputAudio.isEmpty {
let trailingFrame = self.pendingOutputAudio
self.pendingOutputAudio.removeAll(keepingCapacity: true)
guard self.yieldOutputAudioFrame(trailingFrame) else { return }
}
continuation.finish()
self.outputContinuation = nil
}
private func markOutputPlaybackFinished() {
// Only drained playback completes output; elapsed time cannot prove the
// device finished queued audio. Publish the terminal transition once.
guard self.isOutputPlaying else { return }
self.isOutputPlaying = false
self.outputIdentity = nil
self.outputStartedAtMs = nil
self.outputEnvelope?.cancel()
self.onSpeakingChanged(false)
self.acknowledgePlaybackMarks(self.takePendingPlaybackMarks())
}
private func takePendingPlaybackMarks() -> [String] {
let marks = self.pendingPlaybackMarks
self.pendingPlaybackMarks.removeAll()
return marks
}
private func handlePlaybackMark(_ payload: [String: AnyCodable]) {
guard let markName = payload["markName"]?.stringValue?.trimmedNonEmpty else { return }
if self.isOutputPlaying {
self.pendingPlaybackMarks.append(markName)
} else {
self.acknowledgePlaybackMarks([markName])
}
}
private func acknowledgePlaybackMarks(_ marks: [String]) {
guard !marks.isEmpty,
let relaySessionId = self.relaySessionId
else { return }
for markName in marks {
Task { [transport, logger] in
let payload: [String: AnyCodable] = [
"sessionId": AnyCodable(relaySessionId),
"markName": AnyCodable(markName),
]
do {
_ = try await transport.request("talk.session.acknowledgeMark", payload, 8000)
} catch {
let message = Self.safeLogMessage(error.localizedDescription)
logger.warning(
"talk realtime: mark acknowledgement failed=\(message, privacy: .public)")
}
}
}
}
private func stopOutputPlayback() {
self.outputSessionId += 1
self.outputContinuation?.finish()
self.outputContinuation = nil
self.pendingOutputAudio.removeAll(keepingCapacity: true)
self.outputTask?.cancel()
self.outputTask = nil
_ = self.pcmPlayer.stop()
self.isOutputPlaying = false
self.outputIdentity = nil
self.outputStartedAtMs = nil
self.outputEnvelope?.cancel()
self.onSpeakingChanged(false)
}
private nonisolated static func safeLogMessage(_ value: String) -> String {
let singleLine = value
.replacingOccurrences(of: "\n", with: " ")
@ -1157,29 +1097,25 @@ extension RealtimeTalkRelaySession {
}
extension RealtimeTalkRelaySession {
private struct OutputIdentity: Equatable {
let turnId: String?
init(_ payload: [String: AnyCodable]) {
self.turnId = payload["talkEvent"]?.dictionaryValue?["turnId"]?.stringValue?.trimmedNonEmpty
}
}
@discardableResult
public func cancelOutput(reason: String = "user") -> Bool {
guard reason != "barge-in" || self.supportsBargeIn == true else { return false }
guard let relaySessionId,
let outputIdentity = self.outputIdentity,
let turnId = outputIdentity.turnId
else { return false }
guard let relaySessionId else { return false }
let turnId = self.output.withLock { output -> String? in
guard let identity = output.outputIdentity, let turnId = identity.turnId else { return nil }
output.terminalOutputCancellationReason = reason == "barge-in" ? nil : reason
output.suppressedOutputIdentity = identity
output.cancelledOutputTurnId = turnId
output.awaitingOutputClear = true
output.cancellationInFlight = true
output.stopOutputPlayback()
return turnId
}
guard let turnId else { return false }
self.drainOutputEffects()
self.outputCancellationGeneration &+= 1
let cancellationGeneration = self.outputCancellationGeneration
self.outputCancellationTask?.cancel()
self.terminalOutputCancellationReason = reason == "barge-in" ? nil : reason
self.suppressedOutputIdentity = outputIdentity
self.cancelledOutputTurnId = outputIdentity.turnId
self.awaitingOutputClear = true
self.stopOutputPlayback()
self.outputCancellationTask = Task { [weak self, transport] in
let payload: [String: AnyCodable] = [
"sessionId": AnyCodable(relaySessionId),
@ -1193,17 +1129,26 @@ extension RealtimeTalkRelaySession {
guard let self, self.isCurrentOutputCancellation(cancellationGeneration) else { return }
switch result.status?.stringValue {
case "stale", "idle":
self.terminalOutputCancellationReason = nil
self.acknowledgePlaybackMarks(self.takePendingPlaybackMarks())
self.retireOutputCancellation()
self.output.withLock { output in
output.terminalOutputCancellationReason = nil
output.acknowledgePlaybackMarks(output.takePendingPlaybackMarks())
output.retireCancellation()
}
self.retireOutputCancellationTask()
case nil, "applied":
guard result.turnid == nil || result.turnid == turnId else {
throw URLError(.badServerResponse)
}
if self.awaitingOutputClear {
self.outputCancellationTask = nil
let retired = self.output.withLock { output in
output.cancellationInFlight = false
if output.awaitingOutputClear { return false }
output.retireCancellation()
return true
}
if retired {
self.retireOutputCancellationTask()
} else {
self.retireOutputCancellation()
self.outputCancellationTask = nil
}
default:
throw URLError(.badServerResponse)
@ -1231,96 +1176,6 @@ extension RealtimeTalkRelaySession {
generation == self.outputCancellationGeneration && !self.isClosed
}
private func handleOutputAudio(_ payload: [String: AnyCodable]) {
guard !self.isOutputPaused else { return }
let incomingIdentity = OutputIdentity(payload)
guard let incomingTurnId = incomingIdentity.turnId else {
self.handleOutputPlaybackOverflow()
return
}
guard !self.awaitingOutputClear else { return }
guard incomingTurnId != self.cancelledOutputTurnId else { return }
guard let base64 = payload["audioBase64"]?.stringValue else { return }
guard let data = Data(base64Encoded: base64) else {
self.handleOutputPlaybackOverflow()
return
}
self.terminalOutputCancellationReason = nil
if let currentIdentity = self.outputIdentity,
currentIdentity != incomingIdentity
{
let marks = self.takePendingPlaybackMarks()
self.stopOutputPlayback()
self.acknowledgePlaybackMarks(marks)
} else if self.outputContinuation == nil, self.outputTask != nil {
self.stopOutputPlayback()
}
self.outputIdentity = incomingIdentity
self.recordOutputAudioChunk(byteCount: data.count)
self.markOutputAudioStarted(nowMs: ProcessInfo.processInfo.systemUptime * 1000)
self.onSpeakingChanged(true)
self.ensureOutputPlaybackStarted()
self.bufferOutputAudio(data)
}
private func bufferOutputAudio(_ data: Data) {
let frameByteCount = max(2, Int((self.outputSampleRateHz * 0.02).rounded()) * 2)
var offset = data.startIndex
if !self.pendingOutputAudio.isEmpty {
let fillCount = min(frameByteCount - self.pendingOutputAudio.count, data.count)
let fillEnd = data.index(offset, offsetBy: fillCount)
self.pendingOutputAudio.append(data[offset..<fillEnd])
offset = fillEnd
if self.pendingOutputAudio.count == frameByteCount {
let frame = self.pendingOutputAudio
self.pendingOutputAudio.removeAll(keepingCapacity: true)
guard self.yieldOutputAudioFrame(frame) else { return }
}
}
while data.distance(from: offset, to: data.endIndex) >= frameByteCount {
let frameEnd = data.index(offset, offsetBy: frameByteCount)
let frame = Data(data[offset..<frameEnd])
offset = frameEnd
guard self.yieldOutputAudioFrame(frame) else { return }
}
if offset < data.endIndex {
self.pendingOutputAudio.append(data[offset...])
}
}
private func yieldOutputAudioFrame(_ data: Data) -> Bool {
guard let continuation = self.outputContinuation else { return false }
switch continuation.yield(data) {
case .enqueued:
self.outputEnvelope?.append(data)
return true
case .dropped:
self.handleOutputPlaybackOverflow()
return false
case .terminated:
return false
@unknown default:
self.handleOutputPlaybackOverflow()
return false
}
}
private func handleOutputAudioDone(_ payload: [String: AnyCodable]) {
let incomingIdentity = OutputIdentity(payload)
if incomingIdentity.turnId != nil,
let outputIdentity,
outputIdentity != incomingIdentity
{
return
}
self.finishOutputPlaybackStream()
}
private func handleOutputPlaybackOverflow() {
self.handleOutputPlaybackFailure(
String(localized: "Realtime audio playback fell behind. Reconnecting…"))
}
private func handleOutputPlaybackFailure(_ message: String) {
guard !self.isClosed else { return }
let issue = self.issue(
@ -1332,12 +1187,15 @@ extension RealtimeTalkRelaySession {
self.onTermination(.outputPlaybackOverflow)
}
private func retireOutputCancellation() {
private func retireOutputCancellationTask() {
self.outputCancellationGeneration &+= 1
self.outputCancellationTask?.cancel()
self.outputCancellationTask = nil
self.suppressedOutputIdentity = nil
self.awaitingOutputClear = false
}
private func retireOutputCancellation() {
self.output.withLock { $0.retireCancellation() }
self.retireOutputCancellationTask()
}
}
@ -1395,12 +1253,12 @@ extension RealtimeTalkRelaySession {
{
guard self.isCurrentLifecycleLocally(lifecycleGeneration),
self.audioCaptureGeneration == audioCaptureGeneration,
!self.isInputPaused, self.suppressedOutputIdentity == nil,
let audioSender = self.audioSender
!self.isInputPaused, self.output.withLock({ $0.suppressedOutputIdentity }) == nil,
let audioSender
else { return nil }
self.recordMicrophoneFrame(byteCount: encoded.count, rms: rms, timestampMs: timestampMs)
// Continuous providers own interruptions and need input throughout playback.
if self.isOutputPlaying, self.supportsBargeIn == true {
if self.output.withLock({ $0.isOutputPlaying }), self.supportsBargeIn == true {
if self.audioCapture.suppressesInputDuringOutput {
self.recordSuppressedOutputEchoFrame(
byteCount: encoded.count,
@ -1417,7 +1275,7 @@ extension RealtimeTalkRelaySession {
defer { self.audioSendTasks.removeValue(forKey: taskID) }
guard self.isCurrentLifecycleLocally(lifecycleGeneration),
self.audioCaptureGeneration == audioCaptureGeneration,
!self.isInputPaused, self.suppressedOutputIdentity == nil
!self.isInputPaused, self.output.withLock({ $0.suppressedOutputIdentity }) == nil
else { return }
switch await audioSender.send(encoded, timestampMs: timestampMs) {
case .sent, .inactive:
@ -1485,6 +1343,7 @@ extension RealtimeTalkRelaySession {
// periphery:ignore - package tests drive a relay session without a live gateway handshake.
func _test_setRelaySessionId(_ relaySessionId: String) {
self.relaySessionId = relaySessionId
self.output.withLock { $0.relaySessionId = relaySessionId }
}
// periphery:ignore - package tests inject gateway events without a live socket.
@ -1523,22 +1382,23 @@ extension RealtimeTalkRelaySession {
// periphery:ignore - package tests start output playback without decoding real audio.
func _test_markOutputAudioStarted(nowMs: Double) {
self.markOutputAudioStarted(nowMs: nowMs)
self.output.withLock { $0.markOutputAudioStarted(nowMs: nowMs) }
}
// periphery:ignore - package tests finish playback without a real player callback.
func _test_markOutputPlaybackFinished() {
self.markOutputPlaybackFinished()
self.output.withLock { $0.markOutputPlaybackFinished() }
self.drainOutputEffects()
}
// periphery:ignore - package tests observe barge-in timing state.
func _test_outputStartedAtMs() -> Double? {
self.outputStartedAtMs
self.output.withLock { $0.outputStartedAtMs }
}
// periphery:ignore - package tests observe playback state without exposing it publicly.
func _test_isOutputPlaying() -> Bool {
self.isOutputPlaying
nonisolated func _test_isOutputPlaying() -> Bool {
self.output.withLock { $0.isOutputPlaying }
}
// periphery:ignore - package tests exercise the audio sender without a started session.

View file

@ -11,56 +11,84 @@ private struct RealtimePCMPlaybackFailure: Error {}
private let realtimePCMPlaybackWaitTimeoutSeconds = 15.0
@MainActor
private final class RealtimePCMPlaybackBackend {
private final class RealtimePCMPlaybackBackend: @unchecked Sendable {
private let lock = NSRecursiveLock()
private struct Waiter {
let count: Int
let continuation: CheckedContinuation<Void, any Error>
}
private(set) var scheduledFrames: [Data] = []
private(set) var completions: [@Sendable () -> Void] = []
private(set) var activeCount = 0
private(set) var maxActiveCount = 0
private var storedScheduledFrames: [Data] = []
private var storedCompletions: [@Sendable () -> Void] = []
private var storedActiveCount = 0
private var storedMaxActiveCount = 0
private var completedCallbacks = 0
private var scheduledWaiters: [UUID: Waiter] = [:]
private var completionWaiters: [UUID: Waiter] = [:]
func prepare(sampleRate _: Double) throws {}
var scheduledFrames: [Data] {
self.lock.withLock { self.storedScheduledFrames }
}
var completions: [@Sendable () -> Void] {
self.lock.withLock { self.storedCompletions }
}
var activeCount: Int {
self.lock.withLock { self.storedActiveCount }
}
var maxActiveCount: Int {
self.lock.withLock { self.storedMaxActiveCount }
}
func prepare(sampleRate _: Double) throws {
self.lock.withLock {}
}
func schedule(
data: Data,
sampleRate _: Double,
completion: @escaping @Sendable () -> Void) throws
{
self.scheduledFrames.append(data)
self.activeCount += 1
self.maxActiveCount = max(self.maxActiveCount, self.activeCount)
self.resumeScheduledWaiters()
self.completions.append { [weak self] in
Task { @MainActor in
self?.activeCount -= 1
self?.completedCallbacks += 1
self?.resumeCompletionWaiters()
self.lock.withLock {
self.storedScheduledFrames.append(data)
self.storedActiveCount += 1
self.storedMaxActiveCount = max(self.storedMaxActiveCount, self.storedActiveCount)
self.resumeScheduledWaiters()
self.storedCompletions.append { [weak self] in
self?.lock.withLock {
self?.storedActiveCount -= 1
self?.completedCallbacks += 1
self?.resumeCompletionWaiters()
}
completion()
}
}
}
func stop() {
self.activeCount = 0
self.lock.withLock {
self.storedActiveCount = 0
}
}
func complete(at index: Int = 0) {
self.completions.remove(at: index)()
self.lock.withLock {
self.storedCompletions.remove(at: index)()
}
}
func takeCompletion(at index: Int = 0) -> @Sendable () -> Void {
self.completions.remove(at: index)
self.lock.withLock {
self.storedCompletions.remove(at: index)
}
}
func waitForScheduledFrames(_ count: Int) async throws {
if self.scheduledFrames.count >= count { return }
if self.lock.withLock({ self.storedScheduledFrames.count >= count }) {
return
}
try await AsyncTimeout.withTimeout(
seconds: realtimePCMPlaybackWaitTimeoutSeconds,
onTimeout: { RealtimePCMPlaybackWaitTimeout(label: "scheduled frames \(count)") },
@ -68,7 +96,9 @@ private final class RealtimePCMPlaybackBackend {
}
func waitForCompletionCallbacks(_ count: Int) async throws {
if self.completedCallbacks >= count { return }
if self.lock.withLock({ self.completedCallbacks >= count }) {
return
}
try await AsyncTimeout.withTimeout(
seconds: realtimePCMPlaybackWaitTimeoutSeconds,
onTimeout: { RealtimePCMPlaybackWaitTimeout(label: "completion callbacks \(count)") },
@ -79,14 +109,16 @@ private final class RealtimePCMPlaybackBackend {
let id = UUID()
try await withTaskCancellationHandler {
try await withCheckedThrowingContinuation { continuation in
if self.scheduledFrames.count >= count {
continuation.resume()
} else {
self.scheduledWaiters[id] = Waiter(count: count, continuation: continuation)
self.lock.withLock {
if self.storedScheduledFrames.count >= count {
continuation.resume()
} else {
self.scheduledWaiters[id] = Waiter(count: count, continuation: continuation)
}
}
}
} onCancel: {
Task { @MainActor in self.cancelScheduledWaiter(id) }
self.cancelScheduledWaiter(id)
}
}
@ -94,38 +126,48 @@ private final class RealtimePCMPlaybackBackend {
let id = UUID()
try await withTaskCancellationHandler {
try await withCheckedThrowingContinuation { continuation in
if self.completedCallbacks >= count {
continuation.resume()
} else {
self.completionWaiters[id] = Waiter(count: count, continuation: continuation)
self.lock.withLock {
if self.completedCallbacks >= count {
continuation.resume()
} else {
self.completionWaiters[id] = Waiter(count: count, continuation: continuation)
}
}
}
} onCancel: {
Task { @MainActor in self.cancelCompletionWaiter(id) }
self.cancelCompletionWaiter(id)
}
}
private func cancelScheduledWaiter(_ id: UUID) {
self.scheduledWaiters.removeValue(forKey: id)?.continuation.resume(throwing: CancellationError())
self.lock.withLock {
self.scheduledWaiters.removeValue(forKey: id)?.continuation.resume(throwing: CancellationError())
}
}
private func cancelCompletionWaiter(_ id: UUID) {
self.completionWaiters.removeValue(forKey: id)?.continuation.resume(throwing: CancellationError())
self.lock.withLock {
self.completionWaiters.removeValue(forKey: id)?.continuation.resume(throwing: CancellationError())
}
}
private func resumeScheduledWaiters() {
let ready = self.scheduledWaiters.filter { self.scheduledFrames.count >= $0.value.count }
for (id, waiter) in ready {
self.scheduledWaiters.removeValue(forKey: id)
waiter.continuation.resume()
self.lock.withLock {
let ready = self.scheduledWaiters.filter { self.storedScheduledFrames.count >= $0.value.count }
for (id, waiter) in ready {
self.scheduledWaiters.removeValue(forKey: id)
waiter.continuation.resume()
}
}
}
private func resumeCompletionWaiters() {
let ready = self.completionWaiters.filter { self.completedCallbacks >= $0.value.count }
for (id, waiter) in ready {
self.completionWaiters.removeValue(forKey: id)
waiter.continuation.resume()
self.lock.withLock {
let ready = self.completionWaiters.filter { self.completedCallbacks >= $0.value.count }
for (id, waiter) in ready {
self.completionWaiters.removeValue(forKey: id)
waiter.continuation.resume()
}
}
}
}
@ -139,6 +181,18 @@ private final class RealtimePCMPlaybackResultProbe {
}
}
private final class RealtimePCMStartCounter: @unchecked Sendable {
private let lock = NSLock()
private var storedCount = 0
var count: Int {
self.lock.withLock { self.storedCount }
}
func increment() {
self.lock.withLock { self.storedCount += 1 }
}
}
@MainActor
private func makeRealtimePCMPlayer(
backend: RealtimePCMPlaybackBackend) -> RealtimePCMStreamingAudioPlayer
@ -203,7 +257,46 @@ struct RealtimePCMStreamingAudioPlayerTests {
#expect(probe.results.first?.interruptedAt == nil)
}
@Test func `playback starts after the prebuffer or at end of a short reply`() async throws {
for (frames, finish, expectedStarts) in [
(RealtimePCMStreamingAudioPlayer.prebufferFrames - 1, false, 0),
(RealtimePCMStreamingAudioPlayer.prebufferFrames, false, 1),
(2, true, 1),
] {
let backend = RealtimePCMPlaybackBackend()
let starts = RealtimePCMStartCounter()
let started = RealtimeRelayTestSignal<Void>(timeoutSeconds: 5)
let player = RealtimePCMStreamingAudioPlayer(
preparePlayback: backend.prepare,
scheduleFrame: backend.schedule,
startPlayback: {
starts.increment()
started.send(())
},
stopPlayback: backend.stop,
playbackTime: { nil })
let (stream, continuation) = AsyncThrowingStream<Data, Error>.makeStream()
let playback = Task { _ = await player.play(stream: stream, sampleRate: self.sampleRate) }
continuation.yield(Data(repeating: 1, count: self.frameBytes * frames))
if finish {
continuation.finish()
}
try await backend.waitForScheduledFrames(frames)
if expectedStarts > 0 {
// The prebuffer start and the end-of-input start both run after scheduling returns.
_ = try await started.next("playback start frames=\(frames) finish=\(finish)")
}
// Every start decision runs on the backend queue; drain it so no decision is still pending.
await player._test_waitForBackendOperations()
#expect(starts.count == expectedStarts, "frames=\(frames) finish=\(finish)")
_ = player.stop()
continuation.finish()
await playback.value
}
}
@Test func `withheld completions cap scheduling and one completion admits one frame`() async throws {
let cap = RealtimePCMStreamingAudioPlayer.maxScheduledBuffers
let backend = RealtimePCMPlaybackBackend()
let player = makeRealtimePCMPlayer(backend: backend)
let probe = RealtimePCMPlaybackResultProbe()
@ -213,29 +306,29 @@ struct RealtimePCMStreamingAudioPlayerTests {
probe.record(result)
}
continuation.yield(Data(repeating: 1, count: self.frameBytes * 5))
continuation.yield(Data(repeating: 1, count: self.frameBytes * (cap + 2)))
continuation.finish()
try await backend.waitForScheduledFrames(3)
#expect(backend.scheduledFrames.count == 3)
#expect(backend.maxActiveCount == 3)
try await backend.waitForScheduledFrames(cap)
#expect(backend.scheduledFrames.count == cap)
#expect(backend.maxActiveCount == cap)
#expect(probe.results.isEmpty)
backend.complete()
try await backend.waitForScheduledFrames(4)
#expect(backend.scheduledFrames.count == 4)
#expect(backend.maxActiveCount == 3)
try await backend.waitForScheduledFrames(cap + 1)
#expect(backend.scheduledFrames.count == cap + 1)
#expect(backend.maxActiveCount == cap)
#expect(probe.results.isEmpty)
backend.complete()
try await backend.waitForScheduledFrames(5)
for _ in 0..<3 {
try await backend.waitForScheduledFrames(cap + 2)
for _ in 0..<cap {
backend.complete()
}
try await waitForPlayback(playback, label: "five-frame playback")
try await waitForPlayback(playback, label: "cap-plus-two-frame playback")
#expect(probe.results.count == 1)
#expect(probe.results.first?.finished == true)
#expect(probe.results.first?.interruptedAt == nil)
#expect(backend.scheduledFrames.count == 5)
#expect(backend.scheduledFrames.count == cap + 2)
#expect(backend.scheduledFrames.allSatisfy { $0.count == self.frameBytes })
}
@ -266,6 +359,7 @@ struct RealtimePCMStreamingAudioPlayerTests {
}
@Test func `stop restart ignores stale buffer completions`() async throws {
let cap = RealtimePCMStreamingAudioPlayer.maxScheduledBuffers
let backend = RealtimePCMPlaybackBackend()
let player = makeRealtimePCMPlayer(backend: backend)
let (firstStream, firstContinuation) = AsyncThrowingStream<Data, Error>.makeStream()
@ -288,27 +382,27 @@ struct RealtimePCMStreamingAudioPlayerTests {
let result = await player.play(stream: secondStream, sampleRate: self.sampleRate)
probe.record(result)
}
secondContinuation.yield(Data(repeating: 2, count: self.frameBytes * 5))
secondContinuation.yield(Data(repeating: 2, count: self.frameBytes * (cap + 2)))
secondContinuation.finish()
try await backend.waitForScheduledFrames(4)
#expect(backend.activeCount == 3)
try await backend.waitForScheduledFrames(cap + 1)
#expect(backend.activeCount == cap)
#expect(probe.results.isEmpty)
staleCompletion()
try await backend.waitForCompletionCallbacks(1)
#expect(backend.scheduledFrames.count == 4)
#expect(backend.completions.count == 3)
#expect(backend.scheduledFrames.count == cap + 1)
#expect(backend.completions.count == cap)
#expect(firstProbe.results.map(\.finished) == [false])
#expect(probe.results.isEmpty)
backend.complete()
try await backend.waitForScheduledFrames(5)
#expect(backend.scheduledFrames.count == 5)
try await backend.waitForScheduledFrames(cap + 2)
#expect(backend.scheduledFrames.count == cap + 2)
#expect(probe.results.isEmpty)
backend.complete()
try await backend.waitForScheduledFrames(6)
#expect(backend.scheduledFrames.count == 6)
try await backend.waitForScheduledFrames(cap + 3)
#expect(backend.scheduledFrames.count == cap + 3)
#expect(probe.results.isEmpty)
for _ in 0..<3 {
for _ in 0..<cap {
backend.complete()
}
try await waitForPlayback(secondPlayback, label: "replacement B playback")

View file

@ -199,7 +199,7 @@ struct RealtimeTalkRelaySessionAudioInputTests {
#expect(timestamp == timestamp.rounded())
}
@Test func `microphone saturation terminates once without sending the fifth frame`() async throws {
@Test func `microphone saturation terminates once without sending the frame past the cap`() async throws {
let requests = ControlledRealtimeAudioRequests()
let audioCapture = TestRealtimeTalkAudioCapture()
var statuses: [String] = []
@ -228,20 +228,21 @@ struct RealtimeTalkRelaySessionAudioInputTests {
var pending: [Task<Void, Never>] = []
var saturated: Task<Void, Never>?
do {
for index in 0..<4 {
guard let send = session._test_enqueueMicrophoneFrame(Data([UInt8(index)])) else {
let cap = RealtimeTalkRelaySession.maxPendingAudioSends
for index in 0..<cap {
guard let send = session._test_enqueueMicrophoneFrame(Data([UInt8(truncatingIfNeeded: index)])) else {
throw RealtimeRelayTestTimeout(operation: "microphone frame \(index) admission")
}
pending.append(send)
}
try await requests.waitForRequestCount(4)
try await requests.waitForRequestCount(cap)
guard let saturationSend = session._test_enqueueMicrophoneFrame(Data([0xFF])) else {
throw RealtimeRelayTestTimeout(operation: "saturation frame admission")
}
saturated = saturationSend
_ = try await terminationObserved.next("microphone saturation termination")
await saturated?.value
try await requests.waitForRequestCount(5)
try await requests.waitForRequestCount(cap + 1)
} catch {
saturated?.cancel()
pending.forEach { $0.cancel() }
@ -255,13 +256,9 @@ struct RealtimeTalkRelaySessionAudioInputTests {
}
let message = String(localized: "Realtime audio input fell behind. Reconnecting…")
#expect(await requests.snapshot() == [
"talk.session.appendAudio",
"talk.session.appendAudio",
"talk.session.appendAudio",
"talk.session.appendAudio",
"talk.session.close",
])
#expect(await requests.snapshot() ==
Array(repeating: "talk.session.appendAudio", count: RealtimeTalkRelaySession.maxPendingAudioSends)
+ ["talk.session.close"])
#expect(statuses == [message])
#expect(issues.map(\.code) == ["audio_input_unavailable"])
#expect(issues.map(\.message) == [message])
@ -279,6 +276,39 @@ struct RealtimeTalkRelaySessionAudioInputTests {
#expect(await requests.snapshot().filter { $0 == "talk.session.close" }.count == 1)
}
@Test func `a one second gateway stall does not saturate microphone input`() async throws {
// Gateway appendAudio round trips over Wi-Fi/Tailscale reach ~1 s; ~23 mic frames stay in flight.
let requests = ControlledRealtimeAudioRequests()
var terminations: [RealtimeTalkRelayTermination] = []
let session = RealtimeTalkRelaySession(
transport: RealtimeTalkRelayTransport(
subscribeServerEvents: { _ in AsyncStream { $0.finish() } },
request: { method, _, _ in try await requests.request(method: method) }),
options: .init(sessionKey: "main", provider: "xai", model: nil, voice: nil),
audioCapture: TestRealtimeTalkAudioCapture(),
pcmPlayer: UnusedPCMStreamingAudioPlayer(),
onStatus: { _ in },
onIssue: { _ in },
onTermination: { terminations.append($0) },
onSpeakingChanged: { _ in })
session._test_setRelaySessionId("relay-1")
session._test_prepareAudioSender(relaySessionId: "relay-1")
try session._test_startMicrophonePump()
var pending: [Task<Void, Never>] = []
for index in 0..<23 {
if let send = session._test_enqueueMicrophoneFrame(Data([UInt8(index)])) { pending.append(send) }
}
try await requests.waitForRequestCount(23)
#expect(terminations.isEmpty)
session.stop()
await requests.succeedPendingAppends()
for task in pending {
await task.value
}
}
@Test func `active audio request and response failures share the input failure owner`() async throws {
for behavior in [ControlledAudioAppendBehavior.requestFailure, .malformedResponse] {
let requests = ControlledRealtimeAudioRequests(behavior: behavior)

View file

@ -0,0 +1,186 @@
#if Talk && canImport(ElevenLabsKit) && (os(iOS) || os(macOS))
import Foundation
import OpenClawProtocol
import Testing
@testable import OpenClawKit
/// A fake audio device, not an actor: records the actual scheduleBuffer boundary.
private final class OffMainPlaybackProbe: @unchecked Sendable {
private let lock = NSLock()
private var frames: [Data] = []
private var scheduledOnMain = false
private let expectedFrames: Int
/// Signalled from the audio boundary itself, so a blocked main thread can wait on it.
private let allScheduled = DispatchSemaphore(value: 0)
init(expectedFrames: Int = .max) {
self.expectedFrames = expectedFrames
}
func schedule(_ data: Data, sampleRate _: Double, completion _: @escaping @Sendable () -> Void) {
let reachedExpected = self.lock.withLock {
self.frames.append(data)
self.scheduledOnMain = self.scheduledOnMain || Thread.isMainThread
return self.frames.count == self.expectedFrames
}
if reachedExpected { self.allScheduled.signal() }
}
/// Synchronous on purpose: the caller's thread stays blocked, it never suspends or yields.
func blockUntilAllScheduled(timeoutSeconds: Double) -> Bool {
self.allScheduled.wait(timeout: .now() + timeoutSeconds) == .success
}
func snapshot() -> (frames: [Data], scheduledOnMain: Bool) {
self.lock.withLock { (self.frames, self.scheduledOnMain) }
}
}
@MainActor
struct RealtimeTalkRelaySessionOffMainTests {
@Test func `buffered startup clear acknowledges its playback mark`() async throws {
let events = AsyncStream<EventFrame>.makeStream()
let createBarrier = RealtimeRelayStartupBarrier()
let acknowledged = RealtimeRelayTestSignal<[String: AnyCodable]>()
let created = try JSONEncoder().encode(TalkSessionCreateResult(
sessionid: "relay-1",
mode: AnyCodable("realtime"),
transport: AnyCodable("gateway-relay"),
brain: AnyCodable("agent-consult"),
relaysessionid: "relay-1"))
let session = RealtimeTalkRelaySession(
transport: RealtimeTalkRelayTransport(
subscribeServerEvents: { _ in events.stream },
request: { method, params, _ in
if method == "talk.session.create" {
await createBarrier.suspend()
return created
}
if method == "talk.catalog" {
return try realtimeRelayCatalogData()
}
if method == "talk.session.acknowledgeMark", let params {
acknowledged.send(params)
}
return Data("{\"ok\":true}".utf8)
}),
options: .init(sessionKey: "main", provider: nil, model: nil, voice: nil),
audioCapture: TestRealtimeTalkAudioCapture(),
pcmPlayer: UnusedPCMStreamingAudioPlayer(),
onStatus: { _ in },
onSpeakingChanged: { _ in })
defer { session.stop()
events.continuation.finish()
}
let starting = Task { try await session.start() }
try await createBarrier.waitUntilEntered()
// Drive the same main startup consumer deterministically while creation is suspended.
await session._test_handleGatewayEvent(outputAudioEvent(turnId: "startup", data: Data([1, 1])))
await session._test_handleGatewayEvent(EventFrame(
type: "event", event: "talk.event",
payload: AnyCodable([
"relaySessionId": "relay-1", "type": "mark", "markName": "startup-mark",
])))
await session._test_handleGatewayEvent(outputClearEvent(turnId: "startup"))
await session._test_handleGatewayEvent(EventFrame(
type: "event", event: "talk.event",
payload: AnyCodable(["relaySessionId": "relay-1", "type": "ready"])))
await createBarrier.release()
try await starting.value
let params = try await AsyncTimeout.withTimeout(
seconds: 2,
onTimeout: { RealtimeRelayTestTimeout(operation: "startup mark acknowledgment") },
operation: { try await acknowledged.next("startup mark acknowledgment") })
#expect(params["sessionId"]?.stringValue == "relay-1")
#expect(params["markName"]?.stringValue == "startup-mark")
}
@Test func `installed relay identity does not open startup routing`() {
let notifications = AsyncStream<Void>.makeStream()
let output = RealtimeTalkOutput(
player: UnusedPCMStreamingAudioPlayer(), transport: unusedRealtimeRelayTransport(),
notification: notifications.continuation)
output.withLock {
$0.resetRouting(lifecycleGeneration: 1)
$0.relaySessionId = "relay-1"
}
let event = outputAudioEvent(turnId: "startup", data: Data([1, 1]))
let buffered = output.route(event, lifecycleGeneration: 1)
#expect(!buffered.handled && buffered.startup)
output.withLock { $0.startupRoutingReady = true }
let behindBuffer = output.route(event, lifecycleGeneration: 1)
#expect(!behindBuffer.handled && behindBuffer.startup)
output.withLock {
$0.mainEventHandled(startup: true, lifecycleGeneration: 1)
$0.mainEventHandled(startup: true, lifecycleGeneration: 1)
}
let live = output.route(event, lifecycleGeneration: 1)
#expect(live.handled && !live.startup)
output.withLock { $0.stopOutputPlayback() }
}
@Test func `new reply reaches audio device while main actor is stalled`() async throws {
let events = AsyncStream<EventFrame>.makeStream()
let created = try JSONEncoder().encode(TalkSessionCreateResult(
sessionid: "relay-1",
mode: AnyCodable("realtime"),
transport: AnyCodable("gateway-relay"),
brain: AnyCodable("agent-consult"),
relaysessionid: "relay-1"))
let probe = OffMainPlaybackProbe(expectedFrames: 81)
let player = RealtimePCMStreamingAudioPlayer(
preparePlayback: { _ in },
scheduleFrame: probe.schedule,
stopPlayback: {},
playbackTime: { nil })
let session = RealtimeTalkRelaySession(
transport: RealtimeTalkRelayTransport(
subscribeServerEvents: { _ in events.stream },
request: { method, _, _ in
if method == "talk.session.create" {
events.continuation.yield(EventFrame(
type: "event",
event: "talk.event",
payload: AnyCodable(["relaySessionId": "relay-1", "type": "ready"])))
return created
}
if method == "talk.catalog" {
return try realtimeRelayCatalogData()
}
return Data("{\"ok\":true}".utf8)
}),
options: .init(sessionKey: "main", provider: nil, model: nil, voice: nil),
audioCapture: TestRealtimeTalkAudioCapture(),
pcmPlayer: player,
onStatus: { _ in },
onSpeakingChanged: { _ in })
defer { session.stop()
events.continuation.finish()
}
try await session.start()
let stallStarted = RealtimeRelayTestSignal<Void>()
let producer = Task.detached {
_ = try await stallStarted.next("main actor stall")
for index in 0..<80 {
events.continuation.yield(outputAudioEvent(
turnId: "brand-new-turn", data: Data(repeating: UInt8(index + 1), count: 960)))
}
events.continuation.yield(outputAudioEvent(turnId: "brand-new-turn", data: Data([81, 81])))
events.continuation.yield(outputAudioDoneEvent(turnId: "brand-new-turn"))
}
stallStarted.send(())
// Block the main thread until the audio boundary reports every frame, bounded so a
// main-dependent implementation fails instead of hanging. Nothing here yields the main actor.
let scheduledDuringStall = probe.blockUntilAllScheduled(timeoutSeconds: 10)
// Snapshot BEFORE releasing the main actor: after-the-stall delivery cannot pass.
let duringStall = probe.snapshot()
#expect(scheduledDuringStall, "scheduling must complete while the main actor is blocked")
#expect(
duringStall.frames.count == 81,
"all frames, including the first and partial tail, must schedule DURING stall")
#expect(duringStall.frames.first == Data(repeating: 1, count: 960), "the start of the reply must survive")
#expect(!duringStall.scheduledOnMain, "scheduleBuffer must never require main")
try await producer.value
}
}
#endif

View file

@ -75,6 +75,55 @@ struct RealtimeTalkRelaySessionPlaybackTests {
#expect(request.params?["markName"]?.stringValue == "audio-1")
}
@Test func `speaking is reported once per reply, not per audio chunk`() async throws {
var speakingReports: [Bool] = []
let session = RealtimeTalkRelaySession(
transport: RealtimeTalkRelayTransport(
subscribeServerEvents: { _ in AsyncStream { $0.finish() } },
request: { _, _, _ in Data("{\"ok\":true}".utf8) }),
options: .init(sessionKey: "main", provider: "xai", model: nil, voice: nil),
audioCapture: TestRealtimeTalkAudioCapture(),
pcmPlayer: StalledPCMStreamingAudioPlayer(),
onStatus: { _ in },
onIssue: { _ in },
onTermination: { _ in },
onSpeakingChanged: { speakingReports.append($0) })
session._test_setRelaySessionId("relay-1")
for _ in 0..<20 {
await session._test_handleGatewayEvent(
outputAudioEvent(turnId: "turn-1", data: Data(repeating: 1, count: 960)))
}
#expect(speakingReports == [true])
}
@Test func `a whole reply delivered in one burst does not overflow playback`() async throws {
// xAI delivers a reply faster than realtime; 10 s arriving before playback drains is normal.
let player = StalledPCMStreamingAudioPlayer()
var terminations: [RealtimeTalkRelayTermination] = []
let session = RealtimeTalkRelaySession(
transport: RealtimeTalkRelayTransport(
subscribeServerEvents: { _ in AsyncStream { $0.finish() } },
request: { _, _, _ in Data("{\"ok\":true}".utf8) }),
options: .init(sessionKey: "main", provider: "xai", model: nil, voice: nil),
audioCapture: TestRealtimeTalkAudioCapture(),
pcmPlayer: player,
onStatus: { _ in },
onIssue: { _ in },
onTermination: { terminations.append($0) },
onSpeakingChanged: { _ in })
session._test_setRelaySessionId("relay-1")
for _ in 0..<500 {
await session._test_handleGatewayEvent(
outputAudioEvent(turnId: "turn-1", data: Data(repeating: 1, count: 960)))
}
#expect(terminations.isEmpty)
#expect(player.stopCount == 0)
}
@Test func `output buffer cap plus one terminates visibly and requests recovery`() async throws {
let requests = RealtimeRelayStartupRequestLog()
let player = StalledPCMStreamingAudioPlayer()
@ -100,7 +149,7 @@ struct RealtimeTalkRelaySessionPlaybackTests {
onSpeakingChanged: { _ in })
session._test_setRelaySessionId("relay-1")
for _ in 0...32 {
for _ in 0...RealtimeTalkRelaySession.maxBufferedOutputChunks {
await session._test_handleGatewayEvent(
outputAudioEvent(turnId: "turn-1", data: Data(repeating: 1, count: 960)))
}

View file

@ -0,0 +1,352 @@
#if Talk && canImport(ElevenLabsKit) && (os(iOS) || os(macOS))
import Foundation
import OpenClawProtocol
import Testing
@testable import OpenClawKit
/// Exercises the production player, including backend operations and withheld device drains.
final class RealtimeRelayDevice: @unchecked Sendable {
private let lock = NSLock()
private var storedFrames: [Data] = []
private var callbacks: [@Sendable () -> Void] = []
private var storedStops = 0
let scheduled = RealtimeRelayTestSignal<Int>(timeoutSeconds: 5)
let prepared = RealtimeRelayTestSignal<Void>(timeoutSeconds: 5)
let stopped = RealtimeRelayTestSignal<Int>(timeoutSeconds: 5)
let prepareGate: DispatchSemaphore?
let frameGate: DispatchSemaphore?
init(prepareGate: DispatchSemaphore? = nil, frameGate: DispatchSemaphore? = nil) {
self.prepareGate = prepareGate
self.frameGate = frameGate
}
var frames: [Data] {
self.lock.withLock { self.storedFrames }
}
var stopCount: Int {
self.lock.withLock { self.storedStops }
}
func prepare(_: Double) {
self.prepared.send(())
if let prepareGate { _ = prepareGate.wait(timeout: .now() + 2) }
}
func schedule(_ data: Data, _: Double, completion: @escaping @Sendable () -> Void) {
let count = self.lock.withLock {
self.storedFrames.append(data)
self.callbacks.append(completion)
return self.storedFrames.count
}
self.scheduled.send(count)
if let frameGate { _ = frameGate.wait(timeout: .now() + 2) }
}
func stop() {
let count = self.lock.withLock {
self.storedStops += 1
return self.storedStops
}
self.stopped.send(count)
}
func completion(at index: Int) -> @Sendable () -> Void {
self.lock.withLock { self.callbacks[index] }
}
func waitForFrames(_ count: Int) async throws {
while self.frames.count < count {
_ = try await self.scheduled.next("device frames")
}
}
}
@MainActor
final class RealtimeRelayProductionFixture {
let events = AsyncStream<EventFrame>.makeStream()
let requests = RealtimeRelayStartupRequestLog()
let capture = TestRealtimeTalkAudioCapture()
let device: RealtimeRelayDevice
let player: RealtimePCMStreamingAudioPlayer
let speaking = RealtimeRelayTestSignal<Bool>(timeoutSeconds: 5)
let levels = RealtimeRelayTestSignal<Double?>(timeoutSeconds: 5)
let terminated = RealtimeRelayTestSignal<RealtimeTalkRelayTermination>(timeoutSeconds: 5)
var session: RealtimeTalkRelaySession!
init(
device: RealtimeRelayDevice = RealtimeRelayDevice(),
supportsBargeIn: Bool = true,
onSpeaking: @escaping (Bool) -> Void = { _ in }) throws
{
self.device = device
let created = try JSONEncoder().encode(TalkSessionCreateResult(
sessionid: "relay-1", mode: AnyCodable("realtime"), transport: AnyCodable("gateway-relay"),
brain: AnyCodable("agent-consult"), relaysessionid: "relay-1"))
self.player = RealtimePCMStreamingAudioPlayer(
preparePlayback: device.prepare, scheduleFrame: device.schedule,
stopPlayback: device.stop, playbackTime: { nil })
self.session = RealtimeTalkRelaySession(
transport: RealtimeTalkRelayTransport(
subscribeServerEvents: { [events] _ in events.stream },
request: { [requests, events] method, params, _ in
await requests.record(method: method, params: params)
switch method {
case "talk.session.create":
events.continuation.yield(EventFrame(
type: "event",
event: "talk.event",
payload: AnyCodable([
"relaySessionId": "relay-1",
"type": "ready",
])))
return created
case "talk.catalog": return try realtimeRelayCatalogData(supportsBargeIn: supportsBargeIn)
case "talk.session.cancelOutput": return Data(#"{"ok":true,"status":"applied"}"#.utf8)
default: return Data(#"{"ok":true}"#.utf8)
}
}),
options: .init(sessionKey: "main", provider: nil, model: nil, voice: nil),
audioCapture: self.capture, pcmPlayer: self.player, onStatus: { _ in },
onTermination: { [terminated] in terminated.send($0) },
onSpeakingChanged: { [speaking] in speaking.send($0)
onSpeaking($0)
},
onOutputLevel: { [levels] in levels.send($0) })
}
func start() async throws {
try await self.session.start()
}
func send(_ event: EventFrame) {
self.events.continuation.yield(event)
}
func close() {
self.session?.stop()
self.events.continuation.finish()
}
}
extension RealtimeTalkRelaySessionPlaybackTests {
@Test func `production identity switch ignores stale completion and acknowledges marks only after drain`() async throws {
let f = try RealtimeRelayProductionFixture()
defer { f.close() }
try await f.start()
f.send(outputAudioEvent(turnId: "a", data: Data(repeating: 1, count: 960 * 20)))
try await f.device.waitForFrames(20)
f.send(playbackMarkEvent("a-mark"))
f.send(outputAudioEvent(turnId: "b", data: Data(repeating: 2, count: 960 * 20)))
f.send(playbackMarkEvent("b-mark"))
f.send(outputAudioDoneEvent(turnId: "b"))
try await f.device.waitForFrames(40)
try await f.requests.waitForRequestCount(3)
#expect(await f.requests.snapshot().filter { $0.method == "talk.session.acknowledgeMark" }
.map { $0.params?["markName"]?.stringValue } == ["a-mark"])
for index in 0..<20 {
f.device.completion(at: index)()
}
#expect(f.session._test_isOutputPlaying())
for index in 20..<40 {
f.device.completion(at: index)()
}
try await f.requests.waitForRequestCount(4)
#expect(await f.requests.snapshot().filter { $0.method == "talk.session.acknowledgeMark" }
.map { $0.params?["markName"]?.stringValue } == ["a-mark", "b-mark"])
#expect(!f.session._test_isOutputPlaying())
#expect(f.device.frames == Array(repeating: Data(repeating: 1, count: 960), count: 20)
+ Array(repeating: Data(repeating: 2, count: 960), count: 20))
}
}
@MainActor
@Suite(.serialized)
struct RealtimeTalkRelaySessionReviewTests {
@Test func `production overflow stops the device exactly once`() async throws {
let gate = DispatchSemaphore(value: 0)
let f = try RealtimeRelayProductionFixture(device: RealtimeRelayDevice(prepareGate: gate))
defer { gate.signal()
f.close()
}
try await f.start()
f.send(outputAudioEvent(turnId: "a", data: Data(repeating: 1, count: 960)))
_ = try await f.device.prepared.next("prepare entered")
// Preparation must not prevent admission/overflow fencing of a single provider burst.
f.send(outputAudioEvent(turnId: "a", data: Data(
repeating: 1,
count: 960 *
(RealtimeTalkRelaySession.maxBufferedOutputChunks * 3 + 1))))
#expect(try await f.terminated.next("overflow") == .outputPlaybackOverflow)
gate.signal()
_ = try await f.device.stopped.next("overflow device stop")
#expect(f.device.stopCount == 1)
#expect(f.device.frames.isEmpty)
}
@Test func `production stop mid-frame promptly fences remaining frames`() async throws {
let gate = DispatchSemaphore(value: 0)
let f = try RealtimeRelayProductionFixture(device: RealtimeRelayDevice(frameGate: gate))
defer { gate.signal()
f.close()
}
try await f.start()
f.send(outputAudioEvent(turnId: "a", data: Data(repeating: 1, count: 960 * 10)))
_ = try await f.device.scheduled.next("schedule entered")
let started = ProcessInfo.processInfo.systemUptime
f.session.stop()
#expect(ProcessInfo.processInfo.systemUptime - started < 0.25)
gate.signal()
_ = try await f.device.stopped.next("device stop")
#expect(f.device.frames.count == 1)
#expect(!f.session._test_isOutputPlaying())
}
}
extension RealtimeTalkRelaySessionCancellationTests {
@Test(arguments: ["user", "barge-in"])
func `production cancellation fences late audio through response and clear`(reason: String) async throws {
let f = try RealtimeRelayProductionFixture()
defer { f.close() }
try await f.start()
f.send(outputAudioEvent(turnId: "a", data: Data(repeating: 1, count: 960)))
try await f.device.waitForFrames(1)
if reason == "barge-in" {
f.capture.emit(RealtimeTalkAudioFrame(
data: Data([1, 2]),
timestampMs: ProcessInfo.processInfo.systemUptime * 1000 + 1000,
rms: 0.5))
} else { #expect(f.session.cancelOutput()) }
try await f.requests.waitForRequestCount(3)
let cancel = try #require(await f.requests.snapshot().last)
#expect(cancel.method == "talk.session.cancelOutput")
#expect(cancel.params?["turnId"]?.stringValue == "a")
#expect(cancel.params?["reason"]?.stringValue == reason)
let task = f.session._test_outputCancellationTask()
await task?.value
f.send(outputAudioEvent(turnId: "a", data: Data(repeating: 3, count: 960)))
f.send(outputAudioEvent(turnId: "b", data: Data(repeating: 3, count: 960)))
f.send(outputClearEvent(turnId: "a"))
f.send(outputAudioEvent(turnId: "a", data: Data(repeating: 3, count: 960)))
f.send(outputAudioEvent(turnId: "b", data: Data(repeating: 2, count: 960)))
try await f.device.waitForFrames(2)
#expect(f.device.frames == [Data(repeating: 1, count: 960), Data(repeating: 2, count: 960)])
}
}
extension RealtimeTalkRelaySessionReviewTests {
@Test func `production barge-in callbacks do not hold the output lock`() async throws {
var session: RealtimeTalkRelaySession?
let f = try RealtimeRelayProductionFixture(onSpeaking: { speaking in
guard !speaking, let session else { return }
let done = DispatchSemaphore(value: 0)
DispatchQueue.global(qos: .userInitiated).async {
_ = session._test_isOutputPlaying()
done.signal()
}
#expect(
done.wait(timeout: .now() + 0.25) == .success,
"app callback must leave the audio owner accessible off-main")
})
session = f.session
defer { f.close() }
try await f.start()
f.send(outputAudioEvent(turnId: "a", data: Data(repeating: 1, count: 960)))
try await f.device.waitForFrames(1)
f.capture.emit(RealtimeTalkAudioFrame(
data: Data([1, 2]),
timestampMs: ProcessInfo.processInfo.systemUptime * 1000 + 1000,
rms: 0.5))
try await f.requests.waitForRequestCount(3)
}
}
extension RealtimeTalkRelaySessionInterruptionTests {
@Test(arguments: [false, true])
func `production provider owned interruption keeps mic open until clear and explicit stop`(
suppressesInputDuringOutput: Bool) async throws
{
let f = try RealtimeRelayProductionFixture(supportsBargeIn: false)
f.capture.suppressesInputDuringOutput = suppressesInputDuringOutput
defer { f.close() }
try await f.start()
f.send(outputAudioEvent(turnId: "a", data: Data(repeating: 1, count: 960)))
try await f.device.waitForFrames(1)
f.capture.emit(RealtimeTalkAudioFrame(
data: Data([1, 2]),
timestampMs: ProcessInfo.processInfo.systemUptime * 1000 + 1000,
rms: 0.5))
try await f.requests.waitForRequestCount(3)
#expect(await f.requests.snapshot().last?.method == "talk.session.appendAudio")
#expect(!f.session.cancelOutput(reason: "barge-in"))
#expect(f.device.stopCount == 0)
f.send(outputClearEvent(turnId: "a", talkEventType: "output.clear"))
f.send(outputAudioEvent(turnId: "a", data: Data(repeating: 2, count: 960)))
try await f.device.waitForFrames(2)
#expect(f.session.cancelOutput(reason: "user"))
try await f.requests.waitForRequestCount(4)
#expect(await f.requests.snapshot().last?.method == "talk.session.cancelOutput")
}
}
extension RealtimeTalkRelaySessionReviewTests {
@Test func `blocked preparation cannot block synchronous relay controls or start retired audio`() async throws {
let gate = DispatchSemaphore(value: 0)
let f = try RealtimeRelayProductionFixture(device: RealtimeRelayDevice(prepareGate: gate))
defer { gate.signal()
f.close()
}
try await f.start()
f.send(outputAudioEvent(turnId: "a", data: Data(repeating: 1, count: 960)))
_ = try await f.device.prepared.next("prepare entered")
let started = ProcessInfo.processInfo.systemUptime
f.session.setOutputPaused(true)
#expect(ProcessInfo.processInfo.systemUptime - started < 0.25)
gate.signal()
_ = try await f.device.stopped.next("retired generation stopped")
#expect(f.device.frames.isEmpty)
}
@Test func `output levels restart after a mid-reply gap`() async throws {
let f = try RealtimeRelayProductionFixture()
defer { f.close() }
try await f.start()
let audio = Data(repeating: 0x20, count: 960 * 10)
f.send(outputAudioEvent(turnId: "a", data: audio))
try await f.device.waitForFrames(10)
var sawLevel = false
while let level = try await f.levels.next("first envelope") {
sawLevel = sawLevel || level > 0
}
#expect(sawLevel)
f.send(outputAudioEvent(turnId: "a", data: Data(repeating: 0x20, count: 960 * 10)))
try await f.device.waitForFrames(20)
var resumed = false
do {
while let level = try await f.levels.next("resumed envelope") {
if level > 0 { resumed = true
break
}
}
} catch {}
#expect(resumed, "same reply must restart level publishing after silence timeout")
}
@Test func `deinit cancels transport event subscription`() async throws {
let f = try RealtimeRelayProductionFixture()
try await f.start()
let ended = RealtimeRelayTestSignal<Void>(timeoutSeconds: 2)
f.events.continuation.onTermination = { _ in ended.send(()) }
weak var weakSession = f.session
f.session = nil
#expect(weakSession == nil)
var cancelled = false
do { _ = try await ended.next("subscription cancellation")
cancelled = true
} catch {}
#expect(cancelled)
f.events.continuation.finish()
}
}
#endif