@@ -6,9 +6,8 @@ import LiveKit
66
77/// A single-use conversation session, created by and owned by `ConversationClient`.
88///
9- /// Manages the lifecycle of one conversation: network layer
10- /// (`WebRTCConnectionManager`|`WebSocketConnectionManager`), protocol parser
11- /// (`EventParser`), and audio device setup.
9+ /// Owns one conversation's transport (`WebRTCConnectionManager` or
10+ /// `WebSocketConnectionManager`) and audio device setup.
1211@MainActor
1312final class Conversation : ObservableObject {
1413 // MARK: - State
@@ -17,8 +16,6 @@ final class Conversation: ObservableObject {
1716 @Published var chatHistory : [ any ChatHistoryItem ] = [ ]
1817 @Published var agentState : AgentState = . listening
1918
20- private var chatHistoryReconciler = ChatHistoryReconciler ( )
21-
2219 /// Stream of client tool calls that need to be executed by the app
2320 @Published var pendingToolCalls : [ ClientToolCallEvent ] = [ ]
2421
@@ -31,38 +28,35 @@ final class Conversation: ObservableObject {
3128 /// Current MCP connection status for all integrations
3229 @Published var mcpConnectionStatus : MCPConnectionStatusEvent ?
3330
34- /// Pending mute state to apply after connection completes.
35- /// Allows setting mute state during connection phase.
36- private var pendingMuteState : Bool ?
31+ private var chatHistoryReconciler = ChatHistoryReconciler ( )
32+
33+ /// Agent state manager for event-based state tracking
34+ private var agentStateManager : AgentStateManager ?
3735
3836 /// Audio device management
3937 private var audioManager : ConversationAudioManager ?
4038
39+ /// Pending mute state to apply after connection completes.
40+ /// Allows setting mute state during connection phase.
41+ private var pendingMuteState : Bool ?
42+
4143 /// Externally registered audio observers. Kept attached across track swaps.
4244 let agentObserverRegistry = AudioObserverRegistry ( )
4345 let micObserverRegistry = AudioObserverRegistry ( )
4446
45- /// Agent state manager for event-based state tracking
46- var agentStateManager : AgentStateManager ?
47+ private let dependencyProvider : any ConversationDependencyProvider
48+ private var activeConnectionManager : ( any ConnectionManaging ) ?
49+ private var activeWebRTCConnectionManager : ( any WebRTCConnectionManaging ) ? {
50+ activeConnectionManager as? any WebRTCConnectionManaging
51+ }
4752
48- /// Forward a signal to the event-based state manager, or fall back to directly setting `agentState`.
49- func applyStateSignal( _ signal: AgentStateSignal , fallback: AgentState ) {
50- if let manager = agentStateManager {
51- manager. processSignal ( signal)
52- } else {
53- agentState = fallback
54- }
53+ /// Internal LiveKit tracks used to attach ``ConversationAudioObserver``s.
54+ private var inputTrack : ( any AudioTrackProtocol ) ? {
55+ activeWebRTCConnectionManager? . inputTrack
5556 }
5657
57- func handleRemoteSpeakingUpdate( isSpeaking: Bool ) {
58- if let manager = agentStateManager {
59- manager. processSignal ( isSpeaking ? . agentStartedSpeaking : . agentStoppedSpeaking)
60- } else if isSpeaking {
61- speakingTimer? . cancel ( )
62- agentState = . speaking
63- } else {
64- scheduleBackToListening ( delay: 1.0 )
65- }
58+ private var agentAudioTrack : ( any AudioTrackProtocol ) ? {
59+ activeWebRTCConnectionManager? . agentAudioTrack
6660 }
6761
6862 /// Internal logger, accessible from nonisolated contexts.
@@ -71,14 +65,10 @@ final class Conversation: ObservableObject {
7165 /// Context for logging (e.g. agentId)
7266 private var activeContext : [ String : String ] ?
7367
74- /// Internal LiveKit tracks used to attach ``ConversationAudioObserver``s.
75- var inputTrack : ( any AudioTrackProtocol ) ? {
76- activeWebRTCConnectionManager? . inputTrack
77- }
78-
79- var agentAudioTrack : ( any AudioTrackProtocol ) ? {
80- activeWebRTCConnectionManager? . agentAudioTrack
81- }
68+ let config : ConversationConfig
69+ let callbacks : ConversationCallbacks
70+ private var speakingTimer : Task < Void , Never > ?
71+ private var isTearingDown = false
8272
8373 // MARK: - Init
8474
@@ -95,16 +85,6 @@ final class Conversation: ObservableObject {
9585 logger = dependencyProvider. logger
9686 }
9787
98- private func setupAgentStateManager( ) {
99- guard let configuration = config. agentStateConfiguration else { return }
100- let manager = AgentStateManager ( configuration: configuration)
101- manager. onStateChange = { [ weak self] state in
102- self ? . agentState = state
103- self ? . callbacks. onAgentStateChange ? ( state)
104- }
105- agentStateManager = manager
106- }
107-
10888 // MARK: - API
10989
11090 func startVoiceConversation( _ auth: ConversationAuth . Voice ) async throws -> ConversationStartResult {
@@ -156,6 +136,18 @@ final class Conversation: ObservableObject {
156136 return try await setConnected ( result)
157137 }
158138
139+ func startTextOnlyConversation( _ auth: ConversationAuth . TextOnly ) async throws -> ConversationStartResult {
140+ let manager = dependencyProvider. webSocketConnectionManager
141+ let result = try await start ( agentId: auth. agentId, isTextOnly: true , using: manager) { config in
142+ try await manager. connect (
143+ auth: auth,
144+ config: config,
145+ onStartupStateChange: { [ weak self] in self ? . updateStartupStage ( $0) }
146+ )
147+ }
148+ return try await setConnected ( result)
149+ }
150+
159151 // MARK: - Audio observers
160152
161153 /// Register an observer for the agent's decoded output audio.
@@ -187,18 +179,6 @@ final class Conversation: ObservableObject {
187179 micObserverRegistry. attach ( to: inputTrack)
188180 }
189181
190- func startTextOnlyConversation( _ auth: ConversationAuth . TextOnly ) async throws -> ConversationStartResult {
191- let manager = dependencyProvider. webSocketConnectionManager
192- let result = try await start ( agentId: auth. agentId, isTextOnly: true , using: manager) { config in
193- try await manager. connect (
194- auth: auth,
195- config: config,
196- onStartupStateChange: { [ weak self] in self ? . updateStartupStage ( $0) }
197- )
198- }
199- return try await setConnected ( result)
200- }
201-
202182 /// End and clean up.
203183 /// Can be called during connection phase to cancel, or during connected conversation to end.
204184 func endConversation( reason: EndReason = . userEnded) async {
@@ -266,7 +246,7 @@ final class Conversation: ObservableObject {
266246 try await publish ( event)
267247 }
268248
269- /// Contextual update to agent (system prompt-ish) .
249+ /// Inject context for the agent without interrupting or adding a user-visible message .
270250 func updateContext( _ context: String ) async throws {
271251 guard state. isConnected else { throw ConversationError . notConnected }
272252 let event = OutgoingEvent . contextualUpdate ( ContextualUpdateEvent ( text: context) )
@@ -307,23 +287,6 @@ final class Conversation: ObservableObject {
307287
308288 // MARK: - Private
309289
310- private let dependencyProvider : any ConversationDependencyProvider
311- private var activeConnectionManager : ( any ConnectionManaging ) ?
312- private var activeWebRTCConnectionManager : ( any WebRTCConnectionManaging ) ? {
313- activeConnectionManager as? any WebRTCConnectionManaging
314- }
315-
316- let config : ConversationConfig
317- let callbacks : ConversationCallbacks
318-
319- var speakingTimer : Task < Void , Never > ?
320- private var isTearingDown = false
321-
322- private func updateStartupStage( _ stage: ConversationStartupState ) {
323- guard state. isConnecting, state != . connecting( stage) else { return }
324- state = . connecting( stage)
325- }
326-
327290 private func start(
328291 agentId: String ,
329292 isTextOnly: Bool ,
@@ -381,6 +344,41 @@ final class Conversation: ObservableObject {
381344 return result
382345 }
383346
347+ private func updateStartupStage( _ stage: ConversationStartupState ) {
348+ guard state. isConnecting, state != . connecting( stage) else { return }
349+ state = . connecting( stage)
350+ }
351+
352+ private func setupAgentStateManager( ) {
353+ guard let configuration = config. agentStateConfiguration else { return }
354+ let manager = AgentStateManager ( configuration: configuration)
355+ manager. onStateChange = { [ weak self] state in
356+ self ? . agentState = state
357+ self ? . callbacks. onAgentStateChange ? ( state)
358+ }
359+ agentStateManager = manager
360+ }
361+
362+ /// Forward a signal to the event-based state manager, or fall back to directly setting `agentState`.
363+ private func applyStateSignal( _ signal: AgentStateSignal , fallback: AgentState ) {
364+ if let manager = agentStateManager {
365+ manager. processSignal ( signal)
366+ } else {
367+ agentState = fallback
368+ }
369+ }
370+
371+ private func handleRemoteSpeakingUpdate( isSpeaking: Bool ) {
372+ if let manager = agentStateManager {
373+ manager. processSignal ( isSpeaking ? . agentStartedSpeaking : . agentStoppedSpeaking)
374+ } else if isSpeaking {
375+ speakingTimer? . cancel ( )
376+ agentState = . speaking
377+ } else {
378+ scheduleBackToListening ( delay: 1.0 )
379+ }
380+ }
381+
384382 private func handleIncomingEvent( _ event: IncomingEvent , from connectionManagerID: ObjectIdentifier ) async {
385383 guard let activeConnectionManager,
386384 ObjectIdentifier ( activeConnectionManager) == connectionManagerID,
@@ -448,7 +446,7 @@ final class Conversation: ObservableObject {
448446 }
449447 }
450448
451- func publish( _ event: OutgoingEvent ) async throws {
449+ private func publish( _ event: OutgoingEvent ) async throws {
452450 guard let connectionManager = activeConnectionManager else {
453451 throw ConversationError . notConnected
454452 }
0 commit comments