@@ -39,6 +39,8 @@ class DataChannelPair: NSObject, @unchecked Sendable, Loggable {
3939 private struct State {
4040 var lossy : LKRTCDataChannel ?
4141 var reliable : LKRTCDataChannel ?
42+ var reliableDataSequence : UInt32 = 1
43+ var reliableReceivedState : TTLDictionary < String , UInt32 > = TTLDictionary ( ttl: reliableReceivedStateTTL)
4244
4345 var isOpen : Bool {
4446 guard let lossy, let reliable else { return false }
@@ -54,13 +56,59 @@ class DataChannelPair: NSObject, @unchecked Sendable, Loggable {
5456 case lossy, reliable
5557 }
5658
57- private struct BufferingState {
58- var queue : Deque < PublishDataRequest > = [ ]
59- var amount : UInt64 = 0
59+ private struct SendBuffer {
60+ private var queue : Deque < PublishDataRequest > = [ ]
61+ var rtcAmount : UInt64 = 0
62+
63+ mutating func enqueue( _ request: PublishDataRequest ) {
64+ queue. append ( request)
65+ }
66+
67+ @discardableResult
68+ mutating func dequeue( ) -> PublishDataRequest ? {
69+ guard !queue. isEmpty else { return nil }
70+ return queue. removeFirst ( )
71+ }
72+
73+ func canSend( threshold: UInt64 ) -> Bool {
74+ rtcAmount <= threshold
75+ }
76+ }
77+
78+ private struct RetryBuffer {
79+ private var queue : Deque < PublishDataRequest > = [ ]
80+ private var currentAmount : UInt64 = 0
81+ private let minAmount : UInt64
82+
83+ init ( minAmount: UInt64 ) {
84+ self . minAmount = minAmount
85+ }
86+
87+ func peek( ) -> PublishDataRequest ? { queue. first }
88+
89+ mutating func enqueue( _ request: PublishDataRequest ) {
90+ queue. append ( request)
91+ currentAmount += UInt64 ( request. data. data. count)
92+ }
93+
94+ @discardableResult
95+ mutating func dequeue( ) -> PublishDataRequest ? {
96+ guard !queue. isEmpty else { return nil }
97+ let first = queue. removeFirst ( )
98+ currentAmount -= UInt64 ( first. data. data. count)
99+ return first
100+ }
101+
102+ mutating func trim( toAmount: UInt64 ) {
103+ while currentAmount > toAmount + minAmount {
104+ dequeue ( )
105+ }
106+ }
60107 }
61108
62109 private struct PublishDataRequest : Sendable {
63110 let data : LKRTCDataBuffer
111+ let sequence : UInt32
64112 let continuation : CheckedContinuation < Void , any Error > ?
65113 }
66114
@@ -70,41 +118,60 @@ class DataChannelPair: NSObject, @unchecked Sendable, Loggable {
70118
71119 enum Detail : Sendable {
72120 case publishData( PublishDataRequest )
121+ case publishedData( PublishDataRequest )
73122 case bufferedAmountChanged( UInt64 )
123+ case retryRequested( UInt32 )
74124 }
75125 }
76126
127+ // MARK: - Event handling
128+
77129 private func handleEvents(
78130 events: AsyncStream < ChannelEvent >
79131 ) async {
80- var lossyBuffering = BufferingState ( )
81- var reliableBuffering = BufferingState ( )
132+ var lossyBuffer = SendBuffer ( )
133+ var reliableBuffer = SendBuffer ( )
134+
135+ var reliableRetryBuffer = RetryBuffer ( minAmount: Self . reliableRetryAmount)
82136
83137 for await event in events {
84138 switch event. detail {
85139 case let . publishData( request) :
86140 switch event. channelKind {
87- case . lossy: lossyBuffering. queue. append ( request)
88- case . reliable: reliableBuffering. queue. append ( request)
141+ case . lossy: lossyBuffer. enqueue ( request)
142+ case . reliable: reliableBuffer. enqueue ( request)
143+ }
144+ case let . publishedData( request) :
145+ switch event. channelKind {
146+ case . lossy: ( )
147+ case . reliable: reliableRetryBuffer. enqueue ( request)
89148 }
90149 case let . bufferedAmountChanged( amount) :
91150 switch event. channelKind {
92- case . lossy: updateBufferingState ( state: & lossyBuffering, newAmount: amount)
93- case . reliable: updateBufferingState ( state: & reliableBuffering, newAmount: amount)
151+ case . lossy:
152+ updateTarget ( buffer: & lossyBuffer, newAmount: amount)
153+ case . reliable:
154+ updateTarget ( buffer: & reliableBuffer, newAmount: amount)
155+ reliableRetryBuffer. trim ( toAmount: amount)
156+ }
157+ case let . retryRequested( lastSeq) :
158+ switch event. channelKind {
159+ case . lossy: ( )
160+ case . reliable: retry ( buffer: & reliableRetryBuffer, from: lastSeq)
94161 }
95162 }
96163
97164 switch event. channelKind {
98165 case . lossy:
99166 processSendQueue (
100167 threshold: Self . lossyLowThreshold,
101- state : & lossyBuffering ,
168+ buffer : & lossyBuffer ,
102169 kind: . lossy
103170 )
104171 case . reliable:
105172 processSendQueue (
106173 threshold: Self . reliableLowThreshold,
107- state : & reliableBuffering ,
174+ buffer : & reliableBuffer ,
108175 kind: . reliable
109176 )
110177 }
@@ -120,14 +187,11 @@ class DataChannelPair: NSObject, @unchecked Sendable, Loggable {
120187
121188 private func processSendQueue(
122189 threshold: UInt64 ,
123- state : inout BufferingState ,
190+ buffer : inout SendBuffer ,
124191 kind: ChannelKind
125192 ) {
126- while state. amount <= threshold {
127- guard !state. queue. isEmpty else { break }
128- let request = state. queue. removeFirst ( )
129-
130- state. amount += UInt64 ( request. data. data. count)
193+ while buffer. canSend ( threshold: threshold) , let request = buffer. dequeue ( ) {
194+ buffer. rtcAmount += UInt64 ( request. data. data. count)
131195
132196 guard let channel = channel ( for: kind) else {
133197 request. continuation? . resume (
@@ -142,21 +206,43 @@ class DataChannelPair: NSObject, @unchecked Sendable, Loggable {
142206 return
143207 }
144208 request. continuation? . resume ( )
209+
210+ let event = ChannelEvent ( channelKind: kind, detail: . publishedData( request) )
211+ _state. eventContinuation? . yield ( event)
145212 }
146213 }
147214
148- private func updateBufferingState(
149- state: inout BufferingState ,
215+ // MARK: - Cache
216+
217+ private func updateTarget(
218+ buffer: inout SendBuffer ,
150219 newAmount: UInt64
151220 ) {
152- guard state . amount >= newAmount else {
221+ guard buffer . rtcAmount >= newAmount else {
153222 log ( " Unexpected buffer size detected " , . error)
154- state . amount = 0
223+ buffer . rtcAmount = 0
155224 return
156225 }
157- state . amount -= newAmount
226+ buffer . rtcAmount -= newAmount
158227 }
159228
229+ private func retry(
230+ buffer: inout RetryBuffer ,
231+ from lastSeq: UInt32
232+ ) {
233+ if let first = buffer. peek ( ) , first. sequence > lastSeq + 1 {
234+ log ( " Wrong packet sequence while retrying: \( first. sequence) > \( lastSeq + 1 ) , \( first. sequence - lastSeq - 1 ) packets missing " , . warning)
235+ }
236+ while let request = buffer. dequeue ( ) {
237+ if request. sequence > lastSeq {
238+ let event = ChannelEvent ( channelKind: . reliable, detail: . publishData( PublishDataRequest ( data: request. data, sequence: request. sequence, continuation: nil ) ) )
239+ _state. eventContinuation? . yield ( event)
240+ }
241+ }
242+ }
243+
244+ // MARK: - Init
245+
160246 init ( delegate: DataChannelDelegate ? = nil ,
161247 lossyChannel: LKRTCDataChannel ? = nil ,
162248 reliableChannel: LKRTCDataChannel ? = nil )
@@ -207,6 +293,8 @@ class DataChannelPair: NSObject, @unchecked Sendable, Loggable {
207293 let ( lossy, reliable) = _state. mutate {
208294 let result = ( $0. lossy, $0. reliable)
209295 $0. reliable = nil
296+ $0. reliableDataSequence = 1
297+ $0. reliableReceivedState. removeAll ( )
210298 $0. lossy = nil
211299 return result
212300 }
@@ -217,6 +305,8 @@ class DataChannelPair: NSObject, @unchecked Sendable, Loggable {
217305 openCompleter. reset ( )
218306 }
219307
308+ // MARK: - Send
309+
220310 func send( userPacket: Livekit_UserPacket , kind: Livekit_DataPacket . Kind ) async throws {
221311 try await send ( dataPacket: . with {
222312 $0. kind = kind // TODO: field is deprecated
@@ -225,12 +315,14 @@ class DataChannelPair: NSObject, @unchecked Sendable, Loggable {
225315 }
226316
227317 func send( dataPacket packet: Livekit_DataPacket ) async throws {
318+ let packet = withSequence ( packet)
228319 let serializedData = try packet. serializedData ( )
229320 let rtcData = RTC . createDataBuffer ( data: serializedData)
230321
231322 try await withCheckedThrowingContinuation { continuation in
232323 let request = PublishDataRequest (
233324 data: rtcData,
325+ sequence: packet. sequence,
234326 continuation: continuation
235327 )
236328 let event = ChannelEvent (
@@ -241,15 +333,48 @@ class DataChannelPair: NSObject, @unchecked Sendable, Loggable {
241333 }
242334 }
243335
336+ private func withSequence( _ packet: Livekit_DataPacket ) -> Livekit_DataPacket {
337+ guard packet. kind == . reliable, packet. sequence == 0 else { return packet }
338+ var packet = packet
339+ _state. mutate {
340+ packet. sequence = $0. reliableDataSequence
341+ $0. reliableDataSequence += 1
342+ }
343+ return packet
344+ }
345+
346+ func retryReliable( lastSequence: UInt32 ) {
347+ let event = ChannelEvent ( channelKind: . reliable, detail: . retryRequested( lastSequence) )
348+ _state. eventContinuation? . yield ( event)
349+ }
350+
351+ // MARK: - Sync state
352+
244353 func infos( ) -> [ Livekit_DataChannelInfo ] {
245354 _state. read { [ $0. lossy, $0. reliable] }
246355 . compactMap { $0 }
247356 . map { $0. toLKInfoType ( ) }
248357 }
249358
359+ func receiveStates( ) -> [ Livekit_DataChannelReceiveState ] {
360+ _state. reliableReceivedState. map { sid, seq in
361+ Livekit_DataChannelReceiveState . with {
362+ $0. publisherSid = sid
363+ $0. lastSeq = seq
364+ }
365+ }
366+ }
367+
368+ // MARK: - Constants
369+
250370 private static let reliableLowThreshold : UInt64 = 2 * 1024 * 1024 // 2 MB
251371 private static let lossyLowThreshold : UInt64 = reliableLowThreshold
252372
373+ // If rtc drains its buffer to 0, keep at least this amount of data for retry.
374+ // Should be >= the full backpressure amount to avoid losing packets.
375+ private static let reliableRetryAmount : UInt64 = . init( Double ( reliableLowThreshold) * 1.25 )
376+ private static let reliableReceivedStateTTL : TimeInterval = 30
377+
253378 deinit {
254379 _state. eventContinuation? . finish ( )
255380 }
@@ -272,18 +397,30 @@ extension DataChannelPair: LKRTCDataChannelDelegate {
272397 }
273398 }
274399
275- func dataChannel( _: LKRTCDataChannel , didReceiveMessageWith buffer: LKRTCDataBuffer ) {
400+ func dataChannel( _ dataChannel : LKRTCDataChannel , didReceiveMessageWith buffer: LKRTCDataBuffer ) {
276401 guard let dataPacket = try ? Livekit_DataPacket ( serializedBytes: buffer. data) else {
277402 log ( " Could not decode data message " , . error)
278403 return
279404 }
280405
406+ if dataChannel. kind == . reliable, dataPacket. sequence > 0 , !dataPacket. participantSid. isEmpty {
407+ if let lastSeq = _state. reliableReceivedState [ dataPacket. participantSid] , dataPacket. sequence <= lastSeq {
408+ log ( " Ignoring duplicate/out-of-order reliable data message " , . warning)
409+ return
410+ }
411+ _state. mutate {
412+ $0. reliableReceivedState [ dataPacket. participantSid] = dataPacket. sequence
413+ }
414+ }
415+
281416 delegates. notify {
282417 $0. dataChannel ( self , didReceiveDataPacket: dataPacket)
283418 }
284419 }
285420}
286421
422+ // MARK: - Extensions
423+
287424private extension DataChannelPair . ChannelKind {
288425 init ( _ packetKind: Livekit_DataPacket . Kind ) {
289426 guard case . lossy = packetKind else {
0 commit comments