From 5ac5bd236a903a1a5d720a81abec3669db5edf72 Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Sat, 17 Feb 2024 13:40:07 +0530 Subject: [PATCH] Let track events go through after participant close. (#2487) * Let track events go through after participant close. Also, reducing lock scope in telemetry service. * use shadow --- pkg/rtc/participant.go | 2 +- pkg/rtc/subscriptionmanager.go | 2 +- pkg/telemetry/telemetryservice.go | 39 +++++++++++++++++++++++-------- 3 files changed, 31 insertions(+), 12 deletions(-) diff --git a/pkg/rtc/participant.go b/pkg/rtc/participant.go index 5cf2dc299..962b90b1f 100644 --- a/pkg/rtc/participant.go +++ b/pkg/rtc/participant.go @@ -2063,7 +2063,7 @@ func (p *ParticipantImpl) addMediaTrack(signalCid string, sdpCid string, ti *liv p.ID(), p.Identity(), mt.ToProto(), - !p.IsClosed(), + true, ) // re-use Track sid diff --git a/pkg/rtc/subscriptionmanager.go b/pkg/rtc/subscriptionmanager.go index 353936fd1..eb879d1a6 100644 --- a/pkg/rtc/subscriptionmanager.go +++ b/pkg/rtc/subscriptionmanager.go @@ -662,7 +662,7 @@ func (m *SubscriptionManager) handleSubscribedTrackClose(s *trackSubscription, w context.Background(), m.params.Participant.ID(), &livekit.TrackInfo{Sid: string(s.trackID), Type: subTrack.MediaTrack().Kind()}, - !willBeResumed && !m.params.Participant.IsClosed(), + !willBeResumed, ) dt := subTrack.DownTrack() diff --git a/pkg/telemetry/telemetryservice.go b/pkg/telemetry/telemetryservice.go index fe5e329c2..c684d5842 100644 --- a/pkg/telemetry/telemetryservice.go +++ b/pkg/telemetry/telemetryservice.go @@ -24,6 +24,7 @@ import ( "github.com/livekit/protocol/livekit" "github.com/livekit/protocol/logger" "github.com/livekit/protocol/webhook" + "golang.org/x/exp/maps" ) //go:generate go run github.com/maxbrunsfeld/counterfeiter/v6 . TelemetryService @@ -93,8 +94,9 @@ type telemetryService struct { notifier webhook.QueuedNotifier jobsQueue *utils.OpsQueue - lock sync.RWMutex - workers map[livekit.ParticipantID]*StatsWorker + lock sync.RWMutex + workers map[livekit.ParticipantID]*StatsWorker + workersShadow []*StatsWorker } func NewTelemetryService(notifier webhook.QueuedNotifier, analytics AnalyticsService) TelemetryService { @@ -113,10 +115,11 @@ func NewTelemetryService(notifier webhook.QueuedNotifier, analytics AnalyticsSer } func (t *telemetryService) FlushStats() { - t.lock.Lock() - defer t.lock.Unlock() + t.lock.RLock() + workersShadow := t.workersShadow + t.lock.RUnlock() - for _, worker := range t.workers { + for _, worker := range workersShadow { worker.Flush() } } @@ -167,21 +170,37 @@ func (t *telemetryService) createWorker(ctx context.Context, t.lock.Lock() t.workers[participantID] = worker + t.workersShadow = maps.Values(t.workers) t.lock.Unlock() return worker } func (t *telemetryService) cleanupWorkers() { - t.lock.Lock() - defer t.lock.Unlock() + t.lock.RLock() + workersShadow := t.workersShadow + t.lock.RUnlock() - for participantID, worker := range t.workers { + toReap := make([]livekit.ParticipantID, 0, len(workersShadow)) + for _, worker := range workersShadow { closedAt := worker.ClosedAt() if !closedAt.IsZero() && time.Since(closedAt) > workerCleanupWait { - logger.Debugw("reaping analytics worker for participant", "pID", participantID) - delete(t.workers, participantID) + worker.Flush() + + toReap = append(toReap, worker.ParticipantID()) } } + + if len(toReap) == 0 { + return + } + + t.lock.Lock() + logger.Debugw("reaping analytics worker for participants", "pID", toReap) + for _, pID := range toReap { + delete(t.workers, pID) + } + t.workersShadow = maps.Values(t.workers) + t.lock.Unlock() } func (t *telemetryService) LocalRoomState(ctx context.Context, info *livekit.AnalyticsNodeRooms) {