From 2a3de843513413afd3e1945960be7bc38344c285 Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Mon, 29 Jan 2024 13:03:32 +0530 Subject: [PATCH] Reverting participant worker. (#2428) * Reverting participant worker. Reverts https://github.com/livekit/livekit/pull/2420 partially. This did not revert clean. So, reverting manually. Also, keeping the drive-by clean up bits. * fix test --- pkg/rtc/participant.go | 46 ++++++--- pkg/rtc/room.go | 214 ++++++++++++----------------------------- pkg/rtc/room_test.go | 6 +- 3 files changed, 100 insertions(+), 166 deletions(-) diff --git a/pkg/rtc/participant.go b/pkg/rtc/participant.go index df60f272a..ee4f20da8 100644 --- a/pkg/rtc/participant.go +++ b/pkg/rtc/participant.go @@ -690,9 +690,13 @@ func (p *ParticipantImpl) handleMigrateTracks() { } p.pendingTracksLock.Unlock() - for _, t := range addedTracks { - p.handleTrackPublished(t) - } + // launch callbacks in goroutine since they could block. + // callbacks handle webhooks as well as db persistence + go func() { + for _, t := range addedTracks { + p.handleTrackPublished(t) + } + }() } func (p *ParticipantImpl) removePendingMigratedTrack(mt *MediaTrack) { @@ -929,7 +933,7 @@ func (p *ParticipantImpl) SetMigrateState(s types.MigrateState) { } if onMigrateStateChange := p.getOnMigrateStateChange(); onMigrateStateChange != nil { - onMigrateStateChange(p, s) + go onMigrateStateChange(p, s) } } @@ -1337,7 +1341,7 @@ func (p *ParticipantImpl) updateState(state livekit.ParticipantInfo_State) { p.dirty.Store(true) if onStateChange := p.getOnStateChange(); onStateChange != nil { - onStateChange(p, state) + go onStateChange(p, state) } } @@ -1939,12 +1943,14 @@ func (p *ParticipantImpl) mediaTrackReceived(track *webrtc.TrackRemote, rtpRecei } if newTrack { - p.pubLogger.Debugw( - "track published", - "trackID", mt.ID(), - "track", logger.Proto(mt.ToProto()), - ) - p.handleTrackPublished(mt) + go func() { + p.pubLogger.Debugw( + "track published", + "trackID", mt.ID(), + "track", logger.Proto(mt.ToProto()), + ) + p.handleTrackPublished(mt) + }() } return mt, newTrack @@ -2052,6 +2058,15 @@ func (p *ParticipantImpl) addMediaTrack(signalCid string, sdpCid string, ti *liv p.supervisor.ClearPublishedTrack(trackID, mt) } + // not logged when closing + p.params.Telemetry.TrackUnpublished( + context.Background(), + p.ID(), + p.Identity(), + mt.ToProto(), + !p.IsClosed(), + ) + // re-use Track sid p.pendingTracksLock.Lock() if pti := p.pendingTracks[signalCid]; pti != nil { @@ -2080,6 +2095,15 @@ func (p *ParticipantImpl) handleTrackPublished(track types.MediaTrack) { onTrackPublished(p, track) } + // send webhook after callbacks are complete, persistence and state handling happens + // in `onTrackPublished` cb + p.params.Telemetry.TrackPublished( + context.Background(), + p.ID(), + p.Identity(), + track.ToProto(), + ) + p.pendingTracksLock.Lock() delete(p.pendingPublishingTracks, track.ID()) p.pendingTracksLock.Unlock() diff --git a/pkg/rtc/room.go b/pkg/rtc/room.go index c6c7609f5..1daf41cb6 100644 --- a/pkg/rtc/room.go +++ b/pkg/rtc/room.go @@ -73,11 +73,6 @@ type disconnectSignalOnResumeNoMessages struct { closedCount int } -type participantWorker struct { - eventsQueue *sutils.OpsQueue - participants []types.LocalParticipant -} - type Room struct { lock sync.RWMutex @@ -99,7 +94,6 @@ type Room struct { // map of identity -> Participant participants map[livekit.ParticipantIdentity]types.LocalParticipant - participantWorkers map[livekit.ParticipantIdentity]*participantWorker participantOpts map[livekit.ParticipantIdentity]*ParticipantOptions participantRequestSources map[livekit.ParticipantIdentity]routing.MessageSource hasPublished map[livekit.ParticipantIdentity]bool @@ -157,7 +151,6 @@ func NewRoom( trackManager: NewRoomTrackManager(), serverInfo: serverInfo, participants: make(map[livekit.ParticipantIdentity]types.LocalParticipant), - participantWorkers: make(map[livekit.ParticipantIdentity]*participantWorker), participantOpts: make(map[livekit.ParticipantIdentity]*ParticipantOptions), participantRequestSources: make(map[livekit.ParticipantIdentity]routing.MessageSource), hasPublished: make(map[livekit.ParticipantIdentity]bool), @@ -344,100 +337,78 @@ func (r *Room) Join(participant types.LocalParticipant, requestSource routing.Me r.joinedAt.Store(time.Now().Unix()) } - pw := r.addParticipantWorkerLocked(participant) - participant.OnStateChange(func(p types.LocalParticipant, state livekit.ParticipantInfo_State) { - pw.eventsQueue.Enqueue(func() { - if r.onParticipantChanged != nil { - r.onParticipantChanged(p) + if r.onParticipantChanged != nil { + r.onParticipantChanged(p) + } + r.broadcastParticipantState(p, broadcastOptions{skipSource: true}) + + if state == livekit.ParticipantInfo_ACTIVE { + // subscribe participant to existing published tracks + r.subscribeToExistingTracks(p) + + meta := &livekit.AnalyticsClientMeta{ + ClientConnectTime: uint32(time.Since(p.ConnectedAt()).Milliseconds()), } - r.broadcastParticipantState(p, broadcastOptions{skipSource: true}) - - if state == livekit.ParticipantInfo_ACTIVE { - // subscribe participant to existing published tracks - r.subscribeToExistingTracks(p) - - meta := &livekit.AnalyticsClientMeta{ - ClientConnectTime: uint32(time.Since(p.ConnectedAt()).Milliseconds()), + cds := p.GetICEConnectionDetails() + for _, cd := range cds { + if cd.Type != types.ICEConnectionTypeUnknown { + meta.ConnectionType = string(cd.Type) + break } - cds := p.GetICEConnectionDetails() - for _, cd := range cds { - if cd.Type != types.ICEConnectionTypeUnknown { - meta.ConnectionType = string(cd.Type) - break - } - } - r.telemetry.ParticipantActive(context.Background(), - r.ToProto(), - p.ToProto(), - meta, - false, - ) - - p.GetLogger().Infow("participant active", connectionDetailsFields(cds)...) - } else if state == livekit.ParticipantInfo_DISCONNECTED { - // remove participant from room - r.RemoveParticipant(p.Identity(), p.ID(), types.ParticipantCloseReasonStateDisconnected) } - }) + r.telemetry.ParticipantActive(context.Background(), + r.ToProto(), + p.ToProto(), + meta, + false, + ) + + p.GetLogger().Infow("participant active", connectionDetailsFields(cds)...) + } else if state == livekit.ParticipantInfo_DISCONNECTED { + // remove participant from room + go r.RemoveParticipant(p.Identity(), p.ID(), types.ParticipantCloseReasonStateDisconnected) + } }) // it's important to set this before connection, we don't want to miss out on any published tracks - participant.OnTrackPublished(func(p types.LocalParticipant, t types.MediaTrack) { - pw.eventsQueue.Enqueue(func() { - r.onTrackPublished(p, t) - }) - }) - participant.OnTrackUpdated(func(p types.LocalParticipant, t types.MediaTrack) { - pw.eventsQueue.Enqueue(func() { - r.onTrackUpdated(p, t) - }) - }) - participant.OnTrackUnpublished(func(p types.LocalParticipant, t types.MediaTrack) { - pw.eventsQueue.Enqueue(func() { - r.onTrackUnpublished(p, t) - }) - }) - participant.OnParticipantUpdate(func(p types.LocalParticipant) { - pw.eventsQueue.Enqueue(func() { - r.onParticipantUpdate(p) - }) - }) + participant.OnTrackPublished(r.onTrackPublished) + participant.OnTrackUpdated(r.onTrackUpdated) + participant.OnTrackUnpublished(r.onTrackUnpublished) + participant.OnParticipantUpdate(r.onParticipantUpdate) participant.OnDataPacket(r.onDataPacket) participant.OnSubscribeStatusChanged(func(publisherID livekit.ParticipantID, subscribed bool) { - pw.eventsQueue.Enqueue(func() { - if subscribed { - pub := r.GetParticipantByID(publisherID) - if pub != nil && pub.State() == livekit.ParticipantInfo_ACTIVE { - // when a participant subscribes to another participant, - // send speaker update if the subscribed to participant is active. - level, active := pub.GetAudioLevel() - if active { - _ = participant.SendSpeakerUpdate([]*livekit.SpeakerInfo{ - { - Sid: string(pub.ID()), - Level: float32(level), - Active: active, - }, - }, false) - } - - if cq := pub.GetConnectionQuality(); cq != nil { - update := &livekit.ConnectionQualityUpdate{} - update.Updates = append(update.Updates, cq) - _ = participant.SendConnectionQualityUpdate(update) - } + if subscribed { + pub := r.GetParticipantByID(publisherID) + if pub != nil && pub.State() == livekit.ParticipantInfo_ACTIVE { + // when a participant subscribes to another participant, + // send speaker update if the subscribed to participant is active. + level, active := pub.GetAudioLevel() + if active { + _ = participant.SendSpeakerUpdate([]*livekit.SpeakerInfo{ + { + Sid: string(pub.ID()), + Level: float32(level), + Active: active, + }, + }, false) + } + + if cq := pub.GetConnectionQuality(); cq != nil { + update := &livekit.ConnectionQualityUpdate{} + update.Updates = append(update.Updates, cq) + _ = participant.SendConnectionQualityUpdate(update) } - } else { - // no longer subscribed to the publisher, clear speaker status - _ = participant.SendSpeakerUpdate([]*livekit.SpeakerInfo{ - { - Sid: string(publisherID), - Level: 0, - Active: false, - }, - }, true) } - }) + } else { + // no longer subscribed to the publisher, clear speaker status + _ = participant.SendSpeakerUpdate([]*livekit.SpeakerInfo{ + { + Sid: string(publisherID), + Level: 0, + Active: false, + }, + }, true) + } }) r.Logger.Debugw("new participant joined", @@ -582,7 +553,6 @@ func (r *Room) RemoveParticipant(identity livekit.ParticipantIdentity, pID livek } delete(r.participants, identity) - r.removeParticipantWorkerLocked(p) delete(r.participantOpts, identity) delete(r.participantRequestSources, identity) delete(r.hasPublished, identity) @@ -1066,14 +1036,6 @@ func (r *Room) onTrackPublished(participant types.LocalParticipant, track types. } }() } - - // send webhook after callbacks are complete, i.e. after persistence and state handling - r.telemetry.TrackPublished( - context.Background(), - participant.ID(), - participant.Identity(), - track.ToProto(), - ) } func (r *Room) onTrackUpdated(p types.LocalParticipant, _ types.MediaTrack) { @@ -1085,14 +1047,6 @@ func (r *Room) onTrackUpdated(p types.LocalParticipant, _ types.MediaTrack) { } func (r *Room) onTrackUnpublished(p types.LocalParticipant, track types.MediaTrack) { - r.telemetry.TrackUnpublished( - context.Background(), - p.ID(), - p.Identity(), - track.ToProto(), - !p.IsClosed(), - ) - r.trackManager.RemoveTrack(track) if !p.IsClosed() { r.broadcastParticipantState(p, broadcastOptions{skipSource: true}) @@ -1496,50 +1450,6 @@ func (r *Room) DebugInfo() map[string]interface{} { return info } -func (r *Room) addParticipantWorkerLocked(p types.LocalParticipant) *participantWorker { - identity := p.Identity() - pw := r.participantWorkers[identity] - if pw != nil { - found := false - for _, participant := range pw.participants { - if p == participant { - found = true - break - } - } - if !found { - pw.participants = append(pw.participants, p) - } - return pw - } - - pw = &participantWorker{ - eventsQueue: sutils.NewOpsQueue(fmt.Sprintf("participant-worker-%s-%s", r.Name(), identity), 0, true), - participants: []types.LocalParticipant{p}, - } - pw.eventsQueue.Start() - r.participantWorkers[identity] = pw - return pw -} - -func (r *Room) removeParticipantWorkerLocked(p types.LocalParticipant) { - identity := p.Identity() - if pw, ok := r.participantWorkers[identity]; ok { - n := len(pw.participants) - for idx, participant := range pw.participants { - if p == participant { - pw.participants[idx] = pw.participants[n-1] - pw.participants = pw.participants[:n-1] - break - } - } - if len(pw.participants) == 0 { - pw.eventsQueue.Stop() - delete(r.participantWorkers, identity) - } - } -} - // ------------------------------------------------------------ func BroadcastDataPacketForRoom(r types.Room, source types.LocalParticipant, dp *livekit.DataPacket, logger logger.Logger) { diff --git a/pkg/rtc/room_test.go b/pkg/rtc/room_test.go index 4990b64d4..4bf6ff6b6 100644 --- a/pkg/rtc/room_test.go +++ b/pkg/rtc/room_test.go @@ -121,7 +121,7 @@ func TestRoomJoin(t *testing.T) { numTracks += len(op.GetPublishedTracks()) } - require.Eventually(t, func() bool { return p.SubscribeToTrackCallCount() == numTracks }, 5*time.Second, 10*time.Millisecond) + require.Equal(t, numTracks, p.SubscribeToTrackCallCount()) }) t.Run("participant state change is broadcasted to others", func(t *testing.T) { @@ -217,7 +217,7 @@ func TestParticipantUpdate(t *testing.T) { expected += 1 } fp := p.(*typesfakes.FakeLocalParticipant) - require.Eventually(t, func() bool { return fp.SendParticipantUpdateCallCount() == expected }, 5*time.Second, 10*time.Millisecond) + require.Equal(t, expected, fp.SendParticipantUpdateCallCount()) } }) } @@ -423,8 +423,8 @@ func TestNewTrack(t *testing.T) { require.NotNil(t, trackCB) trackCB(pub, track) // only p1 should've been subscribed to - require.Eventually(t, func() bool { return p1.SubscribeToTrackCallCount() == 1 }, 5*time.Second, 10*time.Millisecond) require.Equal(t, 0, p0.SubscribeToTrackCallCount()) + require.Equal(t, 1, p1.SubscribeToTrackCallCount()) }) }