diff --git a/apps/.i18n/native-source.json b/apps/.i18n/native-source.json index 838cfbd5af2e..8baf031c8438 100644 --- a/apps/.i18n/native-source.json +++ b/apps/.i18n/native-source.json @@ -25037,6 +25037,10 @@ { "kind": "ui-localized-call", "path": "apps/ios/Sources/Voice/TalkModeManager.swift" + }, + { + "kind": "ui-localized-call", + "path": "apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift" } ] }, @@ -29897,6 +29901,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 +29934,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", @@ -32253,6 +32279,10 @@ { "kind": "ui-localized-call", "path": "apps/ios/Sources/Voice/TalkModeManager.swift" + }, + { + "kind": "ui-localized-call", + "path": "apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift" } ] }, @@ -39238,6 +39268,10 @@ { "kind": "conditional-branch", "path": "apps/macos/Sources/OpenClaw/SkillsSettings.swift" + }, + { + "kind": "ui-localized-call", + "path": "apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift" } ] }, @@ -39326,6 +39360,17 @@ } ] }, + { + "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.90f64f78a77015d7", "source": "Realtime audio playback fell behind. Reconnecting…", @@ -39337,6 +39382,39 @@ } ] }, + { + "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", @@ -39348,6 +39426,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", @@ -39356,6 +39445,21 @@ { "kind": "ui-localized-call", "path": "apps/ios/Sources/Voice/TalkModeManager.swift" + }, + { + "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" } ] }, @@ -39370,6 +39474,17 @@ } ] }, + { + "id": "native.apple.ddb5344f621cd627", + "source": "Realtime unavailable — using native speech", + "surface": "apple", + "sites": [ + { + "kind": "ui-localized-call", + "path": "apps/macos/Sources/OpenClaw/TalkModeRuntime.swift" + } + ] + }, { "id": "native.apple.4c945b4b74b4de2d", "source": "Realtime voice did not start.", @@ -46452,6 +46567,10 @@ { "kind": "ui-localized-call", "path": "apps/ios/Sources/Voice/TalkModeManager.swift" + }, + { + "kind": "ui-localized-call", + "path": "apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift" } ] }, @@ -48869,6 +48988,10 @@ { "kind": "ui-localized-call", "path": "apps/ios/Sources/Voice/TalkModeManager.swift" + }, + { + "kind": "ui-localized-call", + "path": "apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift" } ] }, diff --git a/apps/macos/Sources/OpenClaw/GatewayConnection+RealtimeTalk.swift b/apps/macos/Sources/OpenClaw/GatewayConnection+RealtimeTalk.swift index eb65cb0b38aa..297a7ed1c0ab 100644 --- a/apps/macos/Sources/OpenClaw/GatewayConnection+RealtimeTalk.swift +++ b/apps/macos/Sources/OpenClaw/GatewayConnection+RealtimeTalk.swift @@ -20,7 +20,16 @@ extension GatewayConnection { let task = Task { for await push in pushes { guard case let .event(event) = push else { continue } - continuation.yield(event) + switch continuation.yield(event) { + case .enqueued: + continue + case .dropped, .terminated: + continuation.finish() + return + @unknown default: + continuation.finish() + return + } } continuation.finish() } @@ -53,7 +62,16 @@ extension GatewayConnection { return } if let snapshot = self.lastSnapshot { - continuation.yield(.snapshot(snapshot)) + switch continuation.yield(.snapshot(snapshot)) { + case .enqueued: + break + case .dropped, .terminated: + continuation.finish() + return + @unknown default: + continuation.finish() + return + } } self.realtimeTalkSubscribers[lease.socketGeneration, default: [:]][id] = continuation continuation.onTermination = { @Sendable _ in @@ -66,7 +84,7 @@ extension GatewayConnection { } } - private func removeRealtimeTalkSubscriber(_ id: UUID, socketGeneration: UInt64) { + func removeRealtimeTalkSubscriber(_ id: UUID, socketGeneration: UInt64) { self.realtimeTalkSubscribers[socketGeneration]?[id] = nil if self.realtimeTalkSubscribers[socketGeneration]?.isEmpty == true { self.realtimeTalkSubscribers[socketGeneration] = nil diff --git a/apps/macos/Sources/OpenClaw/GatewayConnection.swift b/apps/macos/Sources/OpenClaw/GatewayConnection.swift index 8bf3c926d83e..9ab7cde88b1d 100644 --- a/apps/macos/Sources/OpenClaw/GatewayConnection.swift +++ b/apps/macos/Sources/OpenClaw/GatewayConnection.swift @@ -1228,8 +1228,21 @@ extension GatewayConnection { continuation.yield(push) } if let socketGeneration = self.activeSocketGeneration { - for (_, continuation) in self.realtimeTalkSubscribers[socketGeneration] ?? [:] { - continuation.yield(push) + var terminatedSubscriberIDs: [UUID] = [] + for (id, continuation) in self.realtimeTalkSubscribers[socketGeneration] ?? [:] { + switch continuation.yield(push) { + case .enqueued: + break + case .dropped, .terminated: + continuation.finish() + terminatedSubscriberIDs.append(id) + @unknown default: + continuation.finish() + terminatedSubscriberIDs.append(id) + } + } + for id in terminatedSubscriberIDs { + self.removeRealtimeTalkSubscriber(id, socketGeneration: socketGeneration) } } } diff --git a/apps/macos/Sources/OpenClaw/TalkModeRuntime.swift b/apps/macos/Sources/OpenClaw/TalkModeRuntime.swift index 3a85bf2f3187..041f1ac82430 100644 --- a/apps/macos/Sources/OpenClaw/TalkModeRuntime.swift +++ b/apps/macos/Sources/OpenClaw/TalkModeRuntime.swift @@ -22,12 +22,7 @@ actor TalkModeRuntime { case fallback } - typealias StartDependencies = ( - realtime: @Sendable (Int) async throws -> Void, - projectFallback: @Sendable () async -> Void, - fallback: @Sendable (Int) async -> Void) - - private struct NativeFallbackOwner { + struct NativeFallbackOwner { let lifecycleGeneration: Int let recognitionGeneration: Int let realtimeRelayGeneration: UInt64 @@ -107,7 +102,6 @@ actor TalkModeRuntime { private var pendingRealtimeRelayStartLifecycleGeneration: Int? var realtimeRestartGeneration: UInt64 = 0 var realtimeRestartTask: Task? - var startDependencies: StartDependencies? private var speechLocaleID: String? private var lastInterruptedAtSeconds: Double? private var voiceAliases: [String: String] = [:] @@ -261,18 +255,26 @@ actor TalkModeRuntime { self.realtimeSession == nil } + func transitionToNativeFallback( + owner: NativeFallbackOwner, + projectFailure: @Sendable () async -> Void) async -> Bool + { + guard self.canOwnNativeFallback(owner) else { return false } + await projectFailure() + // Projection crosses actors. Revalidate before native capture can replace + // a newer realtime or recognition owner. + return self.canOwnNativeFallback(owner) + } + func start() async { let gen = self.lifecycleGeneration guard voiceWakeSupported else { return } - let startDependencies = self.startDependencies - if startDependencies == nil { - guard await PermissionManager.ensureVoiceWakePermissions(interactive: true) else { - self.logger.error("talk runtime not starting: permissions missing") - return - } - await reloadConfig() + guard await PermissionManager.ensureVoiceWakePermissions(interactive: true) else { + self.logger.error("talk runtime not starting: permissions missing") + return } + await reloadConfig() guard self.isCurrent(gen) else { return } if self.isPaused { self.phase = .idle @@ -290,11 +292,7 @@ actor TalkModeRuntime { recognitionGeneration: self.recognitionGeneration, realtimeRelayGeneration: self.realtimeRelayGeneration &+ 1) do { - if let startDependencies { - try await startDependencies.realtime(gen) - } else { - try await self.startRealtimeRelay(generation: gen) - } + try await self.startRealtimeRelay(generation: gen) return } catch is CancellationError { if self.consumePendingRealtimeRelayStart() { @@ -307,24 +305,18 @@ actor TalkModeRuntime { self.logger.error( "talk realtime unavailable; using native fallback: " + "\(error.localizedDescription, privacy: .public)") - if let startDependencies { - await startDependencies.projectFallback() - } else { - await MainActor.run { - TalkModeController.shared.updatePartialTranscript( - "Realtime unavailable — using native speech") - } - } - // The UI projection crosses actors. A newer relay or recognition owner - // must keep the failed attempt from tearing down or replacing its capture. - guard self.canOwnNativeFallback(fallbackOwner) else { return } + guard await self.transitionToNativeFallback( + owner: fallbackOwner, + projectFailure: { + await MainActor.run { + TalkModeController.shared.updatePartialTranscript( + String(localized: "Realtime unavailable — using native speech")) + } + }) + else { return } } } - if let startDependencies { - await startDependencies.fallback(gen) - } else { - await self.startNativeFallback(generation: gen) - } + await self.startNativeFallback(generation: gen) } private func stop() async { @@ -372,17 +364,6 @@ actor TalkModeRuntime { (self.macOSRealtimeRelayOptIn, self.hasGatewayRealtimeRelayTuple) = (true, true) } - func _test_setStartDependencies( - startRealtimeRelay: @escaping @Sendable (Int) async throws -> Void, - projectNativeFallback: @escaping @Sendable () async -> Void = {}, - startNativeFallback: @escaping @Sendable (Int) async -> Void) - { - self.startDependencies = ( - realtime: startRealtimeRelay, - projectFallback: projectNativeFallback, - fallback: startNativeFallback) - } - func _test_beginRecognitionAttempt(lifecycleGeneration: Int) -> Int? { self.beginRecognitionAttempt(lifecycleGeneration: lifecycleGeneration) } @@ -1461,10 +1442,10 @@ extension TalkModeRuntime { case .speech: "barge-in" case .manual: "shutdown" } - await MainActor.run { + let cancelled = await MainActor.run { realtimeSession.cancelOutput(reason: relayReason) } - if reason != .manual, !self.isPaused { + if cancelled, reason != .manual, !self.isPaused { self.phase = .listening await MainActor.run { TalkModeController.shared.updatePhase(.listening) } } diff --git a/apps/macos/Tests/OpenClawIPCTests/GatewayConnectionControlTests.swift b/apps/macos/Tests/OpenClawIPCTests/GatewayConnectionControlTests.swift index d6b9b587a298..ddcd02b8f850 100644 --- a/apps/macos/Tests/OpenClawIPCTests/GatewayConnectionControlTests.swift +++ b/apps/macos/Tests/OpenClawIPCTests/GatewayConnectionControlTests.swift @@ -296,6 +296,55 @@ private func assertConfigLookupCannotRecreateRoute( } } + @Test func `realtime talk event overflow terminates its bounded subscription`() async throws { + let session = GatewayTestWebSocketSession(taskFactory: { + GatewayTestWebSocketTask( + sendHook: { task, message, sendIndex in + guard sendIndex > 0, + let data = Self.messageData(message), + let frame = try? JSONSerialization.jsonObject(with: data) as? [String: Any], + let id = frame["id"] as? String + else { return } + task.emitReceiveSuccess(.data(GatewayWebSocketTestSupport.okResponseData(id: id))) + }, + receiveHook: { task, receiveIndex in + if receiveIndex == 0 { + return .data(GatewayWebSocketTestSupport.connectChallengeData()) + } + let id = task.snapshotConnectRequestID() ?? "connect" + return .data(GatewayWebSocketTestSupport.connectOkData(id: id)) + }) + }) + let connection = GatewayConnection( + configProvider: { + ( + url: URL(string: "wss://gateway.example.invalid:9443")!, + token: "test-token-placeholder", + password: nil) + }, + sessionBox: WebSocketSessionBox(session: session)) + try await connection.refresh() + let transport = try await connection.acquireRealtimeTalkTransport() + let events = await transport.subscribeServerEvents(1) + let socketGeneration = try #require(await connection._test_activeSocketGeneration()) + + for seq in 1...20 { + await connection._test_handlePush( + .event(EventFrame( + type: "event", + event: "talk.event", + payload: AnyCodable(["seq": seq]), + seq: seq, + stateversion: nil)), + socketGeneration: socketGeneration) + } + + var iterator = events.makeAsyncIterator() + #expect(await iterator.next() != nil) + #expect(await iterator.next() == nil) + await connection.shutdown() + } + @Test func `operator widget capability refresh is shared and retained`() async throws { let rawOldSurface = "http://127.0.0.1:18789/__openclaw__/cap/old-token" let rawNewSurface = "http://127.0.0.1:18789/__openclaw__/cap/new-token" diff --git a/apps/macos/Tests/OpenClawIPCTests/TalkModeRuntimeSpeechTests.swift b/apps/macos/Tests/OpenClawIPCTests/TalkModeRuntimeSpeechTests.swift index e369368c521c..888003908d20 100644 --- a/apps/macos/Tests/OpenClawIPCTests/TalkModeRuntimeSpeechTests.swift +++ b/apps/macos/Tests/OpenClawIPCTests/TalkModeRuntimeSpeechTests.swift @@ -605,23 +605,27 @@ struct TalkModeRuntimeSpeechTests { sessionB.stop() } - @Test @MainActor func `current relay start failure selects native fallback`() async { + @Test @MainActor func `current relay failure owner can transition to native fallback`() async { let runtime = TalkModeRuntime() let lifecycleGeneration = await runtime._test_prepareEnabledLifecycle() - await runtime._test_enableRealtimeRelaySelection() let fallbackProbe = RuntimeCommitProbe() let player = RuntimeTestPCMPlayer() let session = makeRuntimeTestRealtimeSession(player: player) - await runtime._test_setStartDependencies( - startRealtimeRelay: { lifecycleGeneration in - try await runtime._test_startRealtimeRelay( - lifecycleGeneration: lifecycleGeneration, - makeSession: { session }, - start: { _ in throw RuntimeRelayStartError.failed }) - }, - startNativeFallback: { _ in fallbackProbe.record("native") }) + do { + try await runtime._test_startRealtimeRelay( + lifecycleGeneration: lifecycleGeneration, + makeSession: { session }, + start: { _ in throw RuntimeRelayStartError.failed }) + Issue.record("Expected relay start failure") + } catch {} + let owner = await TalkModeRuntime.NativeFallbackOwner( + lifecycleGeneration: lifecycleGeneration, + recognitionGeneration: runtime.recognitionGeneration, + realtimeRelayGeneration: runtime.realtimeRelayGeneration) - await runtime.start() + if await runtime.transitionToNativeFallback(owner: owner, projectFailure: {}) { + fallbackProbe.record("native") + } #expect(fallbackProbe.values() == ["native"]) #expect(await runtime._test_realtimeSessionIsActive() == false) @@ -635,33 +639,41 @@ struct TalkModeRuntimeSpeechTests { @Test @MainActor func `stale relay fallback cannot replace successor recognition owner`() async { let runtime = TalkModeRuntime() let lifecycleGeneration = await runtime._test_prepareEnabledLifecycle() - await runtime._test_enableRealtimeRelaySelection() let fallbackProjectionBarrier = RuntimeContinuationBarrier() let fallbackProbe = RuntimeCommitProbe() let player = RuntimeTestPCMPlayer() let session = makeRuntimeTestRealtimeSession(player: player) - await runtime._test_setStartDependencies( - startRealtimeRelay: { lifecycleGeneration in - try await runtime._test_startRealtimeRelay( - lifecycleGeneration: lifecycleGeneration, - makeSession: { session }, - start: { _ in throw RuntimeRelayStartError.failed }) - }, - projectNativeFallback: { - await MainActor.run { - TalkModeController.shared.updatePartialTranscript("stale fallback") - } - await fallbackProjectionBarrier.wait() - }, - startNativeFallback: { _ in fallbackProbe.record("native") }) + do { + try await runtime._test_startRealtimeRelay( + lifecycleGeneration: lifecycleGeneration, + makeSession: { session }, + start: { _ in throw RuntimeRelayStartError.failed }) + Issue.record("Expected relay start failure") + } catch {} + let owner = await TalkModeRuntime.NativeFallbackOwner( + lifecycleGeneration: lifecycleGeneration, + recognitionGeneration: runtime.recognitionGeneration, + realtimeRelayGeneration: runtime.realtimeRelayGeneration) - let start = Task { await runtime.start() } + let transition = Task { + let accepted = await runtime.transitionToNativeFallback( + owner: owner, + projectFailure: { + await MainActor.run { + TalkModeController.shared.updatePartialTranscript("stale fallback") + } + await fallbackProjectionBarrier.wait() + }) + if accepted { + fallbackProbe.record("native") + } + } await fallbackProjectionBarrier.waitUntilEntered() let successorRecognition = await runtime._test_beginRecognitionAttempt( lifecycleGeneration: lifecycleGeneration) TalkModeController.shared.updatePartialTranscript("successor") await fallbackProjectionBarrier.release() - await start.value + await transition.value #expect(successorRecognition != nil) #expect(await runtime.recognitionGeneration == successorRecognition) diff --git a/apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift b/apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift index 547abb1a85db..a85487dcbe22 100644 --- a/apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift +++ b/apps/shared/OpenClawKit/Sources/OpenClawKit/RealtimeTalkRelaySession.swift @@ -222,7 +222,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. + private nonisolated static let maxEncodedOutputFrameBytes = 1280 + private nonisolated static let maxDecodedOutputFrameBytes = 960 + /// 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 @@ -320,7 +322,7 @@ public final class RealtimeTalkRelaySession { self.startupIssue = nil self.startupWaiter = nil self.pendingPreRelayEvents.removeAll() - self.onStatus("Connecting realtime…") + self.onStatus(String(localized: "Connecting realtime…")) let eventStream = await self.transport.subscribeServerEvents(200) switch await self.lifecycleStatus(lifecycleGeneration) { case .current: break @@ -358,7 +360,8 @@ public final class RealtimeTalkRelaySession { !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 @@ -367,7 +370,7 @@ public final class RealtimeTalkRelaySession { request: self.transport.request) self.configureAudioContract(result.audio) try self.startMicrophonePump(lifecycleGeneration: lifecycleGeneration) - self.onStatus("Waiting for realtime…") + self.onStatus(String(localized: "Waiting for realtime…")) await self.drainPendingPreRelayEvents(lifecycleGeneration: lifecycleGeneration) switch await self.lifecycleStatus(lifecycleGeneration) { case .current: break @@ -438,7 +441,8 @@ public final class RealtimeTalkRelaySession { /// 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: "Gateway connection was replaced before realtime startup finished", + NSLocalizedDescriptionKey: String( + localized: "Gateway connection was replaced before realtime startup finished"), ]) } @@ -534,7 +538,7 @@ extension RealtimeTalkRelaySession { guard self.hasReceivedReady else { guard !self.hasReceivedFailure else { return } let issue = RealtimeTalkRelayIssue( - message: "Realtime connection ended before it became ready.", + message: String(localized: "Realtime connection ended before it became ready."), provider: self.options.provider, model: self.options.model, transport: "gateway-relay", @@ -546,7 +550,7 @@ extension RealtimeTalkRelaySession { self.finishStartupWait(.failed(issue)) return } - self.onStatus("Ready") + self.onStatus(String(localized: "Ready")) self.close(sendClose: false) self.onTermination(.eventStreamEnded) } @@ -571,7 +575,7 @@ extension RealtimeTalkRelaySession { case "ready": self.hasReceivedReady = true self.finishStartupWait(.ready) - self.onStatus("Listening (Realtime)") + self.onStatus(String(localized: "Listening (Realtime)")) case "audio": self.handleOutputAudio(payload) case "audioDone": @@ -585,7 +589,7 @@ extension 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, @@ -600,14 +604,14 @@ extension RealtimeTalkRelaySession { case "close": self.logger.debug("talk realtime: close") if self.hasReceivedReady { - self.onStatus("Ready") + self.onStatus(String(localized: "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", @@ -615,7 +619,7 @@ extension RealtimeTalkRelaySession { self.onIssue(issue) self.startupIssue = issue self.finishStartupWait(.failed(issue)) - self.onStatus("Realtime failed before connecting") + self.onStatus(String(localized: "Realtime failed before connecting")) } default: return @@ -694,7 +698,7 @@ extension 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", @@ -767,9 +771,9 @@ extension RealtimeTalkRelaySession { self.onTranscript(RealtimeTalkTranscript(role: role, text: text, isFinal: isFinal)) guard isFinal else { return } if role == "user" { - self.onStatus("Thinking…") + self.onStatus(String(localized: "Thinking…")) } else if role == "assistant" { - self.onStatus("Listening (Realtime)") + self.onStatus(String(localized: "Listening (Realtime)")) } } @@ -778,7 +782,7 @@ extension RealtimeTalkRelaySession { let callId = payload["callId"]?.stringValue, let name = payload["name"]?.stringValue else { return } - self.onStatus("Thinking…") + self.onStatus(String(localized: "Thinking…")) do { if name == Self.agentControlToolName { try await self.handleAgentControlToolCall( @@ -805,7 +809,8 @@ extension 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( @@ -821,7 +826,7 @@ extension RealtimeTalkRelaySession { result: result, lifecycleGeneration: lifecycleGeneration) try await self.ensureCurrentLifecycle(lifecycleGeneration) - self.onStatus("Listening (Realtime)") + self.onStatus(String(localized: "Listening (Realtime)")) } catch { guard await self.isCurrentLifecycle(lifecycleGeneration) else { return } let errorResult: [String: AnyCodable] = [ @@ -832,7 +837,7 @@ extension RealtimeTalkRelaySession { result: errorResult, lifecycleGeneration: lifecycleGeneration) guard await self.isCurrentLifecycle(lifecycleGeneration) else { return } - self.onStatus("Listening (Realtime)") + self.onStatus(String(localized: "Listening (Realtime)")) } } @@ -879,7 +884,7 @@ extension RealtimeTalkRelaySession { result: result, lifecycleGeneration: lifecycleGeneration) try await self.ensureCurrentLifecycle(lifecycleGeneration) - self.onStatus("Listening (Realtime)") + self.onStatus(String(localized: "Listening (Realtime)")) } private func submitToolResult( @@ -1165,9 +1170,12 @@ extension RealtimeTalkRelaySession { } } - public func cancelOutput(reason: String = "user") { - guard let relaySessionId else { return } - let outputIdentity = self.outputIdentity ?? OutputIdentity([:]) + @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() @@ -1176,13 +1184,11 @@ extension RealtimeTalkRelaySession { self.awaitingOutputClear = true self.stopOutputPlayback() self.outputCancellationTask = Task { [weak self, transport] in - var payload: [String: AnyCodable] = [ + let payload: [String: AnyCodable] = [ "sessionId": AnyCodable(relaySessionId), "reason": AnyCodable(reason), + "turnId": AnyCodable(turnId), ] - if let turnId = outputIdentity.turnId { - payload["turnId"] = AnyCodable(turnId) - } do { _ = try await transport.request("talk.session.cancelOutput", payload, 8000) } catch { @@ -1202,6 +1208,7 @@ extension RealtimeTalkRelaySession { self.onTermination(.outputCancellationFailed) } } + return true } private func isCurrentOutputCancellation(_ generation: UInt64) -> Bool { @@ -1210,30 +1217,33 @@ extension RealtimeTalkRelaySession { private func handleOutputAudio(_ payload: [String: AnyCodable]) { guard !self.isOutputPaused else { return } - guard let base64 = payload["audioBase64"]?.stringValue, - let data = Data(base64Encoded: base64) - else { return } let incomingIdentity = OutputIdentity(payload) + guard let incomingTurnId = incomingIdentity.turnId else { + self.handleOutputPlaybackOverflow() + return + } guard !self.awaitingOutputClear else { return } if let cancelledOutputTurnId { - if incomingIdentity.turnId == cancelledOutputTurnId { + if incomingTurnId == cancelledOutputTurnId { return } - if incomingIdentity.turnId != nil { - self.cancelledOutputTurnId = nil - } + } + guard let base64 = payload["audioBase64"]?.stringValue else { return } + guard base64.utf8.count <= Self.maxEncodedOutputFrameBytes, + let data = Data(base64Encoded: base64), + data.count <= Self.maxDecodedOutputFrameBytes + else { + self.handleOutputPlaybackOverflow() + return } if let currentIdentity = self.outputIdentity, - !incomingIdentity.isEmpty(), currentIdentity.relation(to: incomingIdentity) == .different { self.stopOutputPlayback() } else if self.outputContinuation == nil, self.outputTask != nil { self.stopOutputPlayback() } - if !incomingIdentity.isEmpty() { - self.outputIdentity = incomingIdentity - } + self.outputIdentity = incomingIdentity self.recordOutputAudioChunk(byteCount: data.count) self.markOutputAudioStarted(byteCount: data.count, nowMs: ProcessInfo.processInfo.systemUptime * 1000) self.onSpeakingChanged(true) @@ -1375,7 +1385,7 @@ extension RealtimeTalkRelaySession { self.audioCaptureGeneration == audioCaptureGeneration, !self.isInputPaused else { return } - self.onStatus("Realtime audio failed: \(message)") + self.onStatus(String(format: String(localized: "Realtime audio failed: %@"), message)) } self.audioSendTasks[taskID] = task return task diff --git a/apps/shared/OpenClawKit/Tests/OpenClawKitTests/RealtimePCMStreamingAudioPlayerTests.swift b/apps/shared/OpenClawKit/Tests/OpenClawKitTests/RealtimePCMStreamingAudioPlayerTests.swift index 6a8bea09b05c..ad758abe3af4 100644 --- a/apps/shared/OpenClawKit/Tests/OpenClawKitTests/RealtimePCMStreamingAudioPlayerTests.swift +++ b/apps/shared/OpenClawKit/Tests/OpenClawKitTests/RealtimePCMStreamingAudioPlayerTests.swift @@ -41,6 +41,15 @@ private final class RealtimePCMPlaybackBackend { } } +@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 @@ -75,7 +84,7 @@ struct RealtimePCMStreamingAudioPlayerTests { let stream = AsyncThrowingStream { continuation = $0 } let playback = Task { await player.play(stream: stream, sampleRate: self.sampleRate) } - continuation?.yield(Data(repeating: 1, count: self.frameBytes * 4)) + continuation?.yield(Data(repeating: 1, count: self.frameBytes * 5)) await waitUntil { backend.scheduledFrames.count == 3 } #expect(backend.scheduledFrames.count == 3) #expect(backend.maxActiveCount == 3) @@ -84,6 +93,10 @@ struct RealtimePCMStreamingAudioPlayerTests { await waitUntil { backend.scheduledFrames.count == 4 } #expect(backend.scheduledFrames.count == 4) #expect(backend.maxActiveCount == 3) + backend.complete() + await waitUntil { backend.scheduledFrames.count == 5 } + #expect(backend.scheduledFrames.count == 5) + #expect(backend.maxActiveCount == 3) continuation?.finish() while !backend.completions.isEmpty { @@ -139,5 +152,26 @@ struct RealtimePCMStreamingAudioPlayerTests { backend.complete() #expect(await (secondPlayback.value).finished) } + + @Test func `stop resumes the active playback exactly once`() async { + let backend = RealtimePCMPlaybackBackend() + let player = makeRealtimePCMPlayer(backend: backend) + let probe = RealtimePCMPlaybackResultProbe() + var continuation: AsyncThrowingStream.Continuation? + let stream = AsyncThrowingStream { continuation = $0 } + 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)) + await waitUntil { backend.scheduledFrames.count == 3 } + + _ = player.stop() + _ = player.stop() + await playback.value + + #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 0d809fefa08e..46521db9f303 100644 --- a/apps/shared/OpenClawKit/Tests/OpenClawKitTests/RealtimeTalkRelaySessionTests.swift +++ b/apps/shared/OpenClawKit/Tests/OpenClawKitTests/RealtimeTalkRelaySessionTests.swift @@ -152,14 +152,14 @@ private func unusedRealtimeRelayTransport() -> RealtimeTalkRelayTransport { request: { _, _, _ in throw CancellationError() }) } -private func outputAudioEvent(turnId: String) -> EventFrame { +private func outputAudioEvent(turnId: String, data: Data = Data([0x01])) -> EventFrame { EventFrame( type: "event", event: "talk.event", payload: AnyCodable([ "relaySessionId": "relay-1", "type": "audio", - "audioBase64": Data([0x01]).base64EncodedString(), + "audioBase64": data.base64EncodedString(), "talkEvent": ["turnId": turnId], ]), seq: nil, @@ -564,23 +564,12 @@ extension RealtimeTalkRelaySessionTests { #expect(await requests.snapshot().isEmpty) } - @Test func `idle cancellation waits for clear while cancellation without relay stays unfenced`() async { + @Test func `cancellation without active identified output is a no-op`() async { var speakingStates: [Bool] = [] let session = self.makeIdleCancellationSession { speakingStates.append($0) } - session.cancelOutput(reason: "barge-in") + #expect(!session.cancelOutput(reason: "barge-in")) await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-1")) - #expect(speakingStates == [false]) - await session._test_handleGatewayEvent(EventFrame( - type: "event", - event: "talk.event", - payload: AnyCodable([ - "relaySessionId": "relay-1", - "type": "clear", - ]), - seq: nil, - stateversion: nil)) - await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-1")) - #expect(speakingStates == [false, true]) + #expect(speakingStates == [true]) var unfencedStates: [Bool] = [] let unfenced = RealtimeTalkRelaySession( @@ -590,12 +579,81 @@ extension RealtimeTalkRelaySessionTests { pcmPlayer: DrainingPCMStreamingAudioPlayer(), onStatus: { _ in }, onSpeakingChanged: { unfencedStates.append($0) }) - unfenced.cancelOutput() + #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 `oversized current audio terminates but oversized cancelled audio stays fenced`() async { + var issues: [RealtimeTalkRelayIssue] = [] + var terminations: [RealtimeTalkRelayTermination] = [] + let session = RealtimeTalkRelaySession( + transport: RealtimeTalkRelayTransport( + subscribeServerEvents: { _ in AsyncStream { $0.finish() } }, + request: { _, _, _ in Data("{\"ok\":true}".utf8) }), + 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-a")) + #expect(session.cancelOutput()) + 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)) + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-b")) + await session._test_handleGatewayEvent( + outputAudioEvent(turnId: "turn-a", data: Data(repeating: 1, count: 961))) + #expect(issues.isEmpty) + #expect(terminations.isEmpty) + + await session._test_handleGatewayEvent( + outputAudioEvent(turnId: "turn-b", data: Data(repeating: 1, count: 961))) + #expect(issues.map(\.phase) == ["output-playback"]) + #expect(terminations == [.outputPlaybackOverflow]) + } + @Test func `active output pause cancels the exact turn`() async throws { let requests = RealtimeRelayStartupRequestLog() let session = RealtimeTalkRelaySession( @@ -705,10 +763,22 @@ extension RealtimeTalkRelaySessionTests { onIssue: { issues.append($0) }, onSpeakingChanged: { speakingStates.append($0) }) session._test_setRelaySessionId("relay-1") + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-1")) - session.cancelOutput() + #expect(session.cancelOutput()) await barrier.waitUntilEntered() - session.cancelOutput() + 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)) + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-2")) + #expect(session.cancelOutput()) await barrier.release() for _ in 0..<10 { if await requests.snapshot().count == 2 { break } @@ -718,7 +788,7 @@ extension RealtimeTalkRelaySessionTests { #expect(issues.isEmpty) #expect(await requests.snapshot().count == 2) - #expect(!speakingStates.contains(true)) + #expect(speakingStates == [true, false, true, false]) } @Test(arguments: [CancellationRetirement.clear, .close]) @@ -745,7 +815,8 @@ extension RealtimeTalkRelaySessionTests { onSpeakingChanged: { _ in }) session._test_setRelaySessionId("relay-1") - session.cancelOutput() + await session._test_handleGatewayEvent(outputAudioEvent(turnId: "turn-1")) + #expect(session.cancelOutput()) await barrier.waitUntilEntered() switch retirement { case .clear: @@ -755,6 +826,7 @@ extension RealtimeTalkRelaySessionTests { payload: AnyCodable([ "relaySessionId": "relay-1", "type": "clear", + "talkEvent": ["turnId": "turn-1"], ]), seq: nil, stateversion: nil)) diff --git a/docs/nodes/talk.md b/docs/nodes/talk.md index 4d5dfb650794..c5033dc7a37c 100644 --- a/docs/nodes/talk.md +++ b/docs/nodes/talk.md @@ -91,6 +91,20 @@ then 2 s). If those attempts are exhausted, the overlay reports `Realtime disconnected repeatedly — using native speech` and the next start bypasses realtime. Losing the microphone mid-session closes the relay and takes the same route. +Relay output cancellation is turn-scoped. Clients copy the current `turnId` from the +`talk.event` audio envelope; missing, empty, or stale turn ids are ignored: + +```json +{ + "method": "talk.session.cancelOutput", + "params": { + "sessionId": "relay-session-id", + "turnId": "turn-7", + "reason": "barge-in" + } +} +``` + ## Voice directives in replies The assistant can prefix a reply with a single JSON line to control voice: diff --git a/src/gateway/talk-realtime-relay-operations.ts b/src/gateway/talk-realtime-relay-operations.ts index 91d7e6c6d155..a62c46f0e13f 100644 --- a/src/gateway/talk-realtime-relay-operations.ts +++ b/src/gateway/talk-realtime-relay-operations.ts @@ -522,10 +522,10 @@ export function cancelTalkRealtimeRelayTurn(params: { }): void { const session = getRelaySession(params.relaySessionId, params.connId); const requestedTurnId = normalizeOptionalString(params.turnId); - if (requestedTurnId && session.harness.talk.activeTurnId !== requestedTurnId) { + if (!requestedTurnId || session.harness.talk.activeTurnId !== requestedTurnId) { return; } - const turnId = requestedTurnId ?? ensureRelayTurn(session); + const turnId = requestedTurnId; session.toolResultEpoch += 1; session.forcedTerminalProviderResults.clear(); const reason = params.reason ?? "client-cancelled"; diff --git a/src/gateway/talk-realtime-relay-session-create.ts b/src/gateway/talk-realtime-relay-session-create.ts index 9f6089a97d47..e912fcf0570d 100644 --- a/src/gateway/talk-realtime-relay-session-create.ts +++ b/src/gateway/talk-realtime-relay-session-create.ts @@ -63,6 +63,9 @@ import { registerTalkConnectionCleanup, } from "./talk-session-registry.js"; +const MAX_PROVIDER_OUTPUT_AUDIO_BYTES = 1_048_576; +const RELAY_OUTPUT_AUDIO_FRAME_BYTES = 960; + function isRelayAssistantEchoTranscript(session: RelaySession | undefined, text: string): boolean { return session?.harness.isLikelyAssistantEchoTranscript(text) ?? false; } @@ -218,21 +221,33 @@ export function createTalkRealtimeRelaySession( if (!relay) { return; } + if (audio.byteLength > MAX_PROVIDER_OUTPUT_AUDIO_BYTES) { + failSession( + `Realtime provider audio callback exceeded ${MAX_PROVIDER_OUTPUT_AUDIO_BYTES} bytes`, + ); + return; + } const turnId = ensureRelayTurn(relay); - emit( - { - relaySessionId, - type: "audio", - audioBase64: audio.toString("base64"), - ...(currentOutputItemId ? { itemId: currentOutputItemId } : {}), - ...(currentOutputResponseId ? { responseId: currentOutputResponseId } : {}), - }, - { - type: "output.audio.delta", - turnId, - payload: { byteLength: audio.length }, - }, - ); + for (let offset = 0; offset < audio.byteLength; offset += RELAY_OUTPUT_AUDIO_FRAME_BYTES) { + const frame = audio.subarray( + offset, + Math.min(offset + RELAY_OUTPUT_AUDIO_FRAME_BYTES, audio.byteLength), + ); + emit( + { + relaySessionId, + type: "audio", + audioBase64: frame.toString("base64"), + ...(currentOutputItemId ? { itemId: currentOutputItemId } : {}), + ...(currentOutputResponseId ? { responseId: currentOutputResponseId } : {}), + }, + { + type: "output.audio.delta", + turnId, + payload: { byteLength: frame.byteLength }, + }, + ); + } }, clearAudio: (reason) => { const relay = getActiveRelay(); diff --git a/src/gateway/talk-realtime-relay.test.ts b/src/gateway/talk-realtime-relay.test.ts index d25b5d3bd0b0..7e55826131cd 100644 --- a/src/gateway/talk-realtime-relay.test.ts +++ b/src/gateway/talk-realtime-relay.test.ts @@ -77,6 +77,17 @@ function stopTalkRealtimeRelaySession( activeRelaySessions.delete(params.relaySessionId); } +function ensureActiveRelayTurnId(relaySessionId: string): string { + const relay = relaySessions.get(relaySessionId); + if (!relay) { + throw new Error(`Missing relay test session ${relaySessionId}`); + } + if (!relay.harness.talk.activeTurnId) { + relay.harness.talk.startTurn({ turnId: "turn-1" }); + } + return relay.harness.talk.activeTurnId ?? "turn-1"; +} + describe("talk realtime gateway relay", () => { it.each([ [ @@ -1674,6 +1685,7 @@ describe("talk realtime gateway relay", () => { relaySessionId: session.relaySessionId, connId: "conn-1", reason: "barge-in", + turnId: ensureActiveRelayTurnId(session.relaySessionId), }); stopTalkRealtimeRelaySession({ relaySessionId: session.relaySessionId, connId: "conn-1" }); @@ -2598,6 +2610,7 @@ describe("talk realtime gateway relay", () => { relaySessionId: cancelledSession.relaySessionId, connId: "conn-1", reason: "barge-in", + turnId: ensureActiveRelayTurnId(cancelledSession.relaySessionId), }); await vi.advanceTimersByTimeAsync(250); expect( @@ -2642,7 +2655,7 @@ describe("talk realtime gateway relay", () => { ).toThrow("Unknown realtime relay session"); }); - it("correlates output audio with the active relay turn", () => { + it("splits provider audio into ordered 20 ms frames on one relay turn", () => { let bridgeRequest: RealtimeVoiceBridgeCreateRequest | undefined; const provider: RealtimeVoiceProviderPlugin = { id: "relay-test", @@ -2653,15 +2666,9 @@ describe("talk realtime gateway relay", () => { return makeRelayTransport(); }, }; - const events: Array<{ - event: string; - payload: { talkEvent?: { type?: string; turnId?: string } }; - }> = []; + const events: Array<{ event: string; payload: Record }> = []; const context = { - broadcastToConnIds: ( - event: string, - payload: { talkEvent?: { type?: string; turnId?: string } }, - ) => { + broadcastToConnIds: (event: string, payload: Record) => { events.push({ event, payload }); }, } as never; @@ -2679,24 +2686,88 @@ describe("talk realtime gateway relay", () => { connId: "conn-1", audioBase64: Buffer.from("audio").toString("base64"), }); - bridgeRequest?.onAudio(Buffer.from("reply")); + const audio = Buffer.alloc(1_921); + for (let index = 0; index < audio.length; index += 1) { + audio[index] = index % 251; + } + bridgeRequest?.onEvent?.({ + direction: "server", + type: "response.output_audio.delta", + itemId: "item-1", + responseId: "response-1", + }); + bridgeRequest?.onAudio(audio); + const audioEvents = events + .map((entry) => entry.payload) + .filter((payload) => payload.type === "audio"); + expect(audioEvents).toHaveLength(3); expect( - events.some( - (entry) => - entry.payload.talkEvent?.type === "output.audio.delta" && - entry.payload.talkEvent.turnId === "turn-1", + audioEvents.map((payload) => Buffer.from(String(payload.audioBase64), "base64").byteLength), + ).toEqual([960, 960, 1]); + expect( + Buffer.concat( + audioEvents.map((payload) => Buffer.from(String(payload.audioBase64), "base64")), ), - ).toBe(true); + ).toEqual(audio); + for (const payload of audioEvents) { + expectRecordFields(payload, { + itemId: "item-1", + responseId: "response-1", + }); + expectRecordFields(payload.talkEvent, { + type: "output.audio.delta", + turnId: "turn-1", + }); + expectRecordFields((payload.talkEvent as { payload?: unknown }).payload, { + byteLength: Buffer.from(String(payload.audioBase64), "base64").byteLength, + }); + } + }); + + it("fails an oversized provider audio callback before creating a turn", () => { + let bridgeRequest: RealtimeVoiceBridgeCreateRequest | undefined; + const provider: RealtimeVoiceProviderPlugin = { + id: "relay-test", + label: "Relay Test", + isConfigured: () => true, + createBridge: (req) => { + bridgeRequest = req; + return makeRelayTransport(); + }, + }; + const events: Array<{ event: string; payload: Record }> = []; + const session = createTalkRealtimeRelaySession({ + context: { + broadcastToConnIds: (event: string, payload: Record) => { + events.push({ event, payload }); + }, + } as never, + connId: "conn-1", + provider, + providerConfig: {}, + instructions: "brief", + tools: [], + }); + const relay = relaySessions.get(session.relaySessionId); + expect(relay?.harness.talk.activeTurnId).toBeUndefined(); + + bridgeRequest?.onAudio(Buffer.alloc(1_048_577)); + + expect(relay?.harness.talk.activeTurnId).toBeUndefined(); + expect(events.some((entry) => entry.payload.type === "audio")).toBe(false); + expect(events.some((entry) => entry.payload.type === "error")).toBe(true); }); it("aborts linked agent consult runs when the relay turn is cancelled", () => { const { abortController, broadcast, nodeSendToSession, removeChatRun, chatRunState, session } = createAbortableRelayRunFixture(); + relaySessions.get(session.relaySessionId)?.harness.talk.startTurn({ turnId: "turn-1" }); cancelTalkRealtimeRelayTurn({ relaySessionId: session.relaySessionId, connId: "conn-1", reason: "barge-in", + turnId: "turn-1", }); expect(abortController.signal.aborted).toBe(true); @@ -2734,6 +2805,35 @@ describe("talk realtime gateway relay", () => { expect(broadcast).not.toHaveBeenCalled(); }); + it("ignores missing and empty turn cancellation without mutating relay state", () => { + const { abortController, broadcast, session } = createAbortableRelayRunFixture(); + const relay = relaySessions.get(session.relaySessionId); + expect(relay).toBeDefined(); + relay?.harness.talk.startTurn({ turnId: "turn-b" }); + const epoch = relay?.toolResultEpoch; + const forcedResult = { + result: { status: "cancelled" }, + turnId: "turn-b", + epoch: epoch ?? 0, + }; + relay?.forcedTerminalProviderResults.set("call-1", forcedResult); + + for (const turnId of [undefined, "", " "]) { + cancelTalkRealtimeRelayTurn({ + relaySessionId: session.relaySessionId, + connId: "conn-1", + reason: "barge-in", + turnId, + }); + } + + expect(relay?.harness.talk.activeTurnId).toBe("turn-b"); + expect(relay?.toolResultEpoch).toBe(epoch); + expect(relay?.forcedTerminalProviderResults.get("call-1")).toBe(forcedResult); + expect(abortController.signal.aborted).toBe(false); + expect(broadcast).not.toHaveBeenCalled(); + }); + it("terminally satisfies a late normal result after turn cancellation without a new turn", async () => { const submitToolResult = vi.fn(); const provider = createIdleRelayProvider(); @@ -2746,6 +2846,7 @@ describe("talk realtime gateway relay", () => { relaySessionId: session.relaySessionId, connId: "conn-1", reason: "barge-in", + turnId: ensureActiveRelayTurnId(session.relaySessionId), }); const startedTurns = () => broadcastToConnIds.mock.calls.filter( @@ -2791,6 +2892,7 @@ describe("talk realtime gateway relay", () => { relaySessionId: session.relaySessionId, connId: "conn-1", reason: "barge-in", + turnId: ensureActiveRelayTurnId(session.relaySessionId), }); expect(abortController.signal.aborted).toBe(false); @@ -3005,6 +3107,7 @@ describe("talk realtime gateway relay", () => { relaySessionId: session.relaySessionId, connId: "conn-1", reason: "barge-in", + turnId: ensureActiveRelayTurnId(session.relaySessionId), }); const cancellation = submitTalkRealtimeRelayToolResult({ relaySessionId: session.relaySessionId, @@ -3063,6 +3166,7 @@ describe("talk realtime gateway relay", () => { relaySessionId: session.relaySessionId, connId: "conn-1", reason: "barge-in", + turnId: ensureActiveRelayTurnId(session.relaySessionId), }); workingAccepted.resolve(); await Promise.all([working, final]); @@ -3117,6 +3221,7 @@ describe("talk realtime gateway relay", () => { relaySessionId: session.relaySessionId, connId: "conn-1", reason: "barge-in", + turnId: ensureActiveRelayTurnId(session.relaySessionId), }); const cancellation = submitTalkRealtimeRelayToolResult({ relaySessionId: session.relaySessionId, @@ -3205,6 +3310,7 @@ describe("talk realtime gateway relay", () => { relaySessionId: session.relaySessionId, connId: "conn-1", reason: "barge-in", + turnId: ensureActiveRelayTurnId(session.relaySessionId), }); expect(abortController.signal.aborted).toBe(false); expect(broadcast).not.toHaveBeenCalledWith( @@ -3247,6 +3353,7 @@ describe("talk realtime gateway relay", () => { relaySessionId: session.relaySessionId, connId: "conn-1", reason: "barge-in", + turnId: ensureActiveRelayTurnId(session.relaySessionId), }); expect(abortController.signal.aborted).toBe(true); }); @@ -3386,6 +3493,7 @@ describe("talk realtime gateway relay", () => { relaySessionId: session.relaySessionId, connId: "conn-1", reason: "barge-in", + turnId: ensureActiveRelayTurnId(session.relaySessionId), }); expect(result).toMatchObject({ @@ -3629,6 +3737,7 @@ describe("talk realtime gateway relay", () => { relaySessionId: fixture.session.relaySessionId, connId: "conn-1", reason: "barge-in", + turnId: ensureActiveRelayTurnId(fixture.session.relaySessionId), }); const startedTurns = () => fixture.events.filter( @@ -3703,6 +3812,7 @@ describe("talk realtime gateway relay", () => { relaySessionId: fixture.session.relaySessionId, connId: "conn-1", reason: "barge-in", + turnId: ensureActiveRelayTurnId(fixture.session.relaySessionId), }); workingAccepted.resolve(); await final; @@ -3741,6 +3851,7 @@ describe("talk realtime gateway relay", () => { relaySessionId: fixture.session.relaySessionId, connId: "conn-1", reason: "barge-in", + turnId: ensureActiveRelayTurnId(fixture.session.relaySessionId), }); const cancellation = submitTalkRealtimeRelayToolResult({ relaySessionId: fixture.session.relaySessionId,