fix(talk): enforce bounded realtime turn ownership

This commit is contained in:
Vincent Koc
2026-08-20 02:30:46 -07:00
parent 03a117fc45
commit bcd3c4517a
13 changed files with 625 additions and 173 deletions
+123
View File
@@ -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"
}
]
},
@@ -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
@@ -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)
}
}
}
@@ -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<Void, Never>?
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) }
}
@@ -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"
@@ -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)
@@ -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
@@ -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<Data, Error> { 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<Data, Error>.Continuation?
let stream = AsyncThrowingStream<Data, Error> { 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
@@ -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))
+14
View File
@@ -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:
@@ -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";
@@ -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();
+126 -15
View File
@@ -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<string, unknown> }> = [];
const context = {
broadcastToConnIds: (
event: string,
payload: { talkEvent?: { type?: string; turnId?: string } },
) => {
broadcastToConnIds: (event: string, payload: Record<string, unknown>) => {
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<string, unknown> }> = [];
const session = createTalkRealtimeRelaySession({
context: {
broadcastToConnIds: (event: string, payload: Record<string, unknown>) => {
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<RealtimeVoiceBridge["submitToolResult"]>();
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,