From cde8962709080dca1a343add594d1424c6b86f4c Mon Sep 17 00:00:00 2001 From: Paul Wells Date: Sat, 23 May 2026 17:42:55 -0700 Subject: [PATCH] rtc: emit per-data-track bytes via BytesTrackStats (#4540) Data tracks (the new _data_track datachannel) previously only updated a private dataTrackStats that logged a single summary at Close. Bytes never reached the OnTrackStats -> TelemetryService.TrackStats pipeline that media tracks and signal channels feed. Wire DataTrack (UPSTREAM, publisher-home) and DataDownTrack (DOWNSTREAM, per-subscriber) into BytesTrackStats on the same 5s cadence, mirroring the media-track convention: subscriber's country and ID with publisher's track ID for DOWNSTREAM. Cross-region proxy DataTracks leave the stats pointer nil (no publisher reporter on that node, and relayed bytes would double-count). Legacy dataTrackStats packet-loss/frame counters are preserved. --- pkg/rtc/datadowntrack.go | 8 ++++++++ pkg/rtc/datatrack.go | 18 ++++++++++++++++++ pkg/rtc/participant.go | 9 +++++++++ pkg/rtc/participant_data_track.go | 9 +++++++++ 4 files changed, 44 insertions(+) diff --git a/pkg/rtc/datadowntrack.go b/pkg/rtc/datadowntrack.go index 1cc94e35c..89c069f1b 100644 --- a/pkg/rtc/datadowntrack.go +++ b/pkg/rtc/datadowntrack.go @@ -33,6 +33,7 @@ type DataDownTrackParams struct { PublishDataTrack types.DataTrack Handle uint16 Transport types.DataTrackTransport + BytesTrackStats *BytesTrackStats } type DataDownTrack struct { @@ -59,6 +60,9 @@ func NewDataDownTrack(params DataDownTrackParams) (*DataDownTrack, error) { func (d *DataDownTrack) Close() { d.logger.Infow("closing data down track") + if d.params.BytesTrackStats != nil { + d.params.BytesTrackStats.Stop() + } d.params.PublishDataTrack.DeleteDataDownTrack(d.SubscriberID()) } @@ -93,6 +97,10 @@ func (d *DataDownTrack) WritePacket(data []byte, packet *datatrack.Packet, _arri } if err := d.params.Transport.SendDataTrackMessage(buf); err != nil { d.logger.Warnw("could not send data track message", err) + return + } + if d.params.BytesTrackStats != nil { + d.params.BytesTrackStats.AddBytes(uint64(len(buf)), true) } } diff --git a/pkg/rtc/datatrack.go b/pkg/rtc/datatrack.go index 878fbbba9..1e8c4d55c 100644 --- a/pkg/rtc/datatrack.go +++ b/pkg/rtc/datatrack.go @@ -37,6 +37,7 @@ type DataTrackParams struct { Logger logger.Logger ParticipantID func() livekit.ParticipantID ParticipantIdentity livekit.ParticipantIdentity + BytesTrackStats *BytesTrackStats } type DataTrack struct { @@ -76,6 +77,9 @@ func (d *DataTrack) Close() { d.closed.Break() d.stats.Close() + if d.params.BytesTrackStats != nil { + d.params.BytesTrackStats.Stop() + } } func (d *DataTrack) PublisherID() livekit.ParticipantID { @@ -110,14 +114,25 @@ func (d *DataTrack) AddSubscriber(sub types.LocalParticipant) (types.DataDownTra return nil, errAlreadySubscribed } + bytesStats := NewBytesTrackStats( + sub.GetCountry(), + d.ID(), + sub.ID(), + sub.Kind(), + sub.KindDetails(), + sub.GetTelemetryListener(), + sub.GetReporter(), + ) dataDownTrack, err := NewDataDownTrack(DataDownTrackParams{ Logger: sub.GetLogger().WithValues("trackID", d.ID()), SubscriberID: sub.ID(), PublishDataTrack: d, Handle: sub.GetNextSubscribedDataTrackHandle(), Transport: sub.GetDataTrackTransport(), + BytesTrackStats: bytesStats, }) if err != nil { + bytesStats.Stop() return nil, err } @@ -165,6 +180,9 @@ func (d *DataTrack) DeleteDataDownTrack(subscriberID livekit.ParticipantID) { func (d *DataTrack) HandlePacket(data []byte, packet *datatrack.Packet, arrivalTime int64) { d.stats.Update(packet, arrivalTime, len(data)) + if d.params.BytesTrackStats != nil { + d.params.BytesTrackStats.AddBytes(uint64(len(data)), false) + } d.downTrackSpreader.Broadcast(func(dts types.DataTrackSender) { dts.WritePacket(data, packet, arrivalTime) diff --git a/pkg/rtc/participant.go b/pkg/rtc/participant.go index 4c3461b36..3c409f92a 100644 --- a/pkg/rtc/participant.go +++ b/pkg/rtc/participant.go @@ -1365,6 +1365,15 @@ func (p *ParticipantImpl) SetMigrateInfo( Logger: p.params.Logger.WithValues("trackID", dti.Sid), ParticipantID: p.ID, ParticipantIdentity: p.params.Identity, + BytesTrackStats: NewBytesTrackStats( + p.params.Country, + livekit.TrackID(dti.Sid), + p.ID(), + p.Kind(), + p.KindDetails(), + p.params.TelemetryListener, + p.params.Reporter, + ), }, dti, ) diff --git a/pkg/rtc/participant_data_track.go b/pkg/rtc/participant_data_track.go index f5da89f8a..829beaf63 100644 --- a/pkg/rtc/participant_data_track.go +++ b/pkg/rtc/participant_data_track.go @@ -99,6 +99,15 @@ func (p *ParticipantImpl) HandlePublishDataTrackRequest(req *livekit.PublishData Logger: p.params.Logger.WithValues("trackID", dti.Sid), ParticipantID: p.ID, ParticipantIdentity: p.params.Identity, + BytesTrackStats: NewBytesTrackStats( + p.params.Country, + livekit.TrackID(dti.Sid), + p.ID(), + p.Kind(), + p.KindDetails(), + p.params.TelemetryListener, + p.params.Reporter, + ), }, dti, )