mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-04 02:00:10 +00:00
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:
parent
7555e3f881
commit
646916d1c2
9 changed files with 1616 additions and 483 deletions
4
apps/.i18n/native-source.json
generated
4
apps/.i18n/native-source.json
generated
|
|
@ -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"}]},
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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")
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
@ -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)))
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
Loading…
Add table
Add a link
Reference in a new issue