diff --git a/apps/.i18n/native-source.json b/apps/.i18n/native-source.json index 23810ce0d2aa..d237bf739417 100644 --- a/apps/.i18n/native-source.json +++ b/apps/.i18n/native-source.json @@ -29897,6 +29897,17 @@ } ] }, + { + "id": "native.apple.024d68fa654acf2c", + "source": "Gateway connection was replaced before realtime startup finished", + "surface": "apple", + "sites": [ + { + "kind": "ui-localized-call", + "path": "apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift" + } + ] + }, { "id": "native.apple.165da8a03dd68ee1", "source": "Gateway data changed while you edited. Save will validate the current revision.", @@ -29919,6 +29930,17 @@ } ] }, + { + "id": "native.apple.da84239bbdff6838", + "source": "Gateway did not return a realtime relay session", + "surface": "apple", + "sites": [ + { + "kind": "ui-localized-call", + "path": "apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift" + } + ] + }, { "id": "native.apple.8104be3229e9088b", "source": "Gateway is not connected", @@ -29927,6 +29949,10 @@ { "kind": "ui-localized-call", "path": "apps/ios/Sources/Design/SettingsSystemAgentChat.swift" + }, + { + "kind": "ui-localized-call", + "path": "apps/ios/Sources/Voice/TalkModeManager.swift" } ] }, @@ -39326,6 +39352,83 @@ } ] }, + { + "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.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" + } + ] + }, { "id": "native.apple.afd0d1396586ab0b", "source": "Realtime disconnected", @@ -39337,6 +39440,17 @@ } ] }, + { + "id": "native.apple.ca3a762269827f37", + "source": "Realtime failed", + "surface": "apple", + "sites": [ + { + "kind": "ui-localized-call", + "path": "apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift" + } + ] + }, { "id": "native.apple.f350f10af66f4a87", "source": "Realtime failed before connecting", @@ -39348,6 +39462,28 @@ } ] }, + { + "id": "native.apple.96ece1b431ac3f96", + "source": "Realtime output cancellation failed: %@", + "surface": "apple", + "sites": [ + { + "kind": "ui-localized-call", + "path": "apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift" + } + ] + }, + { + "id": "native.apple.83f84908b64f015c", + "source": "Realtime tool call did not return a run id", + "surface": "apple", + "sites": [ + { + "kind": "ui-localized-call", + "path": "apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift" + } + ] + }, { "id": "native.apple.651a2b218e29ac62", "source": "Realtime unavailable", diff --git a/apps/ios/Sources/Voice/RealtimeTalkRelaySession.swift b/apps/ios/Sources/Voice/RealtimeTalkRelaySession.swift index 2c8148afb557..3fb47fa86278 100644 --- a/apps/ios/Sources/Voice/RealtimeTalkRelaySession.swift +++ b/apps/ios/Sources/Voice/RealtimeTalkRelaySession.swift @@ -36,7 +36,8 @@ final class IOSRealtimeTalkAudioCapture: RealtimeTalkAudioCapturing { func start( targetSampleRate: Double, - onAudio: @escaping @Sendable (RealtimeTalkAudioFrame) -> Void) throws + onAudio: @escaping @Sendable (RealtimeTalkAudioFrame) -> Void, + onFailure _: @escaping @MainActor (String) -> Void) throws { self.stop() let input = self.audioEngine.inputNode diff --git a/apps/ios/Sources/Voice/TalkModeGatewayConfig.swift b/apps/ios/Sources/Voice/TalkModeGatewayConfig.swift index ab0eedb25b6d..dd8ce5205846 100644 --- a/apps/ios/Sources/Voice/TalkModeGatewayConfig.swift +++ b/apps/ios/Sources/Voice/TalkModeGatewayConfig.swift @@ -9,6 +9,8 @@ enum TalkModeExecutionMode: Equatable { struct TalkRuntimeIssue: Equatable { enum Code: String { + case audioInputUnavailable = "audio_input_unavailable" + case realtimeOutputCancelFailed = "realtime_output_cancel_failed" case realtimeUnavailable = "realtime_unavailable" } diff --git a/apps/ios/Sources/Voice/TalkModeManager.swift b/apps/ios/Sources/Voice/TalkModeManager.swift index 4ea8c43b2396..e44e9e4cc04a 100644 --- a/apps/ios/Sources/Voice/TalkModeManager.swift +++ b/apps/ios/Sources/Voice/TalkModeManager.swift @@ -2179,7 +2179,8 @@ final class TalkModeManager: NSObject { return .started } guard let gateway else { - return .unavailable(realtimeIssue(message: "Gateway not connected", phase: "start")) + return .unavailable( + realtimeIssue(message: String(localized: "Gateway is not connected"), phase: "start")) } let startedAt = Self.nowSeconds() if self.prefetchedRealtimeSession == nil, let prefetchTask = realtimePrefetchTask { @@ -2310,7 +2311,8 @@ final class TalkModeManager: NSObject { private func startRealtimeRelayIfAvailable(attemptID: Int) async -> RealtimeStartResult { guard let gateway else { - return .unavailable(realtimeIssue(message: "Gateway not connected", phase: "start")) + return .unavailable( + realtimeIssue(message: String(localized: "Gateway is not connected"), phase: "start")) } guard self.foregroundAudioCaptureAllowed else { self.setStatus( @@ -2322,7 +2324,8 @@ final class TalkModeManager: NSObject { } guard self.isCurrentStartAttempt(attemptID) else { return .ignored } guard let gatewayRoute = await gateway.currentRoute() else { - return .unavailable(realtimeIssue(message: "Gateway not connected", phase: "start")) + return .unavailable( + realtimeIssue(message: String(localized: "Gateway is not connected"), phase: "start")) } guard self.isCurrentStartAttempt(attemptID) else { return .ignored } if self.realtimeRelaySession != nil { @@ -2353,26 +2356,24 @@ final class TalkModeManager: NSObject { model: self.realtimeModelId, voice: self.realtimeVoiceId), audioCapture: IOSRealtimeTalkAudioCapture(), - pcmPlayer: self.pcmPlayer, + pcmPlayer: RealtimePCMStreamingAudioPlayer(), onStatus: { [weak self] status in guard let self, self.realtimeRelayGeneration == relayGeneration else { return } self.handleRealtimeRelayStatus(status) }, onIssue: { [weak self] relayIssue in guard let self, self.realtimeRelayGeneration == relayGeneration else { return } - let issue = TalkRuntimeIssue( - code: .realtimeUnavailable, - message: relayIssue.message, - provider: relayIssue.provider, - model: relayIssue.model, - transport: relayIssue.transport, - phase: relayIssue.phase) + let issue = Self.runtimeIssue(from: relayIssue) self.realtimeRelayStartIssue = issue self.pendingRealtimeIssue = issue self.gatewayTalkLastIssueText = issue.diagnosticSummary self.gatewayTalkActiveModeTitle = String(localized: "Realtime unavailable") self.gatewayTalkActiveModeSubtitle = issue.displayMessage }, + onTermination: { [weak self] termination in + guard let self, self.realtimeRelayGeneration == relayGeneration else { return } + self.handleRealtimeRelayTermination(termination) + }, onSpeakingChanged: { [weak self] speaking in guard let self, self.realtimeRelayGeneration == relayGeneration else { return } self.isSpeaking = speaking @@ -4159,17 +4160,13 @@ extension TalkModeManager { let phase = Self.phase(forRealtimeStatus: status) if status == "Listening (Realtime)" { // Ready can be followed by a buffered close before start() resumes. Commit continuous - // state here so the close still enters bounded recovery. + // state here so the typed terminal callback still enters bounded recovery. self.markRealtimeSessionReady() } else { self.setStatus( Self.presentationText(forRealtimeStatus: status), phase: phase, watchPresentation: Self.watchPresentation(forRealtimeStatus: status)) - if status == "Ready" { - self.realtimeRelaySession = nil - self.handleRealtimeSessionFinish() - } } self.isListening = phase == .listening if phase == .thinking || phase == .connecting { @@ -4179,6 +4176,13 @@ extension TalkModeManager { } } + private func handleRealtimeRelayTermination(_ termination: RealtimeTalkRelayTermination) { + GatewayDiagnostics.log("talk realtime relay terminated reason=\(String(describing: termination))") + self.realtimeRelaySession = nil + guard self.captureMode != .pushToTalk else { return } + self.handleRealtimeSessionFinish() + } + private func prepareRealtimeRelayStart() { self.realtimeRelayStartIssue = nil self.pendingRealtimeIssue = nil @@ -4210,6 +4214,16 @@ extension TalkModeManager { phase: phase) } + static func runtimeIssue(from issue: RealtimeTalkRelayIssue) -> TalkRuntimeIssue { + TalkRuntimeIssue( + code: TalkRuntimeIssue.Code(rawValue: issue.code) ?? .realtimeUnavailable, + message: issue.message, + provider: issue.provider, + model: issue.model, + transport: issue.transport, + phase: issue.phase) + } + private func realtimeIssue(from error: Error, phase: String) -> TalkRuntimeIssue { if let gatewayError = error as? GatewayResponseError, let issue = Self.talkRuntimeIssue( @@ -5097,6 +5111,12 @@ extension TalkModeManager { self.handleRealtimeRelayStatus(status) } + func _test_handleRealtimeRelayTermination( + _ termination: RealtimeTalkRelayTermination = .remoteClose(reason: "completed")) + { + self.handleRealtimeRelayTermination(termination) + } + func _test_prepareEnabledRealtimeSessionForClose() { self.isEnabled = true self.gatewayConnected = true diff --git a/apps/ios/Tests/TalkModeConfigParsingTests.swift b/apps/ios/Tests/TalkModeConfigParsingTests.swift index e977122e9c8e..589f6b10a4b5 100644 --- a/apps/ios/Tests/TalkModeConfigParsingTests.swift +++ b/apps/ios/Tests/TalkModeConfigParsingTests.swift @@ -403,6 +403,16 @@ struct TalkModeManagerTests { #expect(issue.technicalDetails.contains("code: realtime_unavailable")) } + @Test func `relay issue preserves known code and falls back for unknown code`() { + let codes = ["audio_input_unavailable", "realtime_output_cancel_failed", "future_code"].map { rawCode in + TalkModeManager.runtimeIssue(from: RealtimeTalkRelayIssue( + code: rawCode, + message: "failed")).code + } + + #expect(codes == [.audioInputUnavailable, .realtimeOutputCancelFailed, .realtimeUnavailable]) + } + @Test func `native fallback keeps realtime issue visible`() { let manager = TalkModeManager(allowSimulatorCapture: true) let issue = TalkRuntimeIssue( @@ -480,6 +490,7 @@ struct TalkModeManagerTests { #expect(manager._test_gatewayTalkActiveModeTitle() != "Not active") manager._test_handleRealtimeRelayStatus("Ready") + manager._test_handleRealtimeRelayTermination() #expect(manager.statusText == "Ready") #expect(manager._test_gatewayTalkActiveModeTitle() == "Not active") @@ -524,6 +535,7 @@ struct TalkModeManagerTests { manager._test_handleRealtimeRelayStatus("Listening (Realtime)") manager._test_handleRealtimeRelayStatus("Ready") + manager._test_handleRealtimeRelayTermination() #expect(manager.statusText == "Reconnecting") #expect(manager._test_rapidRealtimeRestartCount() == 1) diff --git a/apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimePCMStreamingAudioPlayer.swift b/apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimePCMStreamingAudioPlayer.swift new file mode 100644 index 000000000000..68d3ac561da1 --- /dev/null +++ b/apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimePCMStreamingAudioPlayer.swift @@ -0,0 +1,237 @@ +#if Talk && canImport(ElevenLabsKit) && (os(iOS) || os(macOS)) +import AVFAudio +import Foundation + +@MainActor +public final class RealtimePCMStreamingAudioPlayer: PCMStreamingAudioPlaying { + static let frameDurationSeconds = 0.020 + static let maxScheduledBuffers = 3 + + typealias Completion = @Sendable () -> Void + private let preparePlayback: (Double) throws -> Void + private let scheduleFrame: (Data, Double, @escaping Completion) throws -> Void + private let stopPlayback: () -> Void + private let playbackTime: () -> Double? + + private var generation: UInt64 = 0 + private var nextBufferID: UInt64 = 0 + private var scheduledBufferIDs: Set = [] + private var slotWaiters: [CheckedContinuation] = [] + private var playbackContinuation: CheckedContinuation? + private var inputTask: Task? + private var inputFinished = false + + public convenience init() { + let engine = AVAudioEngine() + let node = AVAudioPlayerNode() + engine.attach(node) + var format: AVAudioFormat? + self.init( + preparePlayback: { sampleRate in + node.stop() + engine.stop() + engine.disconnectNodeOutput(node) + guard let nextFormat = AVAudioFormat( + commonFormat: .pcmFormatInt16, + sampleRate: sampleRate, + channels: 1, + interleaved: false) + else { + throw NSError(domain: "RealtimePCMStreamingAudioPlayer", code: 1) + } + format = nextFormat + engine.connect(node, to: engine.mainMixerNode, format: nextFormat) + engine.prepare() + try engine.start() + node.play() + }, + scheduleFrame: { data, _, completion in + guard let format else { + throw NSError(domain: "RealtimePCMStreamingAudioPlayer", code: 2) + } + let frames = AVAudioFrameCount(data.count / MemoryLayout.size) + guard let buffer = AVAudioPCMBuffer(pcmFormat: format, frameCapacity: frames), + let channel = buffer.int16ChannelData?[0] + else { + throw NSError(domain: "RealtimePCMStreamingAudioPlayer", code: 3) + } + buffer.frameLength = frames + data.copyBytes( + to: UnsafeMutableRawBufferPointer( + start: channel, + count: data.count)) + node.scheduleBuffer( + buffer, + completionCallbackType: .dataPlayedBack) + { _ in completion() } + }, + stopPlayback: { + node.stop() + engine.stop() + }, + playbackTime: { + guard let renderTime = node.lastRenderTime, + let playerTime = node.playerTime(forNodeTime: renderTime) + else { return nil } + return Double(playerTime.sampleTime) / playerTime.sampleRate + }) + } + + init( + preparePlayback: @escaping (Double) throws -> Void, + scheduleFrame: @escaping (Data, Double, @escaping Completion) throws -> Void, + stopPlayback: @escaping () -> Void, + playbackTime: @escaping () -> Double?) + { + self.preparePlayback = preparePlayback + self.scheduleFrame = scheduleFrame + self.stopPlayback = stopPlayback + self.playbackTime = playbackTime + } + + public func play( + stream: AsyncThrowingStream, + 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) + } + } + } + + public func stop() -> Double? { + let interruptedAt = self.playbackTime() + self.generation &+= 1 + self.inputTask?.cancel() + self.inputTask = nil + self.inputFinished = false + self.scheduledBufferIDs.removeAll() + let waiters = self.slotWaiters + self.slotWaiters.removeAll() + for waiter in waiters { + waiter.resume(returning: false) + } + let continuation = self.playbackContinuation + self.playbackContinuation = nil + self.stopPlayback() + continuation?.resume(returning: StreamingPlaybackResult( + finished: false, + interruptedAt: interruptedAt)) + return interruptedAt + } + + private func consume( + stream: AsyncThrowingStream, + sampleRate: Double, + generation: UInt64) async + { + let frameBytes = max( + MemoryLayout.size, + Int((sampleRate * Self.frameDurationSeconds).rounded()) * MemoryLayout.size) + var pending = Data() + do { + for try await chunk in stream { + try Task.checkCancellation() + pending.append(chunk) + while pending.count >= frameBytes { + let frame = Data(pending.prefix(frameBytes)) + pending.removeFirst(frameBytes) + 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) + else { return } + } + guard self.generation == generation else { return } + self.inputFinished = true + self.finishIfDrained(generation: generation) + } catch { + self.finish(generation: generation, finished: false) + } + } + + 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) + } + 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) + } + } + return true + } catch { + self.scheduledBufferIDs.remove(bufferID) + self.finish(generation: generation, finished: false) + return false + } + } + + 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) + } + self.finishIfDrained(generation: generation) + } + + private func finishIfDrained(generation: UInt64) { + guard self.inputFinished, self.scheduledBufferIDs.isEmpty else { return } + self.finish(generation: generation, finished: true) + } + + private func finish(generation: UInt64, finished: Bool) { + guard self.generation == generation else { return } + let interruptedAt = finished ? nil : self.playbackTime() + self.generation &+= 1 + self.scheduledBufferIDs.removeAll() + let waiters = self.slotWaiters + self.slotWaiters.removeAll() + for waiter in waiters { + waiter.resume(returning: false) + } + self.inputTask = nil + self.inputFinished = false + let continuation = self.playbackContinuation + self.playbackContinuation = nil + self.stopPlayback() + continuation?.resume(returning: StreamingPlaybackResult( + finished: finished, + interruptedAt: interruptedAt)) + } +} +#endif diff --git a/apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift b/apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift index 323af502311c..e29921037093 100644 --- a/apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift +++ b/apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift @@ -55,7 +55,8 @@ public protocol RealtimeTalkAudioCapturing: AnyObject { func start( targetSampleRate: Double, - onAudio: @escaping @Sendable (RealtimeTalkAudioFrame) -> Void) throws + onAudio: @escaping @Sendable (RealtimeTalkAudioFrame) -> Void, + onFailure: @escaping @MainActor (String) -> Void) throws func stop() } @@ -113,6 +114,18 @@ public struct RealtimeTalkTranscript: Equatable, Sendable { } } +public enum RealtimeTalkRelayTermination: Equatable, Sendable { + case remoteClose(reason: String?) + case eventStreamEnded + case audioInputFailed(message: String) + case outputCancellationFailed + case outputPlaybackOverflow +} + +private enum RealtimeAudioSendOutcome { + case sent, inactive, saturated, failed(String) +} + private actor RealtimeAudioSender { private let request: @Sendable (String, [String: AnyCodable]?, Double) async throws -> Data private var relaySessionId: String? @@ -131,25 +144,27 @@ private actor RealtimeAudioSender { self.relaySessionId = nil } - func send(_ data: Data, timestampMs: Double) async -> String? { - guard !Task.isCancelled else { return nil } - guard let relaySessionId else { return nil } - guard self.pendingSends < self.maxPendingSends else { return nil } + func send(_ data: Data, timestampMs: Double) async -> RealtimeAudioSendOutcome { + guard !Task.isCancelled, let relaySessionId else { return .inactive } + guard self.pendingSends < self.maxPendingSends else { return .saturated } self.pendingSends += 1 defer { self.pendingSends -= 1 } + // The Gateway carries this straight into the provider's media timeline, and OpenAI rejects + // a `conversation.item.truncate` whose `audio_end_ms` is not an integer -- a fractional + // timestamp here kills the session on the first barge-in. let payload: [String: AnyCodable] = [ "sessionId": AnyCodable(relaySessionId), "audioBase64": AnyCodable(data.base64EncodedString()), - "timestamp": AnyCodable(timestampMs), + "timestamp": AnyCodable(timestampMs.rounded()), ] do { try Task.checkCancellation() let response = try await self.request("talk.session.appendAudio", payload, 8000) try Task.checkCancellation() _ = try JSONDecoder().decode(TalkSessionOkResult.self, from: response) - return nil + return .sent } catch { - return error.localizedDescription + return Task.isCancelled ? .inactive : .failed(error.localizedDescription) } } } @@ -194,6 +209,15 @@ public final class RealtimeTalkRelaySession { case cancelled } + /// Startup abandons for two very different reasons and callers must not treat them alike. + /// Local cancellation is caller-initiated and stays silent; a lost Gateway route is an + /// external failure the runtime has to see, or Talk reports listening with no relay behind it. + private enum LifecycleStatus { + case current + case cancelledLocally + case routeLost + } + private nonisolated static let expectedInputEncoding = "pcm16" private nonisolated static let expectedOutputEncoding = "pcm16" private nonisolated static let defaultSampleRateHz = 24000 @@ -201,6 +225,9 @@ 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 private let transport: RealtimeTalkRelayTransport private let audioCapture: any RealtimeTalkAudioCapturing @@ -209,6 +236,7 @@ public final class RealtimeTalkRelaySession { private let logger = Logger(subsystem: "ai.openclawfoundation.app", category: "RealtimeTalkRelay") private let onStatus: (String) -> Void private let onIssue: (RealtimeTalkRelayIssue) -> Void + private let onTermination: (RealtimeTalkRelayTermination) -> Void private let onSpeakingChanged: (Bool) -> Void private let onInputLevel: (Double) -> Void private let onOutputLevel: (Double?) -> Void @@ -230,17 +258,25 @@ public final class RealtimeTalkRelaySession { private var audioSendTasks: [UUID: Task] = [:] private var outputTask: Task? private var outputContinuation: AsyncThrowingStream.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 outputIdleTask: Task? private var outputSessionId = 0 - private var pendingOutputChunks: [Data] = [] - private var pendingOutputDone = false 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 outputCancellationTask: Task? private var outputStartedAtMs: Double? private var outputPlaybackExpectedEndMs: Double = 0 private var lastBargeInAtMs: Double = 0 @@ -262,6 +298,7 @@ public final class RealtimeTalkRelaySession { pcmPlayer: PCMStreamingAudioPlaying, onStatus: @escaping (String) -> Void, onIssue: @escaping (RealtimeTalkRelayIssue) -> Void = { _ in }, + onTermination: @escaping (RealtimeTalkRelayTermination) -> Void = { _ in }, onSpeakingChanged: @escaping (Bool) -> Void, onInputLevel: @escaping (Double) -> Void = { _ in }, onOutputLevel: @escaping (Double?) -> Void = { _ in }, @@ -273,6 +310,7 @@ public final class RealtimeTalkRelaySession { self.pcmPlayer = pcmPlayer self.onStatus = onStatus self.onIssue = onIssue + self.onTermination = onTermination self.onSpeakingChanged = onSpeakingChanged self.onInputLevel = onInputLevel self.onOutputLevel = onOutputLevel @@ -290,11 +328,16 @@ public final class RealtimeTalkRelaySession { self.pendingPreRelayEvents.removeAll() self.onStatus("Connecting realtime…") let eventStream = await self.transport.subscribeServerEvents(200) - guard await self.isCurrentLifecycle(lifecycleGeneration) else { return } + switch await self.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() - guard await self.isCurrentLifecycle(lifecycleGeneration) else { + let statusAfterCreate = await self.lifecycleStatus(lifecycleGeneration) + if statusAfterCreate != .current { if let relaySessionId = result.relaysessionid?.trimmingCharacters(in: .whitespacesAndNewlines), !relaySessionId.isEmpty { @@ -302,13 +345,27 @@ public final class RealtimeTalkRelaySession { transport: self.transport, relaySessionId: relaySessionId) } + if statusAfterCreate == .routeLost { + throw Self.gatewayRouteLostError() + } return } + if let startupIssue { + if let relaySessionId = result.relaysessionid?.trimmingCharacters(in: .whitespacesAndNewlines), + !relaySessionId.isEmpty + { + await Self.closeRelaySession( + transport: self.transport, + relaySessionId: relaySessionId) + } + throw Self.startupFailureError(startupIssue) + } guard let relaySessionId = result.relaysessionid?.trimmingCharacters(in: .whitespacesAndNewlines), !relaySessionId.isEmpty else { throw NSError(domain: "RealtimeTalkRelay", code: 1, userInfo: [ - NSLocalizedDescriptionKey: "Gateway did not return a realtime relay session", + NSLocalizedDescriptionKey: String( + localized: "Gateway did not return a realtime relay session"), ]) } self.relaySessionId = relaySessionId @@ -319,7 +376,11 @@ public final class RealtimeTalkRelaySession { try self.startMicrophonePump(lifecycleGeneration: lifecycleGeneration) self.onStatus("Waiting for realtime…") await self.drainPendingPreRelayEvents(lifecycleGeneration: lifecycleGeneration) - guard await self.isCurrentLifecycle(lifecycleGeneration) else { return } + switch await self.lifecycleStatus(lifecycleGeneration) { + case .current: break + case .cancelledLocally: return + case .routeLost: throw Self.gatewayRouteLostError() + } switch await self.waitForStartupResult( timeoutSeconds: Self.startupReadyTimeoutSeconds, lifecycleGeneration: lifecycleGeneration) @@ -327,7 +388,6 @@ public final class RealtimeTalkRelaySession { case .ready: return case let .failed(issue): - self.close(sendClose: true) throw NSError(domain: "RealtimeTalkRelay", code: 6, userInfo: [ NSLocalizedDescriptionKey: issue.message, ]) @@ -335,7 +395,9 @@ public final class RealtimeTalkRelaySession { return } } catch { - guard await self.isCurrentLifecycle(lifecycleGeneration) else { return } + // 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 } let createdRelaySessionId = self.relaySessionId self.close(sendClose: false) if let createdRelaySessionId { @@ -362,14 +424,13 @@ public final class RealtimeTalkRelaySession { for task in self.toolCallTasks.values { task.cancel() } - for task in self.audioSendTasks.values { - task.cancel() - } - self.audioSendTasks.removeAll() self.pendingPlaybackMarks.removeAll() let audioSender = self.audioSender self.audioSender = nil Task { await audioSender?.close() } + self.retireOutputCancellation() + self.cancelledOutputTurnId = nil + self.isOutputPaused = false self.stopOutputPlayback() if sendClose, let relaySessionId = self.relaySessionId { Task { [transport] in @@ -380,6 +441,21 @@ public final class RealtimeTalkRelaySession { self.onSpeakingChanged(false) } + /// Deliberately not a `CancellationError`: the runtime treats those as caller-initiated and + /// returns silently, while any other error routes Talk to its native fallback. + private nonisolated static func gatewayRouteLostError() -> NSError { + NSError(domain: "RealtimeTalkRelay", code: 7, userInfo: [ + NSLocalizedDescriptionKey: String( + localized: "Gateway connection was replaced before realtime startup finished"), + ]) + } + + private nonisolated static func startupFailureError(_ issue: RealtimeTalkRelayIssue) -> NSError { + NSError(domain: "RealtimeTalkRelay", code: 6, userInfo: [ + NSLocalizedDescriptionKey: issue.message, + ]) + } + private nonisolated static func closeRelaySession( transport: RealtimeTalkRelayTransport, relaySessionId: String) async @@ -388,18 +464,6 @@ public final class RealtimeTalkRelaySession { _ = try? await transport.request("talk.session.close", payload, 8000) } - public func cancelOutput(reason: String = "user") { - self.stopOutputPlayback() - guard let relaySessionId else { return } - Task { [transport] in - let payload: [String: AnyCodable] = [ - "sessionId": AnyCodable(relaySessionId), - "reason": AnyCodable(reason), - ] - _ = try? await transport.request("talk.session.cancelOutput", payload, 8000) - } - } - public func setInputPaused(_ paused: Bool) throws { guard self.isInputPaused != paused else { return } self.isInputPaused = paused @@ -416,6 +480,14 @@ 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") + } + } + private func createRelaySession() async throws -> TalkSessionCreateResult { var payload: [String: AnyCodable] = [ "sessionKey": AnyCodable(self.options.sessionKey), @@ -457,8 +529,35 @@ public final class RealtimeTalkRelaySession { if Task.isCancelled { return } await self?.handleGatewayEvent(event, lifecycleGeneration: lifecycleGeneration) } + guard !Task.isCancelled else { return } + await self?.handleEventStreamEnded(lifecycleGeneration: lifecycleGeneration) } } +} + +extension RealtimeTalkRelaySession { + private func handleEventStreamEnded(lifecycleGeneration: UInt64) async { + guard self.isCurrentLifecycleLocally(lifecycleGeneration) else { return } + self.logger.debug("talk realtime: event stream ended") + guard self.hasReceivedReady else { + guard !self.hasReceivedFailure else { return } + let issue = RealtimeTalkRelayIssue( + message: String(localized: "Realtime connection ended before it became ready."), + provider: self.options.provider, + model: self.options.model, + transport: "gateway-relay", + phase: "connect") + self.hasReceivedFailure = true + self.startupIssue = issue + self.onIssue(issue) + self.onStatus(issue.message) + self.finishStartupWait(.failed(issue)) + return + } + self.onStatus("Ready") + self.close(sendClose: false) + self.onTermination(.eventStreamEnded) + } private func handleGatewayEvent(_ event: EventFrame, lifecycleGeneration: UInt64) async { guard self.isCurrentLifecycleLocally(lifecycleGeneration) else { return } @@ -482,25 +581,11 @@ public final class RealtimeTalkRelaySession { self.finishStartupWait(.ready) self.onStatus("Listening (Realtime)") case "audio": - guard let base64 = payload["audioBase64"]?.stringValue, - let data = Data(base64Encoded: base64) - else { return } - self.recordOutputAudioChunk(byteCount: data.count) - self.markOutputAudioStarted(byteCount: data.count, nowMs: ProcessInfo.processInfo.systemUptime * 1000) - self.onSpeakingChanged(true) - if self.outputContinuation == nil, self.outputTask != nil { - self.pendingOutputChunks.append(data) - return - } - self.ensureOutputPlaybackStarted() - self.outputEnvelope?.append(data) - self.outputContinuation?.yield(data) + self.handleOutputAudio(payload) case "audioDone": - self.finishOutputPlaybackStream() + self.handleOutputAudioDone(payload) case "clear": - let marks = self.takePendingPlaybackMarks() - self.stopOutputPlayback() - self.acknowledgePlaybackMarks(marks) + self.handleOutputClear(payload) case "mark": self.handlePlaybackMark(payload) case "transcript": @@ -508,7 +593,7 @@ public final class RealtimeTalkRelaySession { case "toolCall": self.startToolCall(payload, lifecycleGeneration: lifecycleGeneration) case "error": - let message = payload["message"]?.stringValue ?? "Realtime failed" + let message = payload["message"]?.stringValue ?? String(localized: "Realtime failed") let issue = Self.issue( payload: payload, fallbackMessage: message, @@ -524,9 +609,13 @@ public final class RealtimeTalkRelaySession { self.logger.debug("talk realtime: close") if self.hasReceivedReady { self.onStatus("Ready") + let reason = self.nonEmpty(payload["reason"]?.stringValue) + self.close(sendClose: false) + self.onTermination(.remoteClose(reason: reason)) + return } else if !self.hasReceivedFailure { let issue = RealtimeTalkRelayIssue( - message: "Realtime closed before it became ready.", + message: String(localized: "Realtime closed before it became ready."), provider: self.options.provider, model: self.options.model, transport: "gateway-relay", @@ -536,12 +625,39 @@ public final class RealtimeTalkRelaySession { self.finishStartupWait(.failed(issue)) self.onStatus("Realtime failed before connecting") } - self.close(sendClose: false) default: return } } + private func handleOutputClear(_ payload: [String: AnyCodable]) { + let clearIdentity = OutputIdentity(payload) + if self.awaitingOutputClear, + let suppressed = self.suppressedOutputIdentity + { + let clearsSuppressed = + clearIdentity.isEmpty() + ? suppressed.isEmpty() + : suppressed.isEmpty() || suppressed.relation(to: clearIdentity) == .same + if clearsSuppressed { + self.awaitingOutputClear = false + if self.outputCancellationTask == nil { self.retireOutputCancellation() } + } + } + let currentMatches = + clearIdentity.isEmpty() + ? self.outputIdentity == nil + : self.outputIdentity?.relation(to: clearIdentity) == .same + guard 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 @@ -587,7 +703,7 @@ public final class RealtimeTalkRelaySession { return } let issue = RealtimeTalkRelayIssue( - message: "Realtime did not become ready in time.", + message: String(localized: "Realtime did not become ready in time."), provider: self.options.provider, model: self.options.model, transport: "gateway-relay", @@ -698,7 +814,8 @@ public final class RealtimeTalkRelaySession { lifecycleGeneration: lifecycleGeneration) guard let runId = startResponse.runId ?? startResponse.idempotencyKey else { throw NSError(domain: "RealtimeTalkRelay", code: 3, userInfo: [ - NSLocalizedDescriptionKey: "Realtime tool call did not return a run id", + NSLocalizedDescriptionKey: String( + localized: "Realtime tool call did not return a run id"), ]) } let completion = await self.waitForChatCompletion( @@ -844,11 +961,15 @@ public final class RealtimeTalkRelaySession { return try JSONDecoder().decode(type, from: response) } - private func isCurrentLifecycle(_ lifecycleGeneration: UInt64) async -> Bool { - guard self.isCurrentLifecycleLocally(lifecycleGeneration) else { return false } + private func lifecycleStatus(_ lifecycleGeneration: UInt64) async -> LifecycleStatus { + guard self.isCurrentLifecycleLocally(lifecycleGeneration) else { return .cancelledLocally } let routeIsCurrent = await self.transport.isCurrent() - return self.isCurrentLifecycleLocally(lifecycleGeneration) && - routeIsCurrent + guard self.isCurrentLifecycleLocally(lifecycleGeneration) else { return .cancelledLocally } + return routeIsCurrent ? .current : .routeLost + } + + private func isCurrentLifecycle(_ lifecycleGeneration: UInt64) async -> Bool { + await self.lifecycleStatus(lifecycleGeneration) == .current } private func isCurrentLifecycleLocally(_ lifecycleGeneration: UInt64) -> Bool { @@ -869,57 +990,41 @@ public final class RealtimeTalkRelaySession { } envelope.begin(sampleRate: self.outputSampleRateHz) self.outputEnvelope = envelope - let stream = AsyncThrowingStream { continuation in - self.outputContinuation = continuation - } + let stream = AsyncThrowingStream( + 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, let interruptedAt = result.interruptedAt { - self.logger.info("realtime output interrupted at \(interruptedAt, privacy: .public)s") + 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() - self.startPendingOutputPlaybackIfNeeded() } } } private func finishOutputPlaybackStream() { - guard let continuation = self.outputContinuation else { - if self.outputTask != nil, !self.pendingOutputChunks.isEmpty { - self.pendingOutputDone = true - } - return + 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 startPendingOutputPlaybackIfNeeded() { - guard !self.pendingOutputChunks.isEmpty else { - self.pendingOutputDone = false - return - } - let chunks = self.pendingOutputChunks - let shouldFinish = self.pendingOutputDone - self.pendingOutputChunks = [] - self.pendingOutputDone = false - self.ensureOutputPlaybackStarted() - for chunk in chunks { - self.markOutputAudioStarted(byteCount: chunk.count, nowMs: ProcessInfo.processInfo.systemUptime * 1000) - self.onSpeakingChanged(true) - self.outputEnvelope?.append(chunk) - self.outputContinuation?.yield(chunk) - } - if shouldFinish { - self.finishOutputPlaybackStream() - } - } - private func scheduleOutputPlaybackIdle(expectedEndMs: Double) { self.outputIdleTask?.cancel() let nowMs = ProcessInfo.processInfo.systemUptime * 1000 @@ -947,6 +1052,7 @@ public final class RealtimeTalkRelaySession { self.outputIdleTask = nil } self.isOutputPlaying = false + self.outputIdentity = nil self.outputStartedAtMs = nil self.outputPlaybackExpectedEndMs = 0 self.outputEnvelope?.cancel() @@ -984,8 +1090,9 @@ public final class RealtimeTalkRelaySession { do { _ = try await transport.request("talk.session.acknowledgeMark", payload, 8000) } catch { + let message = Self.safeLogMessage(error.localizedDescription) logger.warning( - "talk realtime: mark acknowledgement failed=\(Self.safeLogMessage(error.localizedDescription), privacy: .public)") + "talk realtime: mark acknowledgement failed=\(message, privacy: .public)") } } } @@ -995,14 +1102,14 @@ public final class RealtimeTalkRelaySession { self.outputSessionId += 1 self.outputContinuation?.finish() self.outputContinuation = nil + self.pendingOutputAudio.removeAll(keepingCapacity: true) self.outputTask?.cancel() self.outputTask = nil self.outputIdleTask?.cancel() self.outputIdleTask = nil - self.pendingOutputChunks = [] - self.pendingOutputDone = false _ = self.pcmPlayer.stop() self.isOutputPlaying = false + self.outputIdentity = nil self.outputStartedAtMs = nil self.outputPlaybackExpectedEndMs = 0 self.outputEnvelope?.cancel() @@ -1050,22 +1157,260 @@ public final class RealtimeTalkRelaySession { } } +extension RealtimeTalkRelaySession { + private struct OutputIdentity { + enum Relation { + case same + case different + case unknown + } + + let turnId: String? + init(_ payload: [String: AnyCodable]) { + let turnId = payload["talkEvent"]?.dictionaryValue?["turnId"]?.stringValue? + .trimmingCharacters(in: .whitespacesAndNewlines) + self.turnId = turnId?.isEmpty == false ? turnId : nil + } + + func isEmpty() -> Bool { + self.turnId == nil + } + + func relation(to other: OutputIdentity) -> Relation { + if self.isEmpty() || other.isEmpty() { + return self.isEmpty() == other.isEmpty() ? .same : .different + } + if let turnId, let otherTurnId = other.turnId { + return turnId == otherTurnId ? .same : .different + } + return .unknown + } + } + + @discardableResult + public func cancelOutput(reason: String = "user") -> Bool { + guard let relaySessionId, + let outputIdentity = self.outputIdentity, + let turnId = outputIdentity.turnId + else { return false } + self.outputCancellationGeneration &+= 1 + let cancellationGeneration = self.outputCancellationGeneration + self.outputCancellationTask?.cancel() + 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), + "reason": AnyCodable(reason), + "turnId": AnyCodable(turnId), + ] + do { + let response = try await transport.request("talk.session.cancelOutput", payload, 8000) + let result = try JSONDecoder().decode(TalkSessionCancelOutputResult.self, from: response) + guard result.ok else { throw URLError(.badServerResponse) } + guard let self, self.isCurrentOutputCancellation(cancellationGeneration) else { return } + switch result.status?.stringValue { + case "stale", "idle": + self.retireOutputCancellation() + case nil, "applied": + guard result.turnid == nil || result.turnid == turnId else { + throw URLError(.badServerResponse) + } + if self.awaitingOutputClear { + self.outputCancellationTask = nil + } else { + self.retireOutputCancellation() + } + default: + throw URLError(.badServerResponse) + } + } catch { + guard let self, self.isCurrentOutputCancellation(cancellationGeneration) else { return } + let issue = RealtimeTalkRelayIssue( + code: "realtime_output_cancel_failed", + message: String( + format: String(localized: "Realtime output cancellation failed: %@"), + error.localizedDescription), + provider: self.options.provider, + model: self.options.model, + transport: "gateway-relay", + phase: "output-cancel") + self.onIssue(issue) + self.onStatus(issue.message) + // A failed current cancellation leaves remote output ownership unknown. + // Keep the fence until terminal teardown makes late audio impossible. + self.close(sendClose: true) + self.onTermination(.outputCancellationFailed) + } + } + return true + } + + private func isCurrentOutputCancellation(_ generation: UInt64) -> Bool { + 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 } + if let cancelledOutputTurnId { + if incomingTurnId == cancelledOutputTurnId { + return + } + } + guard let base64 = payload["audioBase64"]?.stringValue else { return } + guard let data = Data(base64Encoded: base64) else { + self.handleOutputPlaybackOverflow() + return + } + if let currentIdentity = self.outputIdentity, + currentIdentity.relation(to: incomingIdentity) == .different + { + self.stopOutputPlayback() + } else if self.outputContinuation == nil, self.outputTask != nil { + self.stopOutputPlayback() + } + self.outputIdentity = incomingIdentity + self.recordOutputAudioChunk(byteCount: data.count) + self.markOutputAudioStarted(byteCount: data.count, 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..= frameByteCount { + let frameEnd = data.index(offset, offsetBy: frameByteCount) + let frame = Data(data[offset.. 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.isEmpty(), + let outputIdentity, + outputIdentity.relation(to: incomingIdentity) != .same + { + 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 = RealtimeTalkRelayIssue( + message: message, + provider: self.options.provider, + model: self.options.model, + transport: "gateway-relay", + phase: "output-playback") + self.onIssue(issue) + self.onStatus(message) + self.close(sendClose: true) + self.onTermination(.outputPlaybackOverflow) + } + + private func retireOutputCancellation() { + self.outputCancellationGeneration &+= 1 + self.outputCancellationTask?.cancel() + self.outputCancellationTask = nil + self.suppressedOutputIdentity = nil + self.awaitingOutputClear = false + } +} + extension RealtimeTalkRelaySession { private func startMicrophonePump(lifecycleGeneration: UInt64) throws { self.stopMicrophonePump() guard !self.isInputPaused else { return } - self.audioCaptureGeneration &+= 1 let audioCaptureGeneration = self.audioCaptureGeneration - try self.audioCapture.start(targetSampleRate: self.inputSampleRateHz) { [weak self] frame in - Task { @MainActor [weak self] in - _ = self?.enqueueMicrophoneFrame( - frame.data, - timestampMs: frame.timestampMs, - rms: frame.rms, + try self.audioCapture.start( + targetSampleRate: self.inputSampleRateHz, + onAudio: { [weak self] frame in + Task { @MainActor [weak self] in + _ = self?.enqueueMicrophoneFrame( + frame.data, + timestampMs: frame.timestampMs, + rms: frame.rms, + lifecycleGeneration: lifecycleGeneration, + audioCaptureGeneration: audioCaptureGeneration) + } + }, + onFailure: { [weak self] message in + self?.handleAudioInputFailure( + message, lifecycleGeneration: lifecycleGeneration, audioCaptureGeneration: audioCaptureGeneration) - } - } + }) + } + + private func handleAudioInputFailure( + _ message: String, + lifecycleGeneration: UInt64, + audioCaptureGeneration: UInt64) + { + guard !self.isClosed, self.isCurrentLifecycleLocally(lifecycleGeneration), + self.audioCaptureGeneration == audioCaptureGeneration + else { return } + let issue = RealtimeTalkRelayIssue( + code: "audio_input_unavailable", + message: message, + provider: self.options.provider, + model: self.options.model, + transport: "gateway-relay", + phase: "audio-input") + self.logger.error("talk realtime microphone failed: \(Self.safeLogMessage(message), privacy: .public)") + self.onIssue(issue) + self.onStatus(message) + self.close(sendClose: true) + self.onTermination(.audioInputFailed(message: message)) } @discardableResult @@ -1078,7 +1423,7 @@ extension RealtimeTalkRelaySession { { guard self.isCurrentLifecycleLocally(lifecycleGeneration), self.audioCaptureGeneration == audioCaptureGeneration, - !self.isInputPaused, + !self.isInputPaused, self.suppressedOutputIdentity == nil, let audioSender = self.audioSender else { return nil } self.recordMicrophoneFrame(byteCount: encoded.count, rms: rms, timestampMs: timestampMs) @@ -1100,10 +1445,24 @@ extension RealtimeTalkRelaySession { let task = Task { @MainActor [weak self, audioSender] in guard let self else { return } defer { self.audioSendTasks.removeValue(forKey: taskID) } - guard self.isCurrentLifecycleLocally(lifecycleGeneration) else { return } - guard let message = await audioSender.send(encoded, timestampMs: timestampMs) else { return } - guard self.isCurrentLifecycleLocally(lifecycleGeneration) else { return } - self.onStatus("Realtime audio failed: \(message)") + guard self.isCurrentLifecycleLocally(lifecycleGeneration), + self.audioCaptureGeneration == audioCaptureGeneration, + !self.isInputPaused, self.suppressedOutputIdentity == nil + else { return } + switch await audioSender.send(encoded, timestampMs: timestampMs) { + case .sent, .inactive: + return + case .saturated: + self.handleAudioInputFailure( + String(localized: "Realtime audio input fell behind. Reconnecting…"), + lifecycleGeneration: lifecycleGeneration, + audioCaptureGeneration: audioCaptureGeneration) + case let .failed(message): + self.handleAudioInputFailure( + String(format: String(localized: "Realtime audio failed: %@"), message), + lifecycleGeneration: lifecycleGeneration, + audioCaptureGeneration: audioCaptureGeneration) + } } self.audioSendTasks[taskID] = task return task @@ -1132,8 +1491,10 @@ extension RealtimeTalkRelaySession { guard timestampMs - self.lastSuppressedEchoLogAtMs >= 1000 else { return } self.lastSuppressedEchoLogAtMs = timestampMs let maxRms = String(format: "%.4f", Double(self.suppressedEchoMaxRms)) + let frames = self.suppressedEchoFrameCount + let bytes = self.suppressedEchoByteCount self.logger.debug( - "talk realtime mic suppressed during output: buffers=\(self.suppressedEchoFrameCount) bytes=\(self.suppressedEchoByteCount) maxRms=\(maxRms)") + "talk realtime mic suppressed during output: buffers=\(frames) bytes=\(bytes) maxRms=\(maxRms)") self.suppressedEchoFrameCount = 0 self.suppressedEchoByteCount = 0 self.suppressedEchoMaxRms = 0 @@ -1141,19 +1502,32 @@ extension RealtimeTalkRelaySession { private func stopMicrophonePump() { self.audioCaptureGeneration &+= 1 + for task in self.audioSendTasks.values { + task.cancel() + } + self.audioSendTasks.removeAll() self.audioCapture.stop() } } +#if DEBUG extension RealtimeTalkRelaySession { + // periphery:ignore - package tests drive a relay session without a live gateway handshake. func _test_setRelaySessionId(_ relaySessionId: String) { self.relaySessionId = relaySessionId } + // periphery:ignore - package tests inject gateway events without a live socket. func _test_handleGatewayEvent(_ event: EventFrame) async { await self.handleGatewayEvent(event, lifecycleGeneration: self.lifecycleGeneration) } + // periphery:ignore - package tests end the event stream deterministically. + func _test_handleEventStreamEnded() async { + await self.handleEventStreamEnded(lifecycleGeneration: self.lifecycleGeneration) + } + + // periphery:ignore - package tests observe startup cancellation without waiting out the timeout. func _test_waitForStartupCancelled(timeoutSeconds: Int) async -> Bool { if case .cancelled = await self.waitForStartupResult( timeoutSeconds: timeoutSeconds, @@ -1164,6 +1538,7 @@ extension RealtimeTalkRelaySession { return false } + // periphery:ignore - package tests await in-flight tool calls before asserting. func _test_waitForToolCalls() async { let tasks = self.toolCallTasks.values for task in tasks { @@ -1171,26 +1546,32 @@ extension RealtimeTalkRelaySession { } } - func _test_startupReadyTimeoutSeconds() -> Int { - Self.startupReadyTimeoutSeconds + // periphery:ignore - package tests capture the exact owned cancellation before replacement or stop. + func _test_outputCancellationTask() -> Task? { + self.outputCancellationTask } + // periphery:ignore - package tests start output playback without decoding real audio. func _test_markOutputAudioStarted(nowMs: Double) { self.markOutputAudioStarted(byteCount: 4800, nowMs: nowMs) } + // periphery:ignore - package tests finish playback without a real player callback. func _test_markOutputPlaybackFinished() { self.markOutputPlaybackFinished() } + // periphery:ignore - package tests observe barge-in timing state. func _test_outputStartedAtMs() -> Double? { self.outputStartedAtMs } + // periphery:ignore - package tests observe playback state without exposing it publicly. func _test_isOutputPlaying() -> Bool { self.isOutputPlaying } + // periphery:ignore - package tests exercise the audio sender without a started session. func _test_prepareAudioSender(relaySessionId: String) { self.isClosed = false self.audioSender = RealtimeAudioSender( @@ -1198,13 +1579,23 @@ extension RealtimeTalkRelaySession { request: self.transport.request) } - func _test_enqueueMicrophoneFrame(_ data: Data) -> Task? { + // periphery:ignore - package tests enqueue frames without a live capture device. + func _test_enqueueMicrophoneFrame( + _ data: Data, + timestampMs: Double = 1) -> Task? + { self.enqueueMicrophoneFrame( data, - timestampMs: 1, + timestampMs: timestampMs, rms: 0.01, lifecycleGeneration: self.lifecycleGeneration, audioCaptureGeneration: self.audioCaptureGeneration) } + + // periphery:ignore - package tests start the pump to observe capture failure handling. + func _test_startMicrophonePump() throws { + try self.startMicrophonePump(lifecycleGeneration: self.lifecycleGeneration) + } } #endif +#endif diff --git a/apps/shared/OpenClawKit/Tests/OpenClawKitTests/RealtimePCMStreamingAudioPlayerTests.swift b/apps/shared/OpenClawKit/Tests/OpenClawKitTests/RealtimePCMStreamingAudioPlayerTests.swift new file mode 100644 index 000000000000..8d172d271901 --- /dev/null +++ b/apps/shared/OpenClawKit/Tests/OpenClawKitTests/RealtimePCMStreamingAudioPlayerTests.swift @@ -0,0 +1,340 @@ +#if Talk && canImport(ElevenLabsKit) && (os(iOS) || os(macOS)) +import Foundation +import Testing +@testable import OpenClawKit + +private struct RealtimePCMPlaybackWaitTimeout: Error { + let label: String +} + +private struct RealtimePCMPlaybackFailure: Error {} + +private let realtimePCMPlaybackWaitTimeoutSeconds = 15.0 + +@MainActor +private final class RealtimePCMPlaybackBackend { + private struct Waiter { + let count: Int + let continuation: CheckedContinuation + } + + private(set) var scheduledFrames: [Data] = [] + private(set) var completions: [@Sendable () -> Void] = [] + private(set) var activeCount = 0 + private(set) var maxActiveCount = 0 + private var completedCallbacks = 0 + private var scheduledWaiters: [UUID: Waiter] = [:] + private var completionWaiters: [UUID: Waiter] = [:] + + func prepare(sampleRate _: Double) throws {} + + 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() + completion() + } + } + } + + func stop() { + self.activeCount = 0 + } + + func complete(at index: Int = 0) { + self.completions.remove(at: index)() + } + + func takeCompletion(at index: Int = 0) -> @Sendable () -> Void { + self.completions.remove(at: index) + } + + func waitForScheduledFrames(_ count: Int) async throws { + if self.scheduledFrames.count >= count { return } + try await AsyncTimeout.withTimeout( + seconds: realtimePCMPlaybackWaitTimeoutSeconds, + onTimeout: { RealtimePCMPlaybackWaitTimeout(label: "scheduled frames \(count)") }, + operation: { try await self.waitForScheduledFramesWithoutDeadline(count) }) + } + + func waitForCompletionCallbacks(_ count: Int) async throws { + if self.completedCallbacks >= count { return } + try await AsyncTimeout.withTimeout( + seconds: realtimePCMPlaybackWaitTimeoutSeconds, + onTimeout: { RealtimePCMPlaybackWaitTimeout(label: "completion callbacks \(count)") }, + operation: { try await self.waitForCompletionCallbacksWithoutDeadline(count) }) + } + + private func waitForScheduledFramesWithoutDeadline(_ count: Int) async throws { + 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) + } + } + } onCancel: { + Task { @MainActor in self.cancelScheduledWaiter(id) } + } + } + + private func waitForCompletionCallbacksWithoutDeadline(_ count: Int) async throws { + 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) + } + } + } onCancel: { + Task { @MainActor in self.cancelCompletionWaiter(id) } + } + } + + private func cancelScheduledWaiter(_ id: UUID) { + self.scheduledWaiters.removeValue(forKey: id)?.continuation.resume(throwing: CancellationError()) + } + + private func cancelCompletionWaiter(_ id: UUID) { + 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() + } + } + + 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() + } + } +} + +@MainActor +private final class RealtimePCMPlaybackResultProbe { + private(set) var results: [StreamingPlaybackResult] = [] + + func record(_ result: StreamingPlaybackResult) { + self.results.append(result) + } +} + +@MainActor +private func makeRealtimePCMPlayer( + backend: RealtimePCMPlaybackBackend) -> RealtimePCMStreamingAudioPlayer +{ + RealtimePCMStreamingAudioPlayer( + preparePlayback: backend.prepare, + scheduleFrame: backend.schedule, + stopPlayback: backend.stop, + playbackTime: { nil }) +} + +private func waitForPlayback(_ task: Task, label: String) async throws { + try await AsyncTimeout.withTimeout( + seconds: realtimePCMPlaybackWaitTimeoutSeconds, + onTimeout: { RealtimePCMPlaybackWaitTimeout(label: label) }, + operation: { await task.value }) +} + +@MainActor +struct RealtimePCMStreamingAudioPlayerTests { + private let sampleRate = 8000.0 + private var frameBytes: Int { + Int(self.sampleRate * RealtimePCMStreamingAudioPlayer.frameDurationSeconds) * 2 + } + + @Test func `prepare failure returns unfinished playback`() async { + let player = RealtimePCMStreamingAudioPlayer( + preparePlayback: { _ in throw RealtimePCMPlaybackFailure() }, + scheduleFrame: { _, _, _ in }, + stopPlayback: {}, + playbackTime: { nil }) + let (stream, continuation) = AsyncThrowingStream.makeStream() + continuation.finish() + + let result = await player.play(stream: stream, sampleRate: self.sampleRate) + + #expect(!result.finished) + #expect(result.interruptedAt == nil) + } + + @Test func `schedule failure returns unfinished playback`() async throws { + let player = RealtimePCMStreamingAudioPlayer( + preparePlayback: { _ in }, + scheduleFrame: { _, _, _ in throw RealtimePCMPlaybackFailure() }, + stopPlayback: {}, + playbackTime: { nil }) + let (stream, continuation) = AsyncThrowingStream.makeStream() + let playback = Task { + await player.play(stream: stream, sampleRate: self.sampleRate) + } + continuation.yield(Data(repeating: 1, count: self.frameBytes)) + continuation.finish() + + let probe = RealtimePCMPlaybackResultProbe() + let observed = Task { + await probe.record(playback.value) + } + try await waitForPlayback(observed, label: "schedule failure") + + #expect(probe.results.count == 1) + #expect(probe.results.first?.finished == false) + #expect(probe.results.first?.interruptedAt == nil) + } + + @Test func `withheld completions cap scheduling and one completion admits one frame`() async throws { + let backend = RealtimePCMPlaybackBackend() + let player = makeRealtimePCMPlayer(backend: backend) + let probe = RealtimePCMPlaybackResultProbe() + let (stream, continuation) = AsyncThrowingStream.makeStream() + let playback = Task { + let result = await player.play(stream: stream, sampleRate: self.sampleRate) + probe.record(result) + } + + continuation.yield(Data(repeating: 1, count: self.frameBytes * 5)) + continuation.finish() + try await backend.waitForScheduledFrames(3) + #expect(backend.scheduledFrames.count == 3) + #expect(backend.maxActiveCount == 3) + #expect(probe.results.isEmpty) + + backend.complete() + try await backend.waitForScheduledFrames(4) + #expect(backend.scheduledFrames.count == 4) + #expect(backend.maxActiveCount == 3) + #expect(probe.results.isEmpty) + + backend.complete() + try await backend.waitForScheduledFrames(5) + for _ in 0..<3 { + backend.complete() + } + try await waitForPlayback(playback, label: "five-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.allSatisfy { $0.count == self.frameBytes }) + } + + @Test func `playback finishes only after input and every scheduled frame complete`() async throws { + let backend = RealtimePCMPlaybackBackend() + let player = makeRealtimePCMPlayer(backend: backend) + let probe = RealtimePCMPlaybackResultProbe() + let (stream, continuation) = AsyncThrowingStream.makeStream() + let playback = Task { + let result = await player.play(stream: stream, sampleRate: self.sampleRate) + probe.record(result) + } + + continuation.yield(Data(repeating: 1, count: self.frameBytes * 2)) + continuation.finish() + try await backend.waitForScheduledFrames(2) + #expect(backend.completions.count == 2) + #expect(probe.results.isEmpty) + + backend.complete() + #expect(backend.completions.count == 1) + #expect(probe.results.isEmpty) + backend.complete() + try await waitForPlayback(playback, label: "completed input playback") + #expect(probe.results.count == 1) + #expect(probe.results.first?.finished == true) + #expect(probe.results.first?.interruptedAt == nil) + } + + @Test func `stop restart ignores stale buffer completions`() async throws { + let backend = RealtimePCMPlaybackBackend() + let player = makeRealtimePCMPlayer(backend: backend) + let (firstStream, firstContinuation) = AsyncThrowingStream.makeStream() + let firstProbe = RealtimePCMPlaybackResultProbe() + let firstPlayback = Task { + let result = await player.play(stream: firstStream, sampleRate: self.sampleRate) + firstProbe.record(result) + } + firstContinuation.yield(Data(repeating: 1, count: self.frameBytes)) + try await backend.waitForScheduledFrames(1) + let staleCompletion = backend.takeCompletion() + + _ = player.stop() + try await waitForPlayback(firstPlayback, label: "stopped A playback") + #expect(firstProbe.results.map(\.finished) == [false]) + + let (secondStream, secondContinuation) = AsyncThrowingStream.makeStream() + let probe = RealtimePCMPlaybackResultProbe() + let secondPlayback = Task { + let result = await player.play(stream: secondStream, sampleRate: self.sampleRate) + probe.record(result) + } + secondContinuation.yield(Data(repeating: 2, count: self.frameBytes * 5)) + secondContinuation.finish() + try await backend.waitForScheduledFrames(4) + #expect(backend.activeCount == 3) + #expect(probe.results.isEmpty) + + staleCompletion() + try await backend.waitForCompletionCallbacks(1) + #expect(backend.scheduledFrames.count == 4) + #expect(backend.completions.count == 3) + #expect(firstProbe.results.map(\.finished) == [false]) + #expect(probe.results.isEmpty) + backend.complete() + try await backend.waitForScheduledFrames(5) + #expect(backend.scheduledFrames.count == 5) + #expect(probe.results.isEmpty) + backend.complete() + try await backend.waitForScheduledFrames(6) + #expect(backend.scheduledFrames.count == 6) + #expect(probe.results.isEmpty) + for _ in 0..<3 { + backend.complete() + } + try await waitForPlayback(secondPlayback, label: "replacement B playback") + #expect(probe.results.count == 1) + #expect(probe.results.first?.finished == true) + #expect(probe.results.first?.interruptedAt == nil) + } + + @Test func `stop resumes the active playback exactly once`() async throws { + let backend = RealtimePCMPlaybackBackend() + let player = makeRealtimePCMPlayer(backend: backend) + let probe = RealtimePCMPlaybackResultProbe() + let (stream, continuation) = AsyncThrowingStream.makeStream() + let playback = Task { + let result = await player.play(stream: stream, sampleRate: self.sampleRate) + probe.record(result) + } + continuation.yield(Data(repeating: 1, count: self.frameBytes * 5)) + try await backend.waitForScheduledFrames(3) + + _ = player.stop() + _ = player.stop() + try await waitForPlayback(playback, label: "stopped active playback") + + #expect(probe.results.count == 1) + #expect(probe.results.first?.finished == false) + } +} +#endif diff --git a/apps/shared/OpenClawKit/Tests/OpenClawKitTests/RealtimeTalkRelaySessionTests.swift b/apps/shared/OpenClawKit/Tests/OpenClawKitTests/RealtimeTalkRelaySessionTests.swift index 4e1f6160fb48..e099856a215b 100644 --- a/apps/shared/OpenClawKit/Tests/OpenClawKitTests/RealtimeTalkRelaySessionTests.swift +++ b/apps/shared/OpenClawKit/Tests/OpenClawKitTests/RealtimeTalkRelaySessionTests.swift @@ -1,5 +1,4 @@ import Foundation -import OpenClawKit import OpenClawProtocol import Testing @testable import OpenClawKit @@ -15,52 +14,348 @@ private final class UnusedPCMStreamingAudioPlayer: PCMStreamingAudioPlaying { } } +@MainActor +private final class DrainingPCMStreamingAudioPlayer: PCMStreamingAudioPlaying { + private(set) var frames: [Data] = [] + private(set) var playCount = 0 + private let playbackStarted = RealtimeRelayTestSignal() + private let playbackFinished = RealtimeRelayTestSignal() + + func play(stream: AsyncThrowingStream, sampleRate _: Double) async -> StreamingPlaybackResult { + self.playCount += 1 + self.playbackStarted.send(self.playCount) + do { + for try await frame in stream { + self.frames.append(frame) + } + } catch {} + self.playbackFinished.send(()) + return StreamingPlaybackResult(finished: true, interruptedAt: nil) + } + + func stop() -> Double? { + nil + } + + func waitUntilPlaybackFinished() async throws { + _ = try await self.playbackFinished.next("draining playback to finish") + } + + func waitForPlaybackCount(_ expectedCount: Int) async throws { + while self.playCount < expectedCount { + _ = try await self.playbackStarted.next("\(expectedCount) draining playback starts") + } + } +} + +@MainActor +private final class StalledPCMStreamingAudioPlayer: PCMStreamingAudioPlaying { + private(set) var playCount = 0 + private(set) var stopCount = 0 + private var continuations: [CheckedContinuation] = [] + private let playbackStarted = RealtimeRelayTestSignal() + + func play( + stream _: AsyncThrowingStream, + sampleRate _: Double) async -> StreamingPlaybackResult + { + self.playCount += 1 + self.playbackStarted.send(self.playCount) + return await withCheckedContinuation { self.continuations.append($0) } + } + + func stop() -> Double? { + self.stopCount += 1 + let continuation = self.continuations.isEmpty ? nil : self.continuations.removeFirst() + continuation?.resume(returning: StreamingPlaybackResult(finished: false, interruptedAt: nil)) + return nil + } + + func waitForPlaybackCount(_ expectedCount: Int) async throws { + while self.playCount < expectedCount { + _ = try await self.playbackStarted.next("\(expectedCount) playback starts") + } + } +} + +private struct RealtimeRelayTestTimeout: Error, CustomStringConvertible { + let operation: String + + var description: String { + "timed out waiting for \(self.operation)" + } +} + +private final class RealtimeRelayTestSignal: @unchecked Sendable { + private struct Waiter { + let id: UUID + let continuation: CheckedContinuation + var deadline: Task? + } + + private let lock = NSLock() + private let timeoutSeconds: Double + private var values: [Value] = [] + private var waiters: [Waiter] = [] + + init(timeoutSeconds: Double = 30) { + self.timeoutSeconds = timeoutSeconds + } + + func send(_ value: Value) { + let waiter: Waiter? = self.lock.withLock { + guard !self.waiters.isEmpty else { + self.values.append(value) + return nil + } + return self.waiters.removeFirst() + } + self.resume(waiter, with: .success(value)) + } + + func next(_ operation: String) async throws -> Value { + try Task.checkCancellation() + let id = UUID() + return try await withTaskCancellationHandler { + try await withCheckedThrowingContinuation { continuation in + let registration: Result? = self.lock.withLock { + if Task.isCancelled { + return .failure(CancellationError()) + } + if !self.values.isEmpty { + return .success(self.values.removeFirst()) + } + self.waiters.append(Waiter(id: id, continuation: continuation)) + return nil + } + if let registration { + continuation.resume(with: registration) + return + } + let deadline = Task { + do { + try await Task.sleep(for: .seconds(self.timeoutSeconds)) + self.failWaiter(id, with: RealtimeRelayTestTimeout(operation: operation)) + } catch {} + } + let retained = self.lock.withLock { + guard let index = self.waiters.firstIndex(where: { $0.id == id }) else { + return false + } + self.waiters[index].deadline = deadline + return true + } + if !retained { + deadline.cancel() + } + } + } onCancel: { + self.failWaiter(id, with: CancellationError()) + } + } + + private func failWaiter(_ id: UUID, with error: any Error) { + self.resume(self.claimWaiter(id), with: .failure(error)) + } + + private func claimWaiter(_ id: UUID) -> Waiter? { + self.lock.withLock { + guard let index = self.waiters.firstIndex(where: { $0.id == id }) else { + return nil + } + return self.waiters.remove(at: index) + } + } + + private func resume(_ waiter: Waiter?, with result: Result) { + guard let waiter else { return } + waiter.deadline?.cancel() + waiter.continuation.resume(with: result) + } +} + +@MainActor +private final class IndexedPCMStreamingAudioPlayer: PCMStreamingAudioPlaying { + private(set) var activePlaybackIndexes: Set = [] + private var continuations: [Int: CheckedContinuation] = [:] + private var isShutdown = false + private let playbackStarted = RealtimeRelayTestSignal() + private let mainActorCheckpoint = RealtimeRelayTestSignal() + private var nextPlaybackIndex = 0 + + func play( + stream _: AsyncThrowingStream, + sampleRate _: Double) async -> StreamingPlaybackResult + { + guard !self.isShutdown else { + return StreamingPlaybackResult(finished: false, interruptedAt: nil) + } + let index = self.nextPlaybackIndex + self.nextPlaybackIndex += 1 + self.activePlaybackIndexes.insert(index) + let result = await withCheckedContinuation { continuation in + self.continuations[index] = continuation + self.playbackStarted.send(index) + } + self.mainActorCheckpoint.send(index) + return result + } + + func stop() -> Double? { + nil + } + + func shutdown() { + self.isShutdown = true + self.activePlaybackIndexes.removeAll() + let continuations = Array(self.continuations.values) + self.continuations.removeAll() + for continuation in continuations { + continuation.resume(returning: StreamingPlaybackResult(finished: false, interruptedAt: nil)) + } + } + + func waitForPlayback(_ expectedIndex: Int) async throws { + let index = try await self.playbackStarted.next("playback \(expectedIndex) to start") + guard index == expectedIndex else { + throw RealtimeRelayTestTimeout(operation: "playback \(expectedIndex), got \(index)") + } + } + + func complete(_ index: Int) { + self.activePlaybackIndexes.remove(index) + self.continuations.removeValue(forKey: index)?.resume( + returning: StreamingPlaybackResult(finished: true, interruptedAt: nil)) + } + + func fail(_ index: Int) { + self.activePlaybackIndexes.remove(index) + self.continuations.removeValue(forKey: index)?.resume( + returning: StreamingPlaybackResult(finished: false, interruptedAt: nil)) + } + + func waitUntilCompletionWasHandled(_ expectedIndex: Int) async throws { + let index = try await self.mainActorCheckpoint.next("playback \(expectedIndex) completion") + guard index == expectedIndex else { + throw RealtimeRelayTestTimeout( + operation: "playback \(expectedIndex) completion, got \(index)") + } + } +} + @MainActor private final class TestRealtimeTalkAudioCapture: RealtimeTalkAudioCapturing { var suppressesInputDuringOutput = false private(set) var isStarted = false private(set) var startCount = 0 private(set) var stopCount = 0 + private var onFailure: (@MainActor (String) -> Void)? + private let started = RealtimeRelayTestSignal() func start( targetSampleRate: Double, - onAudio: @escaping @Sendable (RealtimeTalkAudioFrame) -> Void) throws + onAudio: @escaping @Sendable (RealtimeTalkAudioFrame) -> Void, + onFailure: @escaping @MainActor (String) -> Void) throws { self.isStarted = true self.startCount += 1 + self.onFailure = onFailure + self.started.send(()) } func stop() { self.isStarted = false self.stopCount += 1 + self.onFailure = nil + } + + func fail(_ message: String) { + self.onFailure?(message) + } + + func waitUntilStarted() async throws { + _ = try await self.started.next("audio capture to start") } } private actor RealtimeRelayStartupBarrier { private var entered = false - private var enteredWaiter: CheckedContinuation? + private var released = false + private var enteredWaiter: (id: UUID, continuation: CheckedContinuation)? private var releaseWaiter: CheckedContinuation? func suspend() async { self.entered = true - self.enteredWaiter?.resume() - self.enteredWaiter = nil - await withCheckedContinuation { self.releaseWaiter = $0 } + if let waiter = self.enteredWaiter { + self.enteredWaiter = nil + waiter.continuation.resume() + } + guard !self.released else { return } + await withCheckedContinuation { continuation in + if self.released { continuation.resume() } else { self.releaseWaiter = continuation } + } } - func waitUntilEntered() async { - if self.entered { - return + func waitUntilEntered() async throws { + do { + try await AsyncTimeout.withTimeout( + seconds: 30, + onTimeout: { RealtimeRelayTestTimeout(operation: "request barrier entry") }, + operation: { try await self.waitUntilEnteredWithoutDeadline() }) + } catch { + self.release() + throw error } - await withCheckedContinuation { self.enteredWaiter = $0 } + } + + private func waitUntilEnteredWithoutDeadline() async throws { + guard !self.entered else { return } + let id = UUID() + try await withTaskCancellationHandler { + try await withCheckedThrowingContinuation { continuation in + if self.entered { + continuation.resume() + } else if Task.isCancelled { + continuation.resume(throwing: CancellationError()) + } else { + self.enteredWaiter = (id, continuation) + } + } + } onCancel: { + Task { await self.cancelEnteredWaiter(id) } + } + } + + private func cancelEnteredWaiter(_ id: UUID) { + guard let waiter = self.enteredWaiter, waiter.id == id else { return } + self.enteredWaiter = nil + waiter.continuation.resume(throwing: CancellationError()) } func release() { + guard !self.released else { return } + self.released = true self.releaseWaiter?.resume() self.releaseWaiter = nil } } +private func waitForRealtimeRelayEvent( + _ stream: AsyncStream, + operation: String) async throws -> Event +{ + try await AsyncTimeout.withTimeout( + seconds: 30, + onTimeout: { RealtimeRelayTestTimeout(operation: operation) }, + operation: { + var iterator = stream.makeAsyncIterator() + guard let event = await iterator.next() else { + throw RealtimeRelayTestTimeout(operation: "\(operation) before stream ended") + } + return event + }) +} + private struct RealtimeRelayStartupRequest: Sendable { let method: String let params: [String: AnyCodable]? @@ -68,14 +363,97 @@ private struct RealtimeRelayStartupRequest: Sendable { private actor RealtimeRelayStartupRequestLog { private var requests: [RealtimeRelayStartupRequest] = [] + private let requestObserved = RealtimeRelayTestSignal() func record(method: String, params: [String: AnyCodable]?) { self.requests.append(RealtimeRelayStartupRequest(method: method, params: params)) + self.requestObserved.send(self.requests.count) } func snapshot() -> [RealtimeRelayStartupRequest] { self.requests } + + func waitForRequestCount(_ expectedCount: Int) async throws { + while self.requests.count < expectedCount { + _ = try await self.requestObserved.next("\(expectedCount) relay requests") + } + } +} + +private enum ControlledAudioAppendBehavior { + case suspended + case requestFailure + case malformedResponse +} + +private actor ControlledRealtimeAudioRequests { + private let behavior: ControlledAudioAppendBehavior + private var methods: [String] = [] + private var appendContinuations: [CheckedContinuation] = [] + private let requestObserved = RealtimeRelayTestSignal() + + init(behavior: ControlledAudioAppendBehavior = .suspended) { + self.behavior = behavior + } + + func request(method: String) async throws -> Data { + self.methods.append(method) + self.requestObserved.send(self.methods.count) + guard method == "talk.session.appendAudio" else { + return Data("{\"ok\":true}".utf8) + } + switch self.behavior { + case .suspended: + return try await withCheckedThrowingContinuation { continuation in + self.appendContinuations.append(continuation) + } + case .requestFailure: + throw URLError(.badServerResponse) + case .malformedResponse: + return Data("{}".utf8) + } + } + + func waitForRequestCount(_ expectedCount: Int) async throws { + while self.methods.count < expectedCount { + _ = try await self.requestObserved.next("\(expectedCount) relay requests") + } + } + + func snapshot() -> [String] { + self.methods + } + + func succeedPendingAppends() { + let continuations = self.appendContinuations + self.appendContinuations.removeAll() + continuations.forEach { $0.resume(returning: Data("{\"ok\":true}".utf8)) } + } +} + +private actor RealtimeRelayRouteFlag { + private var isCurrent = true + + func expire() { + self.isCurrent = false + } + + func value() -> Bool { + self.isCurrent + } +} + +private actor RealtimeRelayEventSource { + private var continuation: AsyncStream.Continuation? + + func stream() -> AsyncStream { + AsyncStream { self.continuation = $0 } + } + + func finish() { + self.continuation?.finish() + } } private func unusedRealtimeRelayTransport() -> RealtimeTalkRelayTransport { @@ -84,8 +462,69 @@ private func unusedRealtimeRelayTransport() -> RealtimeTalkRelayTransport { request: { _, _, _ in throw CancellationError() }) } +private func outputAudioEvent( + turnId: String, + data: Data = Data([0x01]), + relaySessionId: String = "relay-1") -> EventFrame +{ + EventFrame( + type: "event", + event: "talk.event", + payload: AnyCodable([ + "relaySessionId": relaySessionId, + "type": "audio", + "audioBase64": data.base64EncodedString(), + "talkEvent": ["turnId": turnId], + ]), + seq: nil, + stateversion: nil) +} + +private func outputClearEvent(turnId: String) -> EventFrame { + EventFrame( + type: "event", + event: "talk.event", + payload: AnyCodable([ + "relaySessionId": "relay-1", + "type": "clear", + "talkEvent": ["turnId": turnId], + ]), + seq: nil, + stateversion: nil) +} + +private func outputAudioDoneEvent(turnId: String) -> EventFrame { + EventFrame( + type: "event", + event: "talk.event", + payload: AnyCodable([ + "relaySessionId": "relay-1", + "type": "audioDone", + "talkEvent": ["turnId": turnId], + ]), + seq: nil, + stateversion: nil) +} + @MainActor struct RealtimeTalkRelaySessionTests { + private func makeIdleCancellationSession( + _ onSpeakingChanged: @escaping (Bool) -> Void) -> RealtimeTalkRelaySession + { + let transport = RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in AsyncStream { $0.finish() } }, + request: { _, _, _ in Data("{\"ok\":true}".utf8) }) + let session = RealtimeTalkRelaySession( + transport: transport, + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: TestRealtimeTalkAudioCapture(), + pcmPlayer: DrainingPCMStreamingAudioPlayer(), + onStatus: { _ in }, + onSpeakingChanged: onSpeakingChanged) + session._test_setRelaySessionId("relay-1") + return session + } + @Test func `transcript callback carries typed partial and final values`() async { var transcripts: [RealtimeTalkTranscript] = [] let session = RealtimeTalkRelaySession( @@ -191,14 +630,10 @@ struct RealtimeTalkRelaySessionTests { ]), seq: nil, stateversion: nil)) - await Task.yield() #expect(await requests.snapshot().isEmpty) session._test_markOutputPlaybackFinished() - for _ in 0..<10 { - if !(await requests.snapshot()).isEmpty { break } - await Task.yield() - } + try await requests.waitForRequestCount(1) let recorded = await requests.snapshot() #expect(recorded.count == 1) @@ -207,6 +642,790 @@ struct RealtimeTalkRelaySessionTests { #expect(request.params?["sessionId"]?.stringValue == "relay-1") #expect(request.params?["markName"]?.stringValue == "audio-1") } +} + +extension RealtimeTalkRelaySessionTests { + @Test func `output buffer cap plus one terminates visibly and requests recovery`() async throws { + let requests = RealtimeRelayStartupRequestLog() + let player = StalledPCMStreamingAudioPlayer() + let terminationObserved = RealtimeRelayTestSignal() + var issues: [RealtimeTalkRelayIssue] = [] + var terminations: [RealtimeTalkRelayTermination] = [] + let session = RealtimeTalkRelaySession( + transport: RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in AsyncStream { $0.finish() } }, + request: { method, params, _ in + await requests.record(method: method, params: params) + return Data("{\"ok\":true}".utf8) + }), + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: TestRealtimeTalkAudioCapture(), + pcmPlayer: player, + onStatus: { _ in }, + onIssue: { issues.append($0) }, + onTermination: { + terminations.append($0) + terminationObserved.send($0) + }, + onSpeakingChanged: { _ in }) + session._test_setRelaySessionId("relay-1") + + for _ in 0...32 { + await session._test_handleGatewayEvent( + outputAudioEvent(turnId: "turn-1", data: Data(repeating: 1, count: 960))) + } + #expect(try await terminationObserved.next("output buffer overflow") == .outputPlaybackOverflow) + + #expect(issues.map(\.phase) == ["output-playback"]) + #expect(terminations == [.outputPlaybackOverflow]) + #expect(player.stopCount == 1) + try await requests.waitForRequestCount(1) + #expect(await requests.snapshot().contains(where: { $0.method == "talk.session.close" })) + } + + @Test func `new turn supersedes a stalled prior playback drain`() async throws { + let player = StalledPCMStreamingAudioPlayer() + let session = RealtimeTalkRelaySession( + transport: unusedRealtimeRelayTransport(), + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: TestRealtimeTalkAudioCapture(), + pcmPlayer: player, + onStatus: { _ in }, + onSpeakingChanged: { _ in }) + session._test_setRelaySessionId("relay-1") + + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-a")) + try await player.waitForPlaybackCount(1) + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-b")) + try await player.waitForPlaybackCount(2) + + #expect(player.playCount == 2) + #expect(player.stopCount == 1) + session.stop() + } + + @Test func `unfinished current playback terminates visibly and requests recovery`() async throws { + let requests = RealtimeRelayStartupRequestLog() + let player = IndexedPCMStreamingAudioPlayer() + let terminated = RealtimeRelayTestSignal() + let requestObserved = RealtimeRelayTestSignal() + var statuses: [String] = [] + var issues: [RealtimeTalkRelayIssue] = [] + let session = RealtimeTalkRelaySession( + transport: RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in AsyncStream { $0.finish() } }, + request: { method, params, _ in + await requests.record(method: method, params: params) + requestObserved.send(method) + return Data(#"{"ok":true}"#.utf8) + }), + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: TestRealtimeTalkAudioCapture(), + pcmPlayer: player, + onStatus: { statuses.append($0) }, + onIssue: { issues.append($0) }, + onTermination: { terminated.send($0) }, + onSpeakingChanged: { _ in }) + defer { + session.stop() + player.shutdown() + } + session._test_setRelaySessionId("relay-1") + + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-1")) + try await player.waitForPlayback(0) + player.fail(0) + #expect(try await terminated.next("unfinished playback recovery") == .outputPlaybackOverflow) + + let message = String(localized: "Realtime audio playback failed. Reconnecting…") + #expect(statuses == [message]) + #expect(issues.map(\.phase) == ["output-playback"]) + #expect(issues.map(\.message) == [message]) + #expect(!session._test_isOutputPlaying()) + #expect(try await requestObserved.next("relay close request") == "talk.session.close") + #expect(await requests.snapshot().map(\.method) == ["talk.session.close"]) + } + + @Test func `stale player completion cannot finish replacement turn playback`() async throws { + let player = IndexedPCMStreamingAudioPlayer() + let replacementFinished = AsyncStream.makeStream( + of: Void.self, + bufferingPolicy: .bufferingNewest(1)) + var replacementIsCompleting = false + var speakingStates: [Bool] = [] + let session = RealtimeTalkRelaySession( + transport: unusedRealtimeRelayTransport(), + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: TestRealtimeTalkAudioCapture(), + pcmPlayer: player, + onStatus: { _ in }, + onSpeakingChanged: { + speakingStates.append($0) + if replacementIsCompleting, !$0 { replacementFinished.continuation.yield() } + }) + defer { + replacementFinished.continuation.finish() + session.stop() + player.shutdown() + } + session._test_setRelaySessionId("relay-1") + + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-a")) + try await player.waitForPlayback(0) + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-b")) + try await player.waitForPlayback(1) + + #expect(player.activePlaybackIndexes.contains(1)) + #expect(speakingStates == [true, false, true]) + + player.complete(0) + try await player.waitUntilCompletionWasHandled(0) + + #expect(player.activePlaybackIndexes.contains(1)) + #expect(session._test_isOutputPlaying()) + #expect(speakingStates == [true, false, true]) + + replacementIsCompleting = true + player.complete(1) + _ = try await waitForRealtimeRelayEvent( + replacementFinished.stream, + operation: "replacement playback to finish") + + #expect(!session._test_isOutputPlaying()) + #expect(speakingStates == [true, false, true, false]) + } + + @Test func `cancelled playback task cannot start after its replacement`() async throws { + let player = DrainingPCMStreamingAudioPlayer() + let audioA = Data(repeating: 0x0A, count: 960) + let audioB = Data(repeating: 0x0B, count: 960) + let speakingChanged = RealtimeRelayTestSignal() + var speakingStates: [Bool] = [] + let session = RealtimeTalkRelaySession( + transport: unusedRealtimeRelayTransport(), + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: TestRealtimeTalkAudioCapture(), + pcmPlayer: player, + onStatus: { _ in }, + onSpeakingChanged: { + speakingStates.append($0) + speakingChanged.send($0) + }) + defer { session.stop() } + session._test_setRelaySessionId("relay-1") + + let replace = Task { @MainActor in + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-a", data: audioA)) + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-b", data: audioB)) + await session._test_handleGatewayEvent(outputAudioDoneEvent(turnId: "turn-b")) + } + await replace.value + try await player.waitForPlaybackCount(1) + for expected in [true, false, true, false] { + #expect(try await speakingChanged.next("playback state \(expected)") == expected) + } + + #expect(player.playCount == 1) + #expect(player.frames == [audioB]) + #expect(speakingStates == [true, false, true, false]) + } + + @Test func `stale relay session cannot clear successor playback`() async { + var speakingStates: [Bool] = [] + let session = RealtimeTalkRelaySession( + transport: unusedRealtimeRelayTransport(), + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: TestRealtimeTalkAudioCapture(), + pcmPlayer: StalledPCMStreamingAudioPlayer(), + onStatus: { _ in }, + onSpeakingChanged: { speakingStates.append($0) }) + session._test_setRelaySessionId("relay-1") + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-a")) + session._test_setRelaySessionId("relay-2") + await session._test_handleGatewayEvent( + outputAudioEvent(turnId: "turn-b", relaySessionId: "relay-2")) + + await session._test_handleGatewayEvent(EventFrame( + type: "event", + event: "talk.event", + payload: AnyCodable([ + "relaySessionId": "relay-1", + "type": "clear", + "talkEvent": ["turnId": "turn-a"], + ]), + seq: nil, + stateversion: nil)) + + #expect(session._test_isOutputPlaying()) + #expect(speakingStates == [true, false, true]) + session.stop() + } + + @Test func `output cancellation fences delayed audio and preserves exact identity`() async throws { + let requests = RealtimeRelayStartupRequestLog() + var speakingStates: [Bool] = [] + let transport = RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in AsyncStream { $0.finish() } }, + request: { method, params, _ in + await requests.record(method: method, params: params) + return Data("{\"ok\":true}".utf8) + }) + let session = RealtimeTalkRelaySession( + transport: transport, + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: TestRealtimeTalkAudioCapture(), + pcmPlayer: DrainingPCMStreamingAudioPlayer(), + onStatus: { _ in }, + onSpeakingChanged: { speakingStates.append($0) }) + session._test_setRelaySessionId("relay-1") + let audio: (String) -> EventFrame = { turnId in + EventFrame( + type: "event", + event: "talk.event", + payload: AnyCodable([ + "relaySessionId": "relay-1", + "type": "audio", + "audioBase64": Data([0x01]).base64EncodedString(), + "talkEvent": ["turnId": turnId], + ]), + seq: nil, + stateversion: nil) + } + let clear: (String) -> EventFrame = { turnId in + EventFrame( + type: "event", + event: "talk.event", + payload: AnyCodable([ + "relaySessionId": "relay-1", + "type": "clear", + "talkEvent": ["turnId": turnId], + ]), + seq: nil, + stateversion: nil) + } + + await session._test_handleGatewayEvent(audio("turn-7")) + session.cancelOutput(reason: "barge-in") + try await requests.waitForRequestCount(1) + let request = try #require(await requests.snapshot().first) + #expect(request.method == "talk.session.cancelOutput") + #expect(request.params?["sessionId"]?.stringValue == "relay-1") + #expect(request.params?["turnId"]?.stringValue == "turn-7") + #expect(request.params?["reason"]?.stringValue == "barge-in") + + await session._test_handleGatewayEvent(audio("turn-7")) + await session._test_handleGatewayEvent(audio("turn-8")) + #expect(speakingStates.first == true) + #expect(!speakingStates.dropFirst().contains(true)) + await session._test_handleGatewayEvent(clear("turn-7")) + await session._test_handleGatewayEvent(audio("turn-7")) + #expect(speakingStates == [true, false]) + await session._test_handleGatewayEvent(audio("turn-8")) + #expect(speakingStates == [true, false, true]) + await session._test_handleGatewayEvent(clear("turn-7")) + #expect(speakingStates == [true, false, true]) + await session._test_handleGatewayEvent(clear("turn-8")) + #expect(speakingStates == [true, false, true, false]) + } + + @Test func `idle cancellation and pause retain the relay without false interruption`() async { + let requests = RealtimeRelayStartupRequestLog() + var speakingStates: [Bool] = [] + let transport = RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in AsyncStream { $0.finish() } }, + request: { method, params, _ in + await requests.record(method: method, params: params) + return Data("{\"ok\":true}".utf8) + }) + let session = RealtimeTalkRelaySession( + transport: transport, + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: TestRealtimeTalkAudioCapture(), + pcmPlayer: DrainingPCMStreamingAudioPlayer(), + onStatus: { _ in }, + onSpeakingChanged: { speakingStates.append($0) }) + session._test_setRelaySessionId("relay-1") + + session.setOutputPaused(true) + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-1")) + #expect(speakingStates.isEmpty) + #expect(await requests.snapshot().isEmpty) + session.setOutputPaused(false) + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-2")) + #expect(speakingStates == [true]) + #expect(await requests.snapshot().isEmpty) + } + + @Test(arguments: ["stale", "idle"]) + func `non applied cancellation retires the wait without reopening the old turn`( + status: String) async throws + { + let barrier = RealtimeRelayStartupBarrier() + let speakingChanged = RealtimeRelayTestSignal() + var speakingStates: [Bool] = [] + let requests = RealtimeRelayStartupRequestLog() + let session = RealtimeTalkRelaySession( + transport: RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in AsyncStream { $0.finish() } }, + request: { method, params, _ in + await requests.record(method: method, params: params) + await barrier.suspend() + return Data("{\"ok\":true,\"status\":\"\(status)\"}".utf8) + }), + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: TestRealtimeTalkAudioCapture(), + pcmPlayer: DrainingPCMStreamingAudioPlayer(), + onStatus: { _ in }, + onSpeakingChanged: { + speakingStates.append($0) + speakingChanged.send($0) + }) + defer { session.stop() } + session._test_setRelaySessionId("relay-1") + session._test_prepareAudioSender(relaySessionId: "relay-1") + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-1")) + var successor: Task? + do { + #expect(try await speakingChanged.next("initial output") == true) + + #expect(session.cancelOutput()) + let cancellationTask = try #require(session._test_outputCancellationTask()) + #expect(try await speakingChanged.next("cancellation fence") == false) + try await barrier.waitUntilEntered() + #expect(await requests.snapshot().map(\.method) == ["talk.session.cancelOutput"]) + #expect(session._test_enqueueMicrophoneFrame(Data([0x01])) == nil) + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-2")) + #expect(speakingStates == [true, false]) + + await barrier.release() + await cancellationTask.value + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-1")) + successor = Task { @MainActor in + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-2")) + } + #expect(try await speakingChanged.next("successor output") == true) + await successor?.value + } catch { + await barrier.release() + successor?.cancel() + await successor?.value + throw error + } + + #expect(speakingStates == [true, false, true]) + } + + @Test func `cancellation without active identified output is a no-op`() async { + var speakingStates: [Bool] = [] + let session = self.makeIdleCancellationSession { speakingStates.append($0) } + #expect(!session.cancelOutput(reason: "barge-in")) + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-1")) + #expect(speakingStates == [true]) + + var unfencedStates: [Bool] = [] + let unfenced = RealtimeTalkRelaySession( + transport: unusedRealtimeRelayTransport(), + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: TestRealtimeTalkAudioCapture(), + pcmPlayer: DrainingPCMStreamingAudioPlayer(), + onStatus: { _ in }, + onSpeakingChanged: { unfencedStates.append($0) }) + #expect(!unfenced.cancelOutput()) + unfenced._test_setRelaySessionId("relay-1") + await unfenced._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-1")) + #expect(unfencedStates == [true]) + } + + @Test func `turn-scoped audio without a turn id terminates visibly`() async { + var issues: [RealtimeTalkRelayIssue] = [] + var terminations: [RealtimeTalkRelayTermination] = [] + let session = RealtimeTalkRelaySession( + transport: unusedRealtimeRelayTransport(), + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: TestRealtimeTalkAudioCapture(), + pcmPlayer: StalledPCMStreamingAudioPlayer(), + onStatus: { _ in }, + onIssue: { issues.append($0) }, + onTermination: { terminations.append($0) }, + onSpeakingChanged: { _ in }) + session._test_setRelaySessionId("relay-1") + + await session._test_handleGatewayEvent(EventFrame( + type: "event", + event: "talk.event", + payload: AnyCodable([ + "relaySessionId": "relay-1", + "type": "audio", + "audioBase64": Data([0x01]).base64EncodedString(), + ]), + seq: nil, + stateversion: nil)) + + #expect(issues.map(\.phase) == ["output-playback"]) + #expect(terminations == [.outputPlaybackOverflow]) + } + + @Test func `large provider audio delta is rebuffered into bounded ordered frames`() async throws { + let player = DrainingPCMStreamingAudioPlayer() + var issues: [RealtimeTalkRelayIssue] = [] + var terminations: [RealtimeTalkRelayTermination] = [] + let session = RealtimeTalkRelaySession( + transport: unusedRealtimeRelayTransport(), + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: TestRealtimeTalkAudioCapture(), + pcmPlayer: player, + onStatus: { _ in }, + onIssue: { issues.append($0) }, + onTermination: { terminations.append($0) }, + onSpeakingChanged: { _ in }) + defer { session.stop() } + session._test_setRelaySessionId("relay-1") + let audio = Data((0..<(960 * 2 + 480)).map { UInt8($0 % 251) }) + + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-1", data: audio)) + await session._test_handleGatewayEvent(outputAudioDoneEvent(turnId: "turn-1")) + try await player.waitUntilPlaybackFinished() + + #expect(player.frames.map(\.count) == [960, 960, 480]) + #expect(player.frames.reduce(into: Data()) { $0.append($1) } == audio) + #expect(issues.isEmpty) + #expect(terminations.isEmpty) + } + + @Test func `exact maximum output audio frame is accepted`() async { + var issues: [RealtimeTalkRelayIssue] = [] + var terminations: [RealtimeTalkRelayTermination] = [] + let session = RealtimeTalkRelaySession( + transport: unusedRealtimeRelayTransport(), + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: TestRealtimeTalkAudioCapture(), + pcmPlayer: StalledPCMStreamingAudioPlayer(), + onStatus: { _ in }, + onIssue: { issues.append($0) }, + onTermination: { terminations.append($0) }, + onSpeakingChanged: { _ in }) + session._test_setRelaySessionId("relay-1") + + await session._test_handleGatewayEvent( + outputAudioEvent(turnId: "turn-1", data: Data(repeating: 1, count: 960))) + + #expect(issues.isEmpty) + #expect(terminations.isEmpty) + #expect(session._test_isOutputPlaying()) + session.stop() + } + + @Test func `active output pause cancels the exact turn`() async throws { + let requests = RealtimeRelayStartupRequestLog() + let session = RealtimeTalkRelaySession( + transport: RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in AsyncStream { $0.finish() } }, + request: { method, params, _ in + await requests.record(method: method, params: params) + return Data("{\"ok\":true}".utf8) + }), + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: TestRealtimeTalkAudioCapture(), + pcmPlayer: DrainingPCMStreamingAudioPlayer(), + onStatus: { _ in }, + onSpeakingChanged: { _ in }) + session._test_setRelaySessionId("relay-1") + await session._test_handleGatewayEvent(EventFrame( + type: "event", + event: "talk.event", + payload: AnyCodable([ + "relaySessionId": "relay-1", + "type": "audio", + "audioBase64": Data([0x01]).base64EncodedString(), + "talkEvent": ["turnId": "turn-7"], + ]), + seq: nil, + stateversion: nil)) + + session.setOutputPaused(true) + try await requests.waitForRequestCount(1) + let request = try #require(await requests.snapshot().first) + #expect(request.params?["turnId"]?.stringValue == "turn-7") + #expect(request.params?["reason"]?.stringValue == "pause") + } + + @Test func `current cancellation failure terminates and rejects late audio`() async throws { + let requests = RealtimeRelayStartupRequestLog() + let issueObserved = RealtimeRelayTestSignal() + let terminationObserved = RealtimeRelayTestSignal() + var issues: [RealtimeTalkRelayIssue] = [] + var terminations: [RealtimeTalkRelayTermination] = [] + var speakingStates: [Bool] = [] + let cancellationError = URLError(.cannotConnectToHost) + let transport = RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in AsyncStream { $0.finish() } }, + request: { method, params, _ in + await requests.record(method: method, params: params) + if method == "talk.session.cancelOutput" { + throw cancellationError + } + return Data("{\"ok\":true}".utf8) + }) + let session = RealtimeTalkRelaySession( + transport: transport, + options: .init(sessionKey: "main", provider: "openai", model: "gpt-realtime-2", voice: nil), + audioCapture: TestRealtimeTalkAudioCapture(), + pcmPlayer: DrainingPCMStreamingAudioPlayer(), + onStatus: { _ in }, + onIssue: { + issues.append($0) + issueObserved.send($0) + }, + onTermination: { + terminations.append($0) + terminationObserved.send($0) + }, + onSpeakingChanged: { speakingStates.append($0) }) + session._test_setRelaySessionId("relay-1") + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-1")) + + session.cancelOutput() + _ = try await issueObserved.next("output cancellation issue") + #expect(try await terminationObserved.next("output cancellation termination") == .outputCancellationFailed) + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-2")) + try await requests.waitForRequestCount(2) + + #expect(issues.map(\.code) == ["realtime_output_cancel_failed"]) + #expect(issues.map(\.phase) == ["output-cancel"]) + #expect(issues.first?.message == String( + format: String(localized: "Realtime output cancellation failed: %@"), + cancellationError.localizedDescription)) + #expect(terminations == [.outputCancellationFailed]) + #expect(await requests.snapshot().map(\.method) == [ + "talk.session.cancelOutput", + "talk.session.close", + ]) + #expect(speakingStates.first == true) + #expect(!speakingStates.dropFirst().contains(true)) + } + + @Test(arguments: [ + #"{"ok":true,"turnId":"turn-2"}"#, + #"{"ok":true,"status":"applied","turnId":"turn-2"}"#, + ]) + func `accepted cancellation result with mismatched turn fails closed`( + response: String) async throws + { + let issueObserved = RealtimeRelayTestSignal() + let terminationObserved = RealtimeRelayTestSignal() + let requests = RealtimeRelayStartupRequestLog() + var issues: [RealtimeTalkRelayIssue] = [] + let session = RealtimeTalkRelaySession( + transport: RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in AsyncStream { $0.finish() } }, + request: { method, params, _ in + await requests.record(method: method, params: params) + return method == "talk.session.cancelOutput" + ? Data(response.utf8) + : Data(#"{"ok":true}"#.utf8) + }), + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: TestRealtimeTalkAudioCapture(), + pcmPlayer: DrainingPCMStreamingAudioPlayer(), + onStatus: { _ in }, + onIssue: { + issues.append($0) + issueObserved.send($0) + }, + onTermination: { terminationObserved.send($0) }, + onSpeakingChanged: { _ in }) + defer { session.stop() } + session._test_setRelaySessionId("relay-1") + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-1")) + + #expect(session.cancelOutput()) + let issue = try await issueObserved.next("mismatched cancellation issue") + #expect(try await terminationObserved.next("mismatched cancellation termination") == .outputCancellationFailed) + try await requests.waitForRequestCount(2) + + #expect(issue.code == "realtime_output_cancel_failed") + #expect(issue.phase == "output-cancel") + #expect(issues.count == 1) + #expect(await requests.snapshot().map(\.method) == [ + "talk.session.cancelOutput", + "talk.session.close", + ]) + } + + @Test func `superseded cancellation failure leaves the active fence intact`() async throws { + let barrier = RealtimeRelayStartupBarrier() + let requests = RealtimeRelayStartupRequestLog() + var issues: [RealtimeTalkRelayIssue] = [] + var speakingStates: [Bool] = [] + let transport = RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in AsyncStream { $0.finish() } }, + request: { method, params, _ in + await requests.record(method: method, params: params) + if await requests.snapshot().count == 1 { + await barrier.suspend() + throw URLError(.cancelled) + } + return Data("{\"ok\":true}".utf8) + }) + let session = RealtimeTalkRelaySession( + transport: transport, + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: TestRealtimeTalkAudioCapture(), + pcmPlayer: DrainingPCMStreamingAudioPlayer(), + onStatus: { _ in }, + onIssue: { issues.append($0) }, + onSpeakingChanged: { speakingStates.append($0) }) + session._test_setRelaySessionId("relay-1") + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-1")) + + #expect(session.cancelOutput()) + let staleCancellationTask = try #require(session._test_outputCancellationTask()) + try await barrier.waitUntilEntered() + await session._test_handleGatewayEvent(outputClearEvent(turnId: "turn-1")) + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-2")) + #expect(session.cancelOutput()) + await barrier.release() + await staleCancellationTask.value + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-1")) + + #expect(issues.isEmpty) + #expect(await requests.snapshot().count == 2) + #expect(speakingStates == [true, false, true, false]) + } + + @Test func `clear keeps microphone fenced until cancellation response`() async throws { + let barrier = RealtimeRelayStartupBarrier() + let requests = RealtimeRelayStartupRequestLog() + let session = RealtimeTalkRelaySession( + transport: RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in AsyncStream { $0.finish() } }, + request: { method, params, _ in + await requests.record(method: method, params: params) + if method == "talk.session.cancelOutput" { + await barrier.suspend() + return Data("{\"ok\":true,\"status\":\"applied\",\"turnId\":\"turn-1\"}".utf8) + } + return Data("{\"ok\":true}".utf8) + }), + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: TestRealtimeTalkAudioCapture(), + pcmPlayer: DrainingPCMStreamingAudioPlayer(), + onStatus: { _ in }, + onSpeakingChanged: { _ in }) + session._test_setRelaySessionId("relay-1") + session._test_prepareAudioSender(relaySessionId: "relay-1") + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-1")) + #expect(session.cancelOutput()) + let cancellationTask = try #require(session._test_outputCancellationTask()) + try await barrier.waitUntilEntered() + await session._test_handleGatewayEvent(EventFrame( + type: "event", + event: "talk.event", + payload: AnyCodable([ + "relaySessionId": "relay-1", + "type": "clear", + "talkEvent": ["turnId": "turn-1"], + ]), + seq: nil, + stateversion: nil)) + + #expect(session._test_enqueueMicrophoneFrame(Data([0x01])) == nil) + await barrier.release() + await cancellationTask.value + let admittedTask = try #require(session._test_enqueueMicrophoneFrame(Data([0x02]))) + await admittedTask.value + #expect(await requests.snapshot().map(\.method) == [ + "talk.session.cancelOutput", + "talk.session.appendAudio", + ]) + } + + @Test(arguments: [ + #"{"ok":true}"#, + #"{"ok":true,"status":"applied","turnId":"turn-1"}"#, + ]) + func `accepted cancellation response keeps fence until matching clear`( + response: String) async throws + { + let barrier = RealtimeRelayStartupBarrier() + var speakingStates: [Bool] = [] + let session = RealtimeTalkRelaySession( + transport: RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in AsyncStream { $0.finish() } }, + request: { method, _, _ in + if method == "talk.session.cancelOutput" { + await barrier.suspend() + return Data(response.utf8) + } + return Data(#"{"ok":true}"#.utf8) + }), + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: TestRealtimeTalkAudioCapture(), + pcmPlayer: DrainingPCMStreamingAudioPlayer(), + onStatus: { _ in }, + onSpeakingChanged: { speakingStates.append($0) }) + defer { session.stop() } + session._test_setRelaySessionId("relay-1") + session._test_prepareAudioSender(relaySessionId: "relay-1") + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-1")) + #expect(session.cancelOutput()) + let cancellationTask = try #require(session._test_outputCancellationTask()) + do { + try await barrier.waitUntilEntered() + await barrier.release() + await cancellationTask.value + #expect(session._test_enqueueMicrophoneFrame(Data([0x01])) == nil) + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-2")) + #expect(speakingStates == [true, false]) + + await session._test_handleGatewayEvent(outputClearEvent(turnId: "turn-1")) + let admittedTask = try #require(session._test_enqueueMicrophoneFrame(Data([0x02]))) + await admittedTask.value + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-2")) + } catch { + await barrier.release() + throw error + } + #expect(speakingStates == [true, false, true]) + } + + @Test func `close retires in flight cancellation failure`() async throws { + let barrier = RealtimeRelayStartupBarrier() + var issues: [RealtimeTalkRelayIssue] = [] + let session = RealtimeTalkRelaySession( + transport: RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in AsyncStream { $0.finish() } }, + request: { method, _, _ in + if method == "talk.session.cancelOutput" { + await barrier.suspend() + throw URLError(.cancelled) + } + return Data("{\"ok\":true}".utf8) + }), + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: TestRealtimeTalkAudioCapture(), + pcmPlayer: DrainingPCMStreamingAudioPlayer(), + onStatus: { _ in }, + onIssue: { issues.append($0) }, + onSpeakingChanged: { _ in }) + session._test_setRelaySessionId("relay-1") + + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-1")) + #expect(session.cancelOutput()) + let cancellationTask = try #require(session._test_outputCancellationTask()) + do { + try await barrier.waitUntilEntered() + session.stop() + await barrier.release() + await cancellationTask.value + } catch { + session.stop() + await barrier.release() + throw error + } + + #expect(issues.isEmpty) + } @Test func `close after classified error does not replace issue`() async { var issues: [RealtimeTalkRelayIssue] = [] @@ -251,6 +1470,295 @@ struct RealtimeTalkRelaySessionTests { #expect(statuses == ["OpenAI API key rejected with 401"]) } + @Test func `pre-ready relay failure throws and closes created session`() async throws { + let requests = RealtimeRelayStartupRequestLog() + let result = TalkSessionCreateResult( + sessionid: "talk-session", + mode: AnyCodable("realtime"), + transport: AnyCodable("gateway-relay"), + brain: AnyCodable("agent-consult"), + relaysessionid: "relay-1") + let resultData = try JSONEncoder().encode(result) + let failureEvent = EventFrame( + type: "event", + event: "talk.event", + payload: AnyCodable([ + "relaySessionId": "relay-1", + "type": "error", + "message": "OpenAI API key rejected with 401", + "phase": "connect", + ]), + seq: nil, + stateversion: nil) + let transport = RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in + AsyncStream { continuation in + continuation.yield(failureEvent) + } + }, + request: { method, params, _ in + await requests.record(method: method, params: params) + if method == "talk.session.create" { + return resultData + } + return Data("{\"ok\":true}".utf8) + }) + let audioCapture = TestRealtimeTalkAudioCapture() + let session = RealtimeTalkRelaySession( + transport: transport, + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: audioCapture, + pcmPlayer: UnusedPCMStreamingAudioPlayer(), + onStatus: { _ in }, + onSpeakingChanged: { _ in }) + + do { + try await session.start() + Issue.record("Expected the pre-ready relay failure to throw") + } catch { + #expect(error.localizedDescription == "OpenAI API key rejected with 401") + } + + let recorded = await requests.snapshot() + #expect(recorded.map(\.method) == ["talk.session.create", "talk.session.close"]) + #expect(recorded.last?.params?["sessionId"]?.stringValue == "relay-1") + #expect(!audioCapture.isStarted) + } + + @Test func `pre-ready event stream end promptly fails startup and closes created session once`() async throws { + let requests = RealtimeRelayStartupRequestLog() + let eventChannel = AsyncStream.makeStream() + let result = TalkSessionCreateResult( + sessionid: "talk-session", + mode: AnyCodable("realtime"), + transport: AnyCodable("gateway-relay"), + brain: AnyCodable("agent-consult"), + relaysessionid: "relay-1") + let resultData = try JSONEncoder().encode(result) + let transport = RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in eventChannel.stream }, + request: { method, params, _ in + await requests.record(method: method, params: params) + if method == "talk.session.create" { + return resultData + } + return Data("{\"ok\":true}".utf8) + }) + let audioCapture = TestRealtimeTalkAudioCapture() + var issues: [RealtimeTalkRelayIssue] = [] + let session = RealtimeTalkRelaySession( + transport: transport, + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: audioCapture, + pcmPlayer: UnusedPCMStreamingAudioPlayer(), + onStatus: { _ in }, + onIssue: { issues.append($0) }, + onSpeakingChanged: { _ in }) + let start = Task { @MainActor in try await session.start() } + try await audioCapture.waitUntilStarted() + + let disconnectedAt = ContinuousClock.now + eventChannel.continuation.finish() + do { + try await start.value + Issue.record("Expected the pre-ready event stream end to throw") + } catch { + #expect(error.localizedDescription == "Realtime connection ended before it became ready.") + } + + #expect(disconnectedAt.duration(to: .now) < .seconds(1)) + #expect(issues.map(\.phase) == ["connect"]) + let recorded = await requests.snapshot() + #expect(recorded.map(\.method) == ["talk.session.create", "talk.session.close"]) + #expect(recorded.last?.params?["sessionId"]?.stringValue == "relay-1") + #expect(!audioCapture.isStarted) + } + + @Test func `event stream ending during relay creation closes the late relay`() async throws { + let barrier = RealtimeRelayStartupBarrier() + let events = RealtimeRelayEventSource() + let requests = RealtimeRelayStartupRequestLog() + let audioCapture = TestRealtimeTalkAudioCapture() + let issueNotification = AsyncStream.makeStream( + of: RealtimeTalkRelayIssue.self, bufferingPolicy: .bufferingNewest(1)) + let result = TalkSessionCreateResult( + sessionid: "talk-session", + mode: AnyCodable("realtime"), + transport: AnyCodable("gateway-relay"), + brain: AnyCodable("agent-consult"), + relaysessionid: "relay-1") + let resultData = try JSONEncoder().encode(result) + let session = RealtimeTalkRelaySession( + transport: RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in await events.stream() }, + request: { method, params, _ in + await requests.record(method: method, params: params) + if method == "talk.session.create" { + await barrier.suspend() + return resultData + } + return Data("{\"ok\":true}".utf8) + }), + options: .init(sessionKey: "main", provider: "openai", model: "gpt-realtime-2", voice: nil), + audioCapture: audioCapture, + pcmPlayer: UnusedPCMStreamingAudioPlayer(), + onStatus: { _ in }, + onIssue: { issueNotification.continuation.yield($0) }, + onSpeakingChanged: { _ in }) + let start = Task { @MainActor in try await session.start() } + do { + try await barrier.waitUntilEntered() + await events.finish() + let issue = try await waitForRealtimeRelayEvent( + issueNotification.stream, + operation: "relay startup issue") + await barrier.release() + + var caughtStartupError: NSError? + do { + try await start.value + Issue.record("Expected relay startup to fail") + } catch { + caughtStartupError = error as NSError + } + let startupError = try #require(caughtStartupError) + #expect(startupError.domain == "RealtimeTalkRelay") + #expect(startupError.code == 6) + #expect(issue.code == "realtime_unavailable") + #expect(issue.phase == "connect") + #expect(issue.transport == "gateway-relay") + #expect(!issue.message.isEmpty) + #expect(audioCapture.startCount == 0) + let recorded = await requests.snapshot() + #expect(recorded.map(\.method) == ["talk.session.create", "talk.session.close"]) + #expect(recorded.last?.params?["sessionId"]?.stringValue == "relay-1") + issueNotification.continuation.finish() + } catch { + await barrier.release() + session.stop() + start.cancel() + _ = try? await start.value + issueNotification.continuation.finish() + throw error + } + } + + @Test func `microphone failure terminates relay and reports typed issue`() async throws { + let requests = RealtimeRelayStartupRequestLog() + let issueObserved = RealtimeRelayTestSignal() + let terminationObserved = RealtimeRelayTestSignal() + let transport = RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in AsyncStream { _ in } }, + request: { method, params, _ in + await requests.record(method: method, params: params) + return Data("{\"ok\":true}".utf8) + }) + let audioCapture = TestRealtimeTalkAudioCapture() + var issues: [RealtimeTalkRelayIssue] = [] + var terminations: [RealtimeTalkRelayTermination] = [] + let session = RealtimeTalkRelaySession( + transport: transport, + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: audioCapture, + pcmPlayer: UnusedPCMStreamingAudioPlayer(), + onStatus: { _ in }, + onIssue: { + issues.append($0) + issueObserved.send($0) + }, + onTermination: { + terminations.append($0) + terminationObserved.send($0) + }, + onSpeakingChanged: { _ in }) + session._test_setRelaySessionId("relay-1") + try session._test_startMicrophonePump() + + audioCapture.fail("Realtime microphone became unavailable: no input") + _ = try await issueObserved.next("microphone failure issue") + _ = try await terminationObserved.next("microphone failure termination") + try await requests.waitForRequestCount(1) + + #expect(issues.map(\.code) == ["audio_input_unavailable"]) + #expect(issues.map(\.phase) == ["audio-input"]) + #expect(terminations == [.audioInputFailed( + message: "Realtime microphone became unavailable: no input")]) + #expect(!audioCapture.isStarted) + let recorded = await requests.snapshot() + #expect(recorded.map(\.method) == ["talk.session.close"]) + #expect(recorded.first?.params?["sessionId"]?.stringValue == "relay-1") + } + + @Test func `ready then close publishes one typed termination and releases capture`() async { + var statuses: [String] = [] + var terminations: [RealtimeTalkRelayTermination] = [] + let audioCapture = TestRealtimeTalkAudioCapture() + let session = RealtimeTalkRelaySession( + transport: unusedRealtimeRelayTransport(), + options: .init(sessionKey: "main", provider: "openai", model: "gpt-realtime-2", voice: nil), + audioCapture: audioCapture, + pcmPlayer: UnusedPCMStreamingAudioPlayer(), + onStatus: { statuses.append($0) }, + onTermination: { terminations.append($0) }, + onSpeakingChanged: { _ in }) + session._test_setRelaySessionId("relay-1") + + await session._test_handleGatewayEvent(EventFrame( + type: "event", + event: "talk.event", + payload: AnyCodable([ + "relaySessionId": "relay-1", + "type": "ready", + ]), + seq: nil, + stateversion: nil)) + let closeEvent = EventFrame( + type: "event", + event: "talk.event", + payload: AnyCodable([ + "relaySessionId": "relay-1", + "type": "close", + "reason": "completed", + ]), + seq: nil, + stateversion: nil) + await session._test_handleGatewayEvent(closeEvent) + await session._test_handleGatewayEvent(closeEvent) + + #expect(statuses == ["Listening (Realtime)", "Ready"]) + #expect(terminations == [.remoteClose(reason: "completed")]) + #expect(audioCapture.stopCount == 1) + } + + @Test func `ready then event stream end publishes typed termination`() async { + var terminations: [RealtimeTalkRelayTermination] = [] + let audioCapture = TestRealtimeTalkAudioCapture() + let session = RealtimeTalkRelaySession( + transport: unusedRealtimeRelayTransport(), + options: .init(sessionKey: "main", provider: "openai", model: "gpt-realtime-2", voice: nil), + audioCapture: audioCapture, + pcmPlayer: UnusedPCMStreamingAudioPlayer(), + onStatus: { _ in }, + onTermination: { terminations.append($0) }, + onSpeakingChanged: { _ in }) + session._test_setRelaySessionId("relay-1") + + await session._test_handleGatewayEvent(EventFrame( + type: "event", + event: "talk.event", + payload: AnyCodable([ + "relaySessionId": "relay-1", + "type": "ready", + ]), + seq: nil, + stateversion: nil)) + await session._test_handleEventStreamEnded() + await session._test_handleEventStreamEnded() + + #expect(terminations == [.eventStreamEnded]) + #expect(audioCapture.stopCount == 1) + } + @Test func `closed relay does not wait for startup ready`() async { let session = RealtimeTalkRelaySession( transport: unusedRealtimeRelayTransport(), @@ -265,18 +1773,6 @@ struct RealtimeTalkRelaySessionTests { #expect(await session._test_waitForStartupCancelled(timeoutSeconds: 1)) } - @Test func `startup ready wait covers gateway connect budget`() { - let session = RealtimeTalkRelaySession( - transport: unusedRealtimeRelayTransport(), - options: .init(sessionKey: "main", provider: "openai", model: "gpt-realtime-2", voice: nil), - audioCapture: TestRealtimeTalkAudioCapture(), - pcmPlayer: UnusedPCMStreamingAudioPlayer(), - onStatus: { _ in }, - onSpeakingChanged: { _ in }) - - #expect(session._test_startupReadyTimeoutSeconds() >= 12) - } - @Test func `stop during event subscription prevents relay creation`() async throws { let barrier = RealtimeRelayStartupBarrier() let requests = RealtimeRelayStartupRequestLog() @@ -299,7 +1795,7 @@ struct RealtimeTalkRelaySessionTests { onStatus: { statuses.append($0) }, onSpeakingChanged: { speakingStates.append($0) }) let start = Task { @MainActor in try await session.start() } - await barrier.waitUntilEntered() + try await barrier.waitUntilEntered() session.stop() await barrier.release() @@ -340,7 +1836,7 @@ struct RealtimeTalkRelaySessionTests { onStatus: { statuses.append($0) }, onSpeakingChanged: { speakingStates.append($0) }) let start = Task { @MainActor in try await session.start() } - await barrier.waitUntilEntered() + try await barrier.waitUntilEntered() session.stop() await barrier.release() @@ -353,7 +1849,7 @@ struct RealtimeTalkRelaySessionTests { #expect(!speakingStates.contains(true)) } - @Test func `stop during buffered tool call prevents late relay side effects`() async { + @Test func `stop during buffered tool call prevents late relay side effects`() async throws { let barrier = RealtimeRelayStartupBarrier() let requests = RealtimeRelayStartupRequestLog() var statuses: [String] = [] @@ -389,7 +1885,7 @@ struct RealtimeTalkRelaySessionTests { seq: nil, stateversion: nil)) } - await barrier.waitUntilEntered() + try await barrier.waitUntilEntered() session.stop() await barrier.release() @@ -402,7 +1898,137 @@ struct RealtimeTalkRelaySessionTests { #expect(statuses == ["Thinking…"]) } - @Test func `stop cancels buffered microphone audio before dispatch`() async throws { + @Test func `stop and pause retire buffered audio while resume admits a fresh frame`() async throws { + let stoppedRequests = ControlledRealtimeAudioRequests() + var stoppedStatuses: [String] = [] + var stoppedIssues: [RealtimeTalkRelayIssue] = [] + var stoppedTerminations: [RealtimeTalkRelayTermination] = [] + let stoppedSession = RealtimeTalkRelaySession( + transport: RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in AsyncStream { $0.finish() } }, + request: { method, _, _ in try await stoppedRequests.request(method: method) }), + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: TestRealtimeTalkAudioCapture(), + pcmPlayer: UnusedPCMStreamingAudioPlayer(), + onStatus: { stoppedStatuses.append($0) }, + onIssue: { stoppedIssues.append($0) }, + onTermination: { stoppedTerminations.append($0) }, + onSpeakingChanged: { _ in }) + stoppedSession._test_setRelaySessionId("relay-1") + stoppedSession._test_prepareAudioSender(relaySessionId: "relay-1") + let stoppedSend = try #require(stoppedSession._test_enqueueMicrophoneFrame(Data([0x01]))) + do { + try await stoppedRequests.waitForRequestCount(1) + stoppedSession.stop() + try await stoppedRequests.waitForRequestCount(2) + } catch { + stoppedSession.stop() + await stoppedRequests.succeedPendingAppends() + await stoppedSend.value + throw error + } + await stoppedRequests.succeedPendingAppends() + await stoppedRequests.succeedPendingAppends() + await stoppedSend.value + #expect(await stoppedRequests.snapshot() == [ + "talk.session.appendAudio", + "talk.session.close", + ]) + #expect(stoppedStatuses.isEmpty) + #expect(stoppedIssues.isEmpty) + #expect(stoppedTerminations.isEmpty) + + let requests = ControlledRealtimeAudioRequests() + var statuses: [String] = [] + var issues: [RealtimeTalkRelayIssue] = [] + 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: "openai", model: nil, voice: nil), + audioCapture: TestRealtimeTalkAudioCapture(), + pcmPlayer: UnusedPCMStreamingAudioPlayer(), + onStatus: { statuses.append($0) }, + onIssue: { issues.append($0) }, + onTermination: { terminations.append($0) }, + onSpeakingChanged: { _ in }) + session._test_setRelaySessionId("relay-1") + session._test_prepareAudioSender(relaySessionId: "relay-1") + let pausedSend = try #require(session._test_enqueueMicrophoneFrame(Data([0x01]))) + var resumedSend: Task? + do { + try await requests.waitForRequestCount(1) + try session.setInputPaused(true) + #expect(await requests.snapshot() == ["talk.session.appendAudio"]) + try session.setInputPaused(false) + guard let freshSend = session._test_enqueueMicrophoneFrame(Data([0x02])) else { + throw RealtimeRelayTestTimeout(operation: "resumed microphone frame admission") + } + resumedSend = freshSend + try await requests.waitForRequestCount(2) + } catch { + session.stop() + await requests.succeedPendingAppends() + await pausedSend.value + await resumedSend?.value + throw error + } + await requests.succeedPendingAppends() + await requests.succeedPendingAppends() + await pausedSend.value + await resumedSend?.value + #expect(await requests.snapshot() == [ + "talk.session.appendAudio", + "talk.session.appendAudio", + ]) + session.stop() + try await requests.waitForRequestCount(3) + #expect(await requests.snapshot() == [ + "talk.session.appendAudio", + "talk.session.appendAudio", + "talk.session.close", + ]) + #expect(statuses.isEmpty) + #expect(issues.isEmpty) + #expect(terminations.isEmpty) + } + + @Test func `gateway route lost during startup fails instead of reporting ready`() async throws { + let route = RealtimeRelayRouteFlag() + let audioCapture = TestRealtimeTalkAudioCapture() + let session = RealtimeTalkRelaySession( + transport: RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in + // The Gateway replaces the route immediately after the subscription lands. + await route.expire() + return AsyncStream { $0.finish() } + }, + request: { _, _, _ in Data("{\"ok\":true}".utf8) }, + isCurrent: { await route.value() }), + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: audioCapture, + pcmPlayer: UnusedPCMStreamingAudioPlayer(), + onStatus: { _ in }, + onSpeakingChanged: { _ in }) + + do { + try await session.start() + Issue.record("Expected a lost Gateway route to fail startup") + } catch is CancellationError { + // The runtime returns silently on CancellationError, so classifying route loss as + // cancellation would leave Talk marked listening with no relay and no fallback. + Issue.record("Route loss must not surface as local cancellation") + } catch { + #expect( + error.localizedDescription == + "Gateway connection was replaced before realtime startup finished") + } + + #expect(!audioCapture.isStarted) + } + + @Test func `appended audio timestamps stay whole milliseconds`() async throws { let requests = RealtimeRelayStartupRequestLog() let transport = RealtimeTalkRelayTransport( subscribeServerEvents: { _ in AsyncStream { $0.finish() } }, @@ -418,11 +2044,138 @@ struct RealtimeTalkRelaySessionTests { onStatus: { _ in }, onSpeakingChanged: { _ in }) session._test_prepareAudioSender(relaySessionId: "relay-1") - let send = try #require(session._test_enqueueMicrophoneFrame(Data([0x01, 0x02]))) - session.stop() + // macOS taps stamp frames with `systemUptime * 1000`, so the raw value is fractional. + let send = try #require( + session._test_enqueueMicrophoneFrame(Data([0x01, 0x02]), timestampMs: 4823.617)) await send.value - #expect(await requests.snapshot().isEmpty) + let recorded = await requests.snapshot() + #expect(recorded.map(\.method) == ["talk.session.appendAudio"]) + // A decimal reaches the provider as a non-integer `audio_end_ms` and its + // `conversation.item.truncate` is rejected, ending the session on the first barge-in. + let timestamp = try #require(recorded.first?.params?["timestamp"]?.value as? Double) + #expect(timestamp == 4824) + #expect(timestamp == timestamp.rounded()) + } + + @Test func `microphone saturation terminates once without sending the fifth frame`() async throws { + let requests = ControlledRealtimeAudioRequests() + let audioCapture = TestRealtimeTalkAudioCapture() + var statuses: [String] = [] + var issues: [RealtimeTalkRelayIssue] = [] + var terminations: [RealtimeTalkRelayTermination] = [] + let terminationObserved = RealtimeRelayTestSignal() + let session = RealtimeTalkRelaySession( + transport: RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in AsyncStream { $0.finish() } }, + request: { method, _, _ in try await requests.request(method: method) }), + options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil), + audioCapture: audioCapture, + pcmPlayer: UnusedPCMStreamingAudioPlayer(), + onStatus: { statuses.append($0) }, + onIssue: { issues.append($0) }, + onTermination: { + terminations.append($0) + terminationObserved.send($0) + }, + onSpeakingChanged: { _ in }) + session._test_setRelaySessionId("relay-1") + session._test_prepareAudioSender(relaySessionId: "relay-1") + try session._test_startMicrophonePump() + let stopCountBeforeFailure = audioCapture.stopCount + + var pending: [Task] = [] + var saturated: Task? + do { + for index in 0..<4 { + guard let send = session._test_enqueueMicrophoneFrame(Data([UInt8(index)])) else { + throw RealtimeRelayTestTimeout(operation: "microphone frame \(index) admission") + } + pending.append(send) + } + try await requests.waitForRequestCount(4) + 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) + } catch { + saturated?.cancel() + pending.forEach { $0.cancel() } + session.stop() + await requests.succeedPendingAppends() + await saturated?.value + for task in pending { + await task.value + } + throw error + } + + 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(statuses == [message]) + #expect(issues.map(\.code) == ["audio_input_unavailable"]) + #expect(issues.map(\.message) == [message]) + #expect(terminations == [.audioInputFailed(message: message)]) + #expect(audioCapture.stopCount == stopCountBeforeFailure + 1) + + await requests.succeedPendingAppends() + await requests.succeedPendingAppends() + for task in pending { + await task.value + } + #expect(statuses == [message]) + #expect(issues.count == 1) + #expect(terminations.count == 1) + #expect(await requests.snapshot().filter { $0 == "talk.session.close" }.count == 1) + } + + @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) + let audioCapture = TestRealtimeTalkAudioCapture() + var statuses: [String] = [] + var issues: [RealtimeTalkRelayIssue] = [] + 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: "openai", model: nil, voice: nil), + audioCapture: audioCapture, + pcmPlayer: UnusedPCMStreamingAudioPlayer(), + onStatus: { statuses.append($0) }, + onIssue: { issues.append($0) }, + onTermination: { terminations.append($0) }, + onSpeakingChanged: { _ in }) + session._test_setRelaySessionId("relay-1") + session._test_prepareAudioSender(relaySessionId: "relay-1") + try session._test_startMicrophonePump() + let stopCountBeforeFailure = audioCapture.stopCount + + let send = try #require(session._test_enqueueMicrophoneFrame(Data([0x01]))) + await send.value + try await requests.waitForRequestCount(2) + + let issue = try #require(issues.first) + #expect(issue.code == "audio_input_unavailable") + #expect(issue.phase == "audio-input") + #expect(statuses == [issue.message]) + #expect(terminations == [.audioInputFailed(message: issue.message)]) + #expect(audioCapture.stopCount == stopCountBeforeFailure + 1) + #expect(await requests.snapshot() == [ + "talk.session.appendAudio", + "talk.session.close", + ]) + } } }