From 81b1edb4153c7fcccdcbf6797e96070dc4b76162 Mon Sep 17 00:00:00 2001 From: Jacob Gelman <3182119+ladvoc@users.noreply.github.com> Date: Wed, 2 Sep 2026 16:27:13 -0700 Subject: [PATCH 1/7] Create the data track data channel --- engine.go | 23 ++++++++++++++++++++--- 1 file changed, 20 insertions(+), 3 deletions(-) diff --git a/engine.go b/engine.go index 6356e687..4a34770b 100644 --- a/engine.go +++ b/engine.go @@ -92,8 +92,9 @@ var ( // ------------------------------------------- const ( - reliableDataChannelName = "_reliable" - lossyDataChannelName = "_lossy" + reliableDataChannelName = "_reliable" + lossyDataChannelName = "_lossy" + dataTrackDataChannelName = "_data_track" maxReconnectCount = 10 initialReconnectInterval = 300 * time.Millisecond @@ -127,6 +128,8 @@ type RTCEngine struct { lossyDC *webrtc.DataChannel reliableDCSub *webrtc.DataChannel lossyDCSub *webrtc.DataChannel + dataTrackDC *webrtc.DataChannel + dataTrackDCSub *webrtc.DataChannel reliableMsgLock sync.Mutex reliableMsgSeq uint32 @@ -557,6 +560,15 @@ func (e *RTCEngine) createPublisherPCLocked(configuration webrtc.Configuration) return err } e.reliableDC.OnMessage(e.handleDataPacket) + + e.dataTrackDC, err = e.publisher.pc.CreateDataChannel(dataTrackDataChannelName, &webrtc.DataChannelInit{ + Ordered: &falseVal, + MaxRetransmits: new(uint16), + }) + if err != nil { + e.dclock.Unlock() + return err + } e.dclock.Unlock() return nil @@ -643,6 +655,9 @@ func (e *RTCEngine) createSubscriberPCLocked(configuration webrtc.Configuration) e.reliableDCSub = c } else if c.Label() == lossyDataChannelName { e.lossyDCSub = c + } else if c.Label() == dataTrackDataChannelName { + e.dataTrackDCSub = c + return } else { return } @@ -750,7 +765,9 @@ func (e *RTCEngine) ensurePublisherConnected(ensureDataReady bool) error { func (e *RTCEngine) dataPubChannelReady() bool { e.dclock.RLock() defer e.dclock.RUnlock() - return e.reliableDC.ReadyState() == webrtc.DataChannelStateOpen && e.lossyDC.ReadyState() == webrtc.DataChannelStateOpen + return e.reliableDC.ReadyState() == webrtc.DataChannelStateOpen && + e.lossyDC.ReadyState() == webrtc.DataChannelStateOpen && + e.dataTrackDC.ReadyState() == webrtc.DataChannelStateOpen } func (e *RTCEngine) RegisterTrackPublishedListener(cid string, c chan *livekit.TrackPublishedResponse) { From debc9f3cbac0418fdac7f5ca813d49d600f94a1a Mon Sep 17 00:00:00 2001 From: Jacob Gelman <3182119+ladvoc@users.noreply.github.com> Date: Wed, 2 Sep 2026 16:27:16 -0700 Subject: [PATCH 2/7] Forward data track packets to the engine handler --- engine.go | 10 ++++++++++ room.go | 3 +++ 2 files changed, 13 insertions(+) diff --git a/engine.go b/engine.go index 4a34770b..74e68a05 100644 --- a/engine.go +++ b/engine.go @@ -80,6 +80,7 @@ type engineHandler interface { OnPublishDataTrackResponse(publishDataTrackResponse *livekit.PublishDataTrackResponse) OnUnpublishDataTrackResponse(unpublishDataTrackResponse *livekit.UnpublishDataTrackResponse) OnDataTrackSubscriberHandles(dataTrackSubscriberHandles *livekit.DataTrackSubscriberHandles) + OnDataTrackPacket(data []byte) } // ------------------------------------------- @@ -569,6 +570,7 @@ func (e *RTCEngine) createPublisherPCLocked(configuration webrtc.Configuration) e.dclock.Unlock() return err } + e.dataTrackDC.OnMessage(e.handleDataTrackPacket) e.dclock.Unlock() return nil @@ -657,6 +659,7 @@ func (e *RTCEngine) createSubscriberPCLocked(configuration webrtc.Configuration) e.lossyDCSub = c } else if c.Label() == dataTrackDataChannelName { e.dataTrackDCSub = c + c.OnMessage(e.handleDataTrackPacket) return } else { return @@ -926,6 +929,13 @@ func (e *RTCEngine) handleDataPacket(msg webrtc.DataChannelMessage) { } } +func (e *RTCEngine) handleDataTrackPacket(msg webrtc.DataChannelMessage) { + if msg.IsString { + return + } + e.engineHandler.OnDataTrackPacket(msg.Data) +} + func (e *RTCEngine) readDataPacket(msg webrtc.DataChannelMessage) (*livekit.DataPacket, error) { dataPacket := &livekit.DataPacket{} if msg.IsString { diff --git a/room.go b/room.go index ca5cc6d0..4cddd5e5 100644 --- a/room.go +++ b/room.go @@ -1415,6 +1415,9 @@ func (r *Room) OnUnpublishDataTrackResponse(unpublishDataTrackResponse *livekit. func (r *Room) OnDataTrackSubscriberHandles(dataTrackSubscriberHandles *livekit.DataTrackSubscriberHandles) { } +func (r *Room) OnDataTrackPacket(data []byte) { +} + func (r *Room) OnStreamHeader(streamHeader *livekit.DataStream_Header, participantIdentity string) { switch header := streamHeader.ContentHeader.(type) { case *livekit.DataStream_Header_TextHeader: From 7f908f288f18fd3f2ec0cec423022c42d21ceb47 Mon Sep 17 00:00:00 2001 From: Jacob Gelman <3182119+ladvoc@users.noreply.github.com> Date: Wed, 2 Sep 2026 16:27:27 -0700 Subject: [PATCH 3/7] Add data track frame sender --- datatracksender.go | 125 ++++++++++++++++++++++++++++++++++++++++ datatracksender_test.go | 56 ++++++++++++++++++ engine.go | 18 ++++++ 3 files changed, 199 insertions(+) create mode 100644 datatracksender.go create mode 100644 datatracksender_test.go diff --git a/datatracksender.go b/datatracksender.go new file mode 100644 index 00000000..ff25115d --- /dev/null +++ b/datatracksender.go @@ -0,0 +1,125 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package lksdk + +import ( + "sync" + + protoLogger "github.com/livekit/protocol/logger" + "github.com/pion/webrtc/v4" +) + +const dataTrackBufferedAmountLowThreshold = 8 * 1024 + +// dataTrackFramePackets is one frame serialized into data track packets. +type dataTrackFramePackets [][]byte + +// dataTrackSender paces frames onto the data track channel, keeping only the freshest frame +// while the channel's send buffer is above the low threshold. +type dataTrackSender struct { + dc func() *webrtc.DataChannel + log protoLogger.Logger + + lock sync.Mutex + frame dataTrackFramePackets + notify chan struct{} + done chan struct{} + once sync.Once +} + +func newDataTrackSender(dc func() *webrtc.DataChannel, log protoLogger.Logger) *dataTrackSender { + return &dataTrackSender{ + dc: dc, + log: log, + notify: make(chan struct{}, 1), + done: make(chan struct{}), + } +} + +func (s *dataTrackSender) setLogger(log protoLogger.Logger) { + s.log = log +} + +func (s *dataTrackSender) send(frame dataTrackFramePackets) { + s.once.Do(func() { go s.run() }) + if dropped := s.push(frame); dropped != nil { + s.log.Debugw("dropping data track frame", "packets", len(dropped)) + } +} + +func (s *dataTrackSender) wake() { + select { + case s.notify <- struct{}{}: + default: + } +} + +func (s *dataTrackSender) stop() { + close(s.done) +} + +func (s *dataTrackSender) push(frame dataTrackFramePackets) (dropped dataTrackFramePackets) { + if len(frame) == 0 { + return nil + } + + s.lock.Lock() + dropped, s.frame = s.frame, frame + s.lock.Unlock() + + s.wake() + return dropped +} + +func (s *dataTrackSender) pop() dataTrackFramePackets { + s.lock.Lock() + defer s.lock.Unlock() + + frame := s.frame + s.frame = nil + return frame +} + +func (s *dataTrackSender) run() { + var ( + dc *webrtc.DataChannel + inFlight dataTrackFramePackets + ) + for { + select { + case <-s.done: + return + case <-s.notify: + } + + for { + if current := s.dc(); current != dc { + dc, inFlight = current, nil + } + if dc == nil || dc.ReadyState() != webrtc.DataChannelStateOpen || dc.BufferedAmount() > dataTrackBufferedAmountLowThreshold { + break + } + if len(inFlight) == 0 { + if inFlight = s.pop(); inFlight == nil { + break + } + } + if err := dc.Send(inFlight[0]); err != nil { + s.log.Debugw("could not send data track packet", "error", err) + } + inFlight = inFlight[1:] + } + } +} diff --git a/datatracksender_test.go b/datatracksender_test.go new file mode 100644 index 00000000..e7ddbbfa --- /dev/null +++ b/datatracksender_test.go @@ -0,0 +1,56 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package lksdk + +import ( + "testing" + + "github.com/pion/webrtc/v4" + "github.com/stretchr/testify/require" +) + +func testFrame(marker byte, packets int) dataTrackFramePackets { + frame := make(dataTrackFramePackets, packets) + for i := range frame { + frame[i] = []byte{marker, byte(i)} + } + return frame +} + +func TestDataTrackSenderQueue(t *testing.T) { + s := newDataTrackSender(func() *webrtc.DataChannel { return nil }, logger) + t.Cleanup(s.stop) + + require.Nil(t, s.push(nil)) + require.Nil(t, s.pop()) + + multi := testFrame(0xaa, 13) + require.Nil(t, s.push(multi)) + require.Equal(t, multi, s.pop()) + + older, newer := testFrame(0x01, 4), testFrame(0x02, 3) + require.Nil(t, s.push(older)) + require.Equal(t, older, s.push(newer)) + require.Equal(t, newer, s.pop()) + require.Nil(t, s.pop()) +} + +func TestSendDataTrackFrameKeepsFreshest(t *testing.T) { + engine := newTestEngine(t) + + engine.sendDataTrackFrame(testFrame(0x01, 2)) + engine.sendDataTrackFrame(testFrame(0x02, 2)) + require.Equal(t, testFrame(0x02, 2), engine.dataTrackSender.pop()) +} diff --git a/engine.go b/engine.go index 74e68a05..180a3d5f 100644 --- a/engine.go +++ b/engine.go @@ -134,6 +134,8 @@ type RTCEngine struct { reliableMsgLock sync.Mutex reliableMsgSeq uint32 + dataTrackSender *dataTrackSender + trackPublishedListenersLock sync.Mutex trackPublishedListeners map[string]chan *livekit.TrackPublishedResponse @@ -174,6 +176,7 @@ func NewRTCEngine( Logger: e.log, Processor: e, }) + e.dataTrackSender = newDataTrackSender(e.dataTrackDataChannel, e.log) e.configureSignalling(useSinglePeerConnection) return e @@ -206,6 +209,7 @@ func (e *RTCEngine) configureSignalling(useSinglePeerConnection bool) { // SetLogger overrides default logger. func (e *RTCEngine) SetLogger(l protoLogger.Logger) { e.log = l + e.dataTrackSender.setLogger(l) e.connectionManager.setLogger(l) e.signalling.SetLogger(l) e.signalHandler.SetLogger(l) @@ -373,6 +377,7 @@ func (e *RTCEngine) Close() { e.connectionManager.setClosed() e.abortPendingRequests() + e.dataTrackSender.stop() e.pclock.Lock() e.pendingPublisherOffer = webrtc.SessionDescription{} @@ -571,6 +576,9 @@ func (e *RTCEngine) createPublisherPCLocked(configuration webrtc.Configuration) return err } e.dataTrackDC.OnMessage(e.handleDataTrackPacket) + e.dataTrackDC.SetBufferedAmountLowThreshold(dataTrackBufferedAmountLowThreshold) + e.dataTrackDC.OnBufferedAmountLow(e.dataTrackSender.wake) + e.dataTrackDC.OnOpen(e.dataTrackSender.wake) e.dclock.Unlock() return nil @@ -1836,3 +1844,13 @@ func waitUntilConnected(d time.Duration, test func() bool) error { } } } + +func (e *RTCEngine) dataTrackDataChannel() *webrtc.DataChannel { + e.dclock.RLock() + defer e.dclock.RUnlock() + return e.dataTrackDC +} + +func (e *RTCEngine) sendDataTrackFrame(frame dataTrackFramePackets) { + e.dataTrackSender.send(frame) +} From fc4cba3b4b463b01da4734148a84bf85bccb0635 Mon Sep 17 00:00:00 2001 From: Jacob Gelman <3182119+ladvoc@users.noreply.github.com> Date: Fri, 4 Sep 2026 15:59:18 -0700 Subject: [PATCH 4/7] "packets" -> "numPackets" --- datatracksender.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/datatracksender.go b/datatracksender.go index ff25115d..eeccf434 100644 --- a/datatracksender.go +++ b/datatracksender.go @@ -55,7 +55,7 @@ func (s *dataTrackSender) setLogger(log protoLogger.Logger) { func (s *dataTrackSender) send(frame dataTrackFramePackets) { s.once.Do(func() { go s.run() }) if dropped := s.push(frame); dropped != nil { - s.log.Debugw("dropping data track frame", "packets", len(dropped)) + s.log.Debugw("dropping data track frame", "numPackets", len(dropped)) } } From 3ec8fe78d6434b2f69d3d78ddb2e923c2d768f69 Mon Sep 17 00:00:00 2001 From: Jacob Gelman <3182119+ladvoc@users.noreply.github.com> Date: Fri, 4 Sep 2026 16:12:47 -0700 Subject: [PATCH 5/7] Run sender in constructor --- datatracksender.go | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/datatracksender.go b/datatracksender.go index eeccf434..1c73139e 100644 --- a/datatracksender.go +++ b/datatracksender.go @@ -36,16 +36,17 @@ type dataTrackSender struct { frame dataTrackFramePackets notify chan struct{} done chan struct{} - once sync.Once } func newDataTrackSender(dc func() *webrtc.DataChannel, log protoLogger.Logger) *dataTrackSender { - return &dataTrackSender{ + s := &dataTrackSender{ dc: dc, log: log, notify: make(chan struct{}, 1), done: make(chan struct{}), } + go s.run() + return s } func (s *dataTrackSender) setLogger(log protoLogger.Logger) { @@ -53,7 +54,6 @@ func (s *dataTrackSender) setLogger(log protoLogger.Logger) { } func (s *dataTrackSender) send(frame dataTrackFramePackets) { - s.once.Do(func() { go s.run() }) if dropped := s.push(frame); dropped != nil { s.log.Debugw("dropping data track frame", "numPackets", len(dropped)) } From d0ab90e5c4b8bd09176d0b0646e2b268d5162f3d Mon Sep 17 00:00:00 2001 From: Jacob Gelman <3182119+ladvoc@users.noreply.github.com> Date: Tue, 8 Sep 2026 11:52:50 -0700 Subject: [PATCH 6/7] Simplify data track sender --- datatracksender.go | 25 +++++++++++++++++++------ datatracksender_test.go | 3 +-- engine.go | 9 ++------- 3 files changed, 22 insertions(+), 15 deletions(-) diff --git a/datatracksender.go b/datatracksender.go index 1c73139e..849c4817 100644 --- a/datatracksender.go +++ b/datatracksender.go @@ -29,18 +29,17 @@ type dataTrackFramePackets [][]byte // dataTrackSender paces frames onto the data track channel, keeping only the freshest frame // while the channel's send buffer is above the low threshold. type dataTrackSender struct { - dc func() *webrtc.DataChannel log protoLogger.Logger lock sync.Mutex + dc *webrtc.DataChannel frame dataTrackFramePackets notify chan struct{} done chan struct{} } -func newDataTrackSender(dc func() *webrtc.DataChannel, log protoLogger.Logger) *dataTrackSender { +func newDataTrackSender(log protoLogger.Logger) *dataTrackSender { s := &dataTrackSender{ - dc: dc, log: log, notify: make(chan struct{}, 1), done: make(chan struct{}), @@ -49,6 +48,15 @@ func newDataTrackSender(dc func() *webrtc.DataChannel, log protoLogger.Logger) * return s } +// setDataChannel points the sender at the channel it should write to. +func (s *dataTrackSender) setDataChannel(dc *webrtc.DataChannel) { + s.lock.Lock() + s.dc = dc + s.lock.Unlock() + + s.wake() +} + func (s *dataTrackSender) setLogger(log protoLogger.Logger) { s.log = log } @@ -104,10 +112,15 @@ func (s *dataTrackSender) run() { case <-s.notify: } + s.lock.Lock() + current := s.dc + s.lock.Unlock() + if current != dc { + // A partially sent frame cannot be completed on a new channel. + dc, inFlight = current, nil + } + for { - if current := s.dc(); current != dc { - dc, inFlight = current, nil - } if dc == nil || dc.ReadyState() != webrtc.DataChannelStateOpen || dc.BufferedAmount() > dataTrackBufferedAmountLowThreshold { break } diff --git a/datatracksender_test.go b/datatracksender_test.go index e7ddbbfa..a76ae0e3 100644 --- a/datatracksender_test.go +++ b/datatracksender_test.go @@ -17,7 +17,6 @@ package lksdk import ( "testing" - "github.com/pion/webrtc/v4" "github.com/stretchr/testify/require" ) @@ -30,7 +29,7 @@ func testFrame(marker byte, packets int) dataTrackFramePackets { } func TestDataTrackSenderQueue(t *testing.T) { - s := newDataTrackSender(func() *webrtc.DataChannel { return nil }, logger) + s := newDataTrackSender(logger) t.Cleanup(s.stop) require.Nil(t, s.push(nil)) diff --git a/engine.go b/engine.go index 180a3d5f..62bcf93d 100644 --- a/engine.go +++ b/engine.go @@ -176,7 +176,7 @@ func NewRTCEngine( Logger: e.log, Processor: e, }) - e.dataTrackSender = newDataTrackSender(e.dataTrackDataChannel, e.log) + e.dataTrackSender = newDataTrackSender(e.log) e.configureSignalling(useSinglePeerConnection) return e @@ -579,6 +579,7 @@ func (e *RTCEngine) createPublisherPCLocked(configuration webrtc.Configuration) e.dataTrackDC.SetBufferedAmountLowThreshold(dataTrackBufferedAmountLowThreshold) e.dataTrackDC.OnBufferedAmountLow(e.dataTrackSender.wake) e.dataTrackDC.OnOpen(e.dataTrackSender.wake) + e.dataTrackSender.setDataChannel(e.dataTrackDC) e.dclock.Unlock() return nil @@ -1845,12 +1846,6 @@ func waitUntilConnected(d time.Duration, test func() bool) error { } } -func (e *RTCEngine) dataTrackDataChannel() *webrtc.DataChannel { - e.dclock.RLock() - defer e.dclock.RUnlock() - return e.dataTrackDC -} - func (e *RTCEngine) sendDataTrackFrame(frame dataTrackFramePackets) { e.dataTrackSender.send(frame) } From b0427c63e2e2d5b710732350433492ac8ecbfbcf Mon Sep 17 00:00:00 2001 From: Jacob Gelman <3182119+ladvoc@users.noreply.github.com> Date: Tue, 8 Sep 2026 12:11:03 -0700 Subject: [PATCH 7/7] Fix lint --- datatracksender.go | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/datatracksender.go b/datatracksender.go index 849c4817..495a0153 100644 --- a/datatracksender.go +++ b/datatracksender.go @@ -120,10 +120,7 @@ func (s *dataTrackSender) run() { dc, inFlight = current, nil } - for { - if dc == nil || dc.ReadyState() != webrtc.DataChannelStateOpen || dc.BufferedAmount() > dataTrackBufferedAmountLowThreshold { - break - } + for dc != nil && dc.ReadyState() == webrtc.DataChannelStateOpen && dc.BufferedAmount() <= dataTrackBufferedAmountLowThreshold { if len(inFlight) == 0 { if inFlight = s.pop(); inFlight == nil { break