@@ -24,6 +24,7 @@ enum WebRTCConnectionManagerError: Error {
2424///
2525/// LiveKit room/track types stay internal to this manager; consumers observe
2626/// audio via ``ConversationAudioObserver`` rather than LiveKit track APIs.
27+ @MainActor
2728final class WebRTCConnectionManager : WebRTCConnectionManaging {
2829 /// Fired when the remote agent leaves, the room disconnects, or all remote participants are gone.
2930 var onDisconnected : ( ( ) async -> Void ) ?
@@ -62,6 +63,7 @@ final class WebRTCConnectionManager: WebRTCConnectionManaging {
6263 private var eventDelegate : LiveKitRoomEventDelegate ?
6364 private var readinessDelegate : LiveKitReadinessDelegate ?
6465 private var initiationMetadataWaiter : ConversationInitiationMetadataWaiter ?
66+ private var dataTask : Task < Void , Never > ?
6567
6668 private static let reliableDataPublishOptions = DataPublishOptions ( reliable: true )
6769
@@ -185,7 +187,7 @@ final class WebRTCConnectionManager: WebRTCConnectionManaging {
185187 }
186188 let first = await group. next ( ) !
187189 if case . timedOut = first {
188- await delegate. release ( )
190+ delegate. release ( )
189191 }
190192 group. cancelAll ( )
191193 return first
@@ -231,23 +233,42 @@ final class WebRTCConnectionManager: WebRTCConnectionManaging {
231233 networkConfiguration: WebRTCConfiguration ,
232234 metadataWaiter: ConversationInitiationMetadataWaiter
233235 ) async throws {
234- await readinessDelegate? . release ( )
236+ dataTask? . cancel ( )
237+ readinessDelegate? . release ( )
235238
236- let readinessDelegate = await LiveKitReadinessDelegate ( logger: logger)
239+ let readinessDelegate = LiveKitReadinessDelegate ( logger: logger)
237240 self . readinessDelegate = readinessDelegate
238241
239242 let logger = logger
243+ let ( dataStream, dataContinuation) = AsyncStream . makeStream ( of: Data . self)
240244 let eventDelegate = LiveKitRoomEventDelegate (
241- onData: { [ weak self] data in
242- self ? . handleIncomingData ( data, metadataWaiter: metadataWaiter, logger: logger)
245+ onData: { dataContinuation. yield ( $0) } ,
246+ onRemoteSpeaking: { [ weak self] isSpeaking in
247+ Task { @MainActor [ weak self] in
248+ self ? . onRemoteSpeakingChanged ? ( isSpeaking)
249+ }
250+ } ,
251+ onRemoteDisconnect: { [ weak self] in
252+ await self ? . handleRemoteDisconnect ( )
243253 } ,
244- onRemoteSpeaking: { [ weak self] isSpeaking in self ? . onRemoteSpeakingChanged ? ( isSpeaking) } ,
245- onRemoteDisconnect: { [ weak self] in await self ? . onDisconnected ? ( ) } ,
246254 onTracksChanged: { [ weak self] in
247255 Task { @MainActor in self ? . onTracksChanged ? ( ) }
248256 }
249257 )
250258 self . eventDelegate = eventDelegate
259+ // One consumer preserves packet order while parsing off the main actor.
260+ dataTask = Task { [ weak self] in
261+ for await data in dataStream {
262+ await Self . handleIncomingData (
263+ data,
264+ metadataWaiter: metadataWaiter,
265+ logger: logger,
266+ onEvent: { [ weak self] event in
267+ self ? . onEventReceived ? ( event)
268+ }
269+ )
270+ }
271+ }
251272
252273 let room = Room ( roomOptions: RoomOptions ( singlePeerConnection: true ) )
253274 self . room = room
@@ -299,21 +320,27 @@ final class WebRTCConnectionManager: WebRTCConnectionManaging {
299320
300321 await initiationMetadataWaiter? . cancel ( )
301322 initiationMetadataWaiter = nil
302- await readinessDelegate? . release ( )
323+ dataTask? . cancel ( )
324+ dataTask = nil
325+ readinessDelegate? . release ( )
303326 readinessDelegate = nil
304327
305328 await room? . disconnect ( )
306329 room = nil
307330 eventDelegate = nil
308331 }
309332
333+ private func handleRemoteDisconnect( ) async {
334+ await onDisconnected ? ( )
335+ }
336+
310337 // MARK: – Private helpers
311338
312339 /// Run one timed startup phase: record its duration into `metrics[keyPath:]`,
313340 /// let `CancellationError` propagate unwrapped, and wrap any other error as
314341 /// `ConversationError` (stamping `total`).
315342 @MainActor
316- private func runPhase< T> (
343+ private func runPhase< T: Sendable > (
317344 timing keyPath: WritableKeyPath < ConversationStartupMetrics , TimeInterval ? > ,
318345 metrics: inout ConversationStartupMetrics ,
319346 startTime: Date ,
0 commit comments