refactor(talk): move Apple relay into OpenClawKit

This commit is contained in:
Zhilong Zheng
2026-08-21 05:46:39 -07:00
committed by Vincent Koc
parent 9bf2fbae1b
commit 697b2967f6
4 changed files with 1395 additions and 1196 deletions
File diff suppressed because it is too large Load Diff
+14 -2
View File
@@ -2321,6 +2321,10 @@ final class TalkModeManager: NSObject {
return .ignored
}
guard self.isCurrentStartAttempt(attemptID) else { return .ignored }
guard let gatewayRoute = await gateway.currentRoute() else {
return .unavailable(realtimeIssue(message: "Gateway not connected", phase: "start"))
}
guard self.isCurrentStartAttempt(attemptID) else { return .ignored }
if self.realtimeRelaySession != nil {
GatewayDiagnostics.log("talk realtime ignored: already active")
return .started
@@ -2342,19 +2346,27 @@ final class TalkModeManager: NSObject {
GatewayDiagnostics.log("talk.timeline realtime relay start attempt sessionKey=\(sessionKey)")
let startedAt = Self.nowSeconds()
let relaySession = RealtimeTalkRelaySession(
gateway: gateway,
transport: .ios(gateway: gateway, route: gatewayRoute),
options: RealtimeTalkRelaySession.Options(
sessionKey: sessionKey,
provider: self.realtimeProvider,
model: self.realtimeModelId,
voice: self.realtimeVoiceId),
audioCapture: IOSRealtimeTalkAudioCapture(),
pcmPlayer: self.pcmPlayer,
onStatus: { [weak self] status in
guard let self, self.realtimeRelayGeneration == relayGeneration else { return }
self.handleRealtimeRelayStatus(status)
},
onIssue: { [weak self] issue in
onIssue: { [weak self] relayIssue in
guard let self, self.realtimeRelayGeneration == relayGeneration else { return }
let issue = TalkRuntimeIssue(
code: .realtimeUnavailable,
message: relayIssue.message,
provider: relayIssue.provider,
model: relayIssue.model,
transport: relayIssue.transport,
phase: relayIssue.phase)
self.realtimeRelayStartIssue = issue
self.pendingRealtimeIssue = issue
self.gatewayTalkLastIssueText = issue.diagnosticSummary
File diff suppressed because it is too large Load Diff
@@ -2,7 +2,7 @@ import Foundation
import OpenClawKit
import OpenClawProtocol
import Testing
@testable import OpenClaw
@testable import OpenClawKit
@MainActor
private final class UnusedPCMStreamingAudioPlayer: PCMStreamingAudioPlaying {
@@ -15,6 +15,27 @@ private final class UnusedPCMStreamingAudioPlayer: PCMStreamingAudioPlaying {
}
}
@MainActor
private final class TestRealtimeTalkAudioCapture: RealtimeTalkAudioCapturing {
var suppressesInputDuringOutput = false
private(set) var isStarted = false
private(set) var startCount = 0
private(set) var stopCount = 0
func start(
targetSampleRate: Double,
onAudio: @escaping @Sendable (RealtimeTalkAudioFrame) -> Void) throws
{
self.isStarted = true
self.startCount += 1
}
func stop() {
self.isStarted = false
self.stopCount += 1
}
}
private actor RealtimeRelayStartupBarrier {
private var entered = false
private var enteredWaiter: CheckedContinuation<Void, Never>?
@@ -42,14 +63,14 @@ private actor RealtimeRelayStartupBarrier {
private struct RealtimeRelayStartupRequest: Sendable {
let method: String
let paramsJSON: String?
let params: [String: AnyCodable]?
}
private actor RealtimeRelayStartupRequestLog {
private var requests: [RealtimeRelayStartupRequest] = []
func record(method: String, paramsJSON: String?) {
self.requests.append(RealtimeRelayStartupRequest(method: method, paramsJSON: paramsJSON))
func record(method: String, params: [String: AnyCodable]?) {
self.requests.append(RealtimeRelayStartupRequest(method: method, params: params))
}
func snapshot() -> [RealtimeRelayStartupRequest] {
@@ -57,13 +78,74 @@ private actor RealtimeRelayStartupRequestLog {
}
}
private func unusedRealtimeRelayTransport() -> RealtimeTalkRelayTransport {
RealtimeTalkRelayTransport(
subscribeServerEvents: { _ in AsyncStream { $0.finish() } },
request: { _, _, _ in throw CancellationError() })
}
@MainActor
struct RealtimeTalkRelaySessionTests {
@Test func `transcript callback carries typed partial and final values`() async {
var transcripts: [RealtimeTalkTranscript] = []
let session = RealtimeTalkRelaySession(
transport: unusedRealtimeRelayTransport(),
options: .init(sessionKey: "main", provider: nil, model: nil, voice: nil),
audioCapture: TestRealtimeTalkAudioCapture(),
pcmPlayer: UnusedPCMStreamingAudioPlayer(),
onStatus: { _ in },
onSpeakingChanged: { _ in },
onTranscript: { transcripts.append($0) })
session._test_setRelaySessionId("relay-1")
for isFinal in [false, true] {
await session._test_handleGatewayEvent(EventFrame(
type: "event",
event: "talk.event",
payload: AnyCodable([
"relaySessionId": "relay-1",
"type": "transcript",
"role": "user",
"text": isFinal ? "hello" : "hel",
"final": isFinal,
]),
seq: nil,
stateversion: nil))
}
#expect(transcripts == [
RealtimeTalkTranscript(role: "user", text: "hel", isFinal: false),
RealtimeTalkTranscript(role: "user", text: "hello", isFinal: true),
])
}
@Test func `input pause and resume are idempotent and keep relay alive`() throws {
let audioCapture = TestRealtimeTalkAudioCapture()
let session = RealtimeTalkRelaySession(
transport: unusedRealtimeRelayTransport(),
options: .init(sessionKey: "main", provider: nil, model: nil, voice: nil),
audioCapture: audioCapture,
pcmPlayer: UnusedPCMStreamingAudioPlayer(),
onStatus: { _ in },
onSpeakingChanged: { _ in })
session._test_setRelaySessionId("relay-1")
try session.setInputPaused(true)
try session.setInputPaused(true)
try session.setInputPaused(false)
try session.setInputPaused(false)
#expect(audioCapture.stopCount == 2)
#expect(audioCapture.startCount == 1)
#expect(audioCapture.isStarted)
}
@Test func `output playback finish clears barge in start time`() {
var speakingStates: [Bool] = []
let session = RealtimeTalkRelaySession(
gateway: GatewayNodeSession(),
transport: unusedRealtimeRelayTransport(),
options: .init(sessionKey: "main", provider: nil, model: nil, voice: nil),
audioCapture: TestRealtimeTalkAudioCapture(),
pcmPlayer: UnusedPCMStreamingAudioPlayer(),
onStatus: { _ in },
onSpeakingChanged: { speakingStates.append($0) })
@@ -83,19 +165,19 @@ struct RealtimeTalkRelaySessionTests {
@Test func `playback mark is acknowledged after output finishes`() async throws {
let requests = RealtimeRelayStartupRequestLog()
let transport = RealtimeTalkRelaySession.StartupTransport(
let transport = RealtimeTalkRelayTransport(
subscribeServerEvents: { _ in AsyncStream { $0.finish() } },
request: { method, paramsJSON, _ in
await requests.record(method: method, paramsJSON: paramsJSON)
request: { method, params, _ in
await requests.record(method: method, params: params)
return Data("{\"ok\":true}".utf8)
})
let session = RealtimeTalkRelaySession(
gateway: GatewayNodeSession(),
transport: transport,
options: .init(sessionKey: "main", provider: "xai", model: nil, voice: nil),
audioCapture: TestRealtimeTalkAudioCapture(),
pcmPlayer: UnusedPCMStreamingAudioPlayer(),
onStatus: { _ in },
onSpeakingChanged: { _ in },
startupTransport: transport)
onSpeakingChanged: { _ in })
session._test_setRelaySessionId("relay-1")
session._test_markOutputAudioStarted(nowMs: 100)
@@ -122,17 +204,17 @@ struct RealtimeTalkRelaySessionTests {
#expect(recorded.count == 1)
let request = try #require(recorded.first)
#expect(request.method == "talk.session.acknowledgeMark")
let paramsData = try #require(request.paramsJSON?.data(using: .utf8))
let params = try #require(JSONSerialization.jsonObject(with: paramsData) as? [String: String])
#expect(params == ["sessionId": "relay-1", "markName": "audio-1"])
#expect(request.params?["sessionId"]?.stringValue == "relay-1")
#expect(request.params?["markName"]?.stringValue == "audio-1")
}
@Test func `close after classified error does not replace issue`() async {
var issues: [TalkRuntimeIssue] = []
var issues: [RealtimeTalkRelayIssue] = []
var statuses: [String] = []
let session = RealtimeTalkRelaySession(
gateway: GatewayNodeSession(),
transport: unusedRealtimeRelayTransport(),
options: .init(sessionKey: "main", provider: "openai", model: "gpt-realtime-2", voice: nil),
audioCapture: TestRealtimeTalkAudioCapture(),
pcmPlayer: UnusedPCMStreamingAudioPlayer(),
onStatus: { statuses.append($0) },
onIssue: { issues.append($0) },
@@ -165,14 +247,15 @@ struct RealtimeTalkRelaySessionTests {
seq: nil,
stateversion: nil))
#expect(issues.map(\.code) == [.realtimeUnavailable])
#expect(issues.map(\.code) == ["realtime_unavailable"])
#expect(statuses == ["OpenAI API key rejected with 401"])
}
@Test func `closed relay does not wait for startup ready`() async {
let session = RealtimeTalkRelaySession(
gateway: GatewayNodeSession(),
transport: unusedRealtimeRelayTransport(),
options: .init(sessionKey: "main", provider: "openai", model: "gpt-realtime-2", voice: nil),
audioCapture: TestRealtimeTalkAudioCapture(),
pcmPlayer: UnusedPCMStreamingAudioPlayer(),
onStatus: { _ in },
onSpeakingChanged: { _ in })
@@ -184,8 +267,9 @@ struct RealtimeTalkRelaySessionTests {
@Test func `startup ready wait covers gateway connect budget`() {
let session = RealtimeTalkRelaySession(
gateway: GatewayNodeSession(),
transport: unusedRealtimeRelayTransport(),
options: .init(sessionKey: "main", provider: "openai", model: "gpt-realtime-2", voice: nil),
audioCapture: TestRealtimeTalkAudioCapture(),
pcmPlayer: UnusedPCMStreamingAudioPlayer(),
onStatus: { _ in },
onSpeakingChanged: { _ in })
@@ -198,22 +282,22 @@ struct RealtimeTalkRelaySessionTests {
let requests = RealtimeRelayStartupRequestLog()
var statuses: [String] = []
var speakingStates: [Bool] = []
let transport = RealtimeTalkRelaySession.StartupTransport(
let transport = RealtimeTalkRelayTransport(
subscribeServerEvents: { _ in
await barrier.suspend()
return AsyncStream { $0.finish() }
},
request: { method, paramsJSON, _ in
await requests.record(method: method, paramsJSON: paramsJSON)
request: { method, params, _ in
await requests.record(method: method, params: params)
throw URLError(.badServerResponse)
})
let session = RealtimeTalkRelaySession(
gateway: GatewayNodeSession(),
transport: transport,
options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil),
audioCapture: TestRealtimeTalkAudioCapture(),
pcmPlayer: UnusedPCMStreamingAudioPlayer(),
onStatus: { statuses.append($0) },
onSpeakingChanged: { speakingStates.append($0) },
startupTransport: transport)
onSpeakingChanged: { speakingStates.append($0) })
let start = Task { @MainActor in try await session.start() }
await barrier.waitUntilEntered()
@@ -238,10 +322,10 @@ struct RealtimeTalkRelaySessionTests {
brain: AnyCodable("agent-consult"),
relaysessionid: "relay-1")
let resultData = try JSONEncoder().encode(result)
let transport = RealtimeTalkRelaySession.StartupTransport(
let transport = RealtimeTalkRelayTransport(
subscribeServerEvents: { _ in AsyncStream { $0.finish() } },
request: { method, paramsJSON, _ in
await requests.record(method: method, paramsJSON: paramsJSON)
request: { method, params, _ in
await requests.record(method: method, params: params)
if method == "talk.session.create" {
await barrier.suspend()
return resultData
@@ -249,12 +333,12 @@ struct RealtimeTalkRelaySessionTests {
return Data("{}".utf8)
})
let session = RealtimeTalkRelaySession(
gateway: GatewayNodeSession(),
transport: transport,
options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil),
audioCapture: TestRealtimeTalkAudioCapture(),
pcmPlayer: UnusedPCMStreamingAudioPlayer(),
onStatus: { statuses.append($0) },
onSpeakingChanged: { speakingStates.append($0) },
startupTransport: transport)
onSpeakingChanged: { speakingStates.append($0) })
let start = Task { @MainActor in try await session.start() }
await barrier.waitUntilEntered()
@@ -264,9 +348,7 @@ struct RealtimeTalkRelaySessionTests {
let recorded = await requests.snapshot()
#expect(recorded.map(\.method) == ["talk.session.create", "talk.session.close"])
let closeJSON = try #require(recorded.last?.paramsJSON?.data(using: .utf8))
let closeParams = try #require(JSONSerialization.jsonObject(with: closeJSON) as? [String: String])
#expect(closeParams["sessionId"] == "relay-1")
#expect(recorded.last?.params?["sessionId"]?.stringValue == "relay-1")
#expect(!statuses.contains("Waiting for realtime…"))
#expect(!speakingStates.contains(true))
}
@@ -275,10 +357,10 @@ struct RealtimeTalkRelaySessionTests {
let barrier = RealtimeRelayStartupBarrier()
let requests = RealtimeRelayStartupRequestLog()
var statuses: [String] = []
let transport = RealtimeTalkRelaySession.StartupTransport(
let transport = RealtimeTalkRelayTransport(
subscribeServerEvents: { _ in AsyncStream { $0.finish() } },
request: { method, paramsJSON, _ in
await requests.record(method: method, paramsJSON: paramsJSON)
request: { method, params, _ in
await requests.record(method: method, params: params)
if method == "talk.client.toolCall" {
await barrier.suspend()
return Data("{\"runId\":\"run-1\"}".utf8)
@@ -286,12 +368,12 @@ struct RealtimeTalkRelaySessionTests {
return Data("{}".utf8)
})
let session = RealtimeTalkRelaySession(
gateway: GatewayNodeSession(),
transport: transport,
options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil),
audioCapture: TestRealtimeTalkAudioCapture(),
pcmPlayer: UnusedPCMStreamingAudioPlayer(),
onStatus: { statuses.append($0) },
onSpeakingChanged: { _ in },
startupTransport: transport)
onSpeakingChanged: { _ in })
session._test_setRelaySessionId("relay-1")
let handling = Task { @MainActor in
await session._test_handleGatewayEvent(EventFrame(
@@ -322,19 +404,19 @@ struct RealtimeTalkRelaySessionTests {
@Test func `stop cancels buffered microphone audio before dispatch`() async throws {
let requests = RealtimeRelayStartupRequestLog()
let transport = RealtimeTalkRelaySession.StartupTransport(
let transport = RealtimeTalkRelayTransport(
subscribeServerEvents: { _ in AsyncStream { $0.finish() } },
request: { method, paramsJSON, _ in
await requests.record(method: method, paramsJSON: paramsJSON)
request: { method, params, _ in
await requests.record(method: method, params: params)
return Data("{\"ok\":true}".utf8)
})
let session = RealtimeTalkRelaySession(
gateway: GatewayNodeSession(),
transport: transport,
options: .init(sessionKey: "main", provider: "openai", model: nil, voice: nil),
audioCapture: TestRealtimeTalkAudioCapture(),
pcmPlayer: UnusedPCMStreamingAudioPlayer(),
onStatus: { _ in },
onSpeakingChanged: { _ in },
startupTransport: transport)
onSpeakingChanged: { _ in })
session._test_prepareAudioSender(relaySessionId: "relay-1")
let send = try #require(session._test_enqueueMicrophoneFrame(Data([0x01, 0x02])))