mirror of
https://github.com/livekit/livekit.git
synced 2026-07-21 06:51:07 +00:00
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.
This commit is contained in:
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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,
|
||||
)
|
||||
|
||||
@@ -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,
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user