Improved subscription manager logging to help with debugging (#1337)

This commit is contained in:
David Zhao
2023-01-26 23:31:03 -08:00
committed by GitHub
parent 2d6c896bba
commit c146398f32
2 changed files with 37 additions and 39 deletions
+36 -39
View File
@@ -68,7 +68,7 @@ func NewSubscriptionManager(params SubscriptionManagerParams) *SubscriptionManag
params: params,
subscriptions: make(map[livekit.TrackID]*trackSubscription),
subscribedTo: make(map[livekit.ParticipantID]map[livekit.TrackID]struct{}),
reconcileCh: make(chan livekit.TrackID, 10),
reconcileCh: make(chan livekit.TrackID, 50),
closeCh: make(chan struct{}),
doneCh: make(chan struct{}),
}
@@ -116,18 +116,21 @@ func (m *SubscriptionManager) SubscribeToTrack(trackID livekit.TrackID, publishe
m.lock.Lock()
sub, ok := m.subscriptions[trackID]
if !ok {
sub = newTrackSubscription(m.params.Participant.ID(), trackID)
m.subscriptions[trackID] = sub
}
m.lock.Unlock()
sub.setPublisher(publisherIdentity, publisherID)
if sub.setDesired(true) {
m.params.Logger.Infow("subscribing to track",
sLogger := m.params.Logger.WithValues(
"trackID", trackID,
"publisherID", publisherID,
"publisherIdentity", publisherIdentity,
)
sub = newTrackSubscription(m.params.Participant.ID(), trackID, sLogger)
m.subscriptions[trackID] = sub
}
desireChanged := sub.setDesired(true)
m.lock.Unlock()
sub.setPublisher(publisherIdentity, publisherID)
if desireChanged {
sub.logger.Infow("subscribing to track")
}
// always reconcile, since SubscribeToTrack could be called when the track is ready
m.queueReconcile(trackID)
}
@@ -141,11 +144,7 @@ func (m *SubscriptionManager) UnsubscribeFromTrack(trackID livekit.TrackID) {
}
if sub.setDesired(false) {
m.params.Logger.Infow("unsubscribing from track",
"trackID", trackID,
"publisherID", sub.getPublisherID(),
"publisherIdentity", sub.getPublisherIdentity(),
)
sub.logger.Infow("unsubscribing from track")
m.queueReconcile(trackID)
}
}
@@ -198,7 +197,10 @@ func (m *SubscriptionManager) UpdateSubscribedTrackSettings(trackID livekit.Trac
m.lock.Lock()
sub, ok := m.subscriptions[trackID]
if !ok {
sub = newTrackSubscription(m.params.Participant.ID(), trackID)
sLogger := m.params.Logger.WithValues(
"trackID", trackID,
)
sub = newTrackSubscription(m.params.Participant.ID(), trackID, sLogger)
m.subscriptions[trackID] = sub
}
m.lock.Unlock()
@@ -293,19 +295,14 @@ func (m *SubscriptionManager) reconcileSubscription(s *trackSubscription) {
// from it. this is the *only* case we'd change desired state
if s.durationSinceStart() > notFoundTimeout {
s.maybeRecordError(m.params.Telemetry, m.params.Participant.ID(), err, true)
m.params.Logger.Infow("unsubscribing track since track isn't available",
"trackID", s.trackID,
"publisherID", s.getPublisherID(),
"publisherIdentity", s.getPublisherIdentity(),
)
s.logger.Infow("unsubscribing track since track isn't available")
s.setDesired(false)
m.queueReconcile(s.trackID)
}
default:
// all other errors
m.params.Logger.Warnw("failed to subscribe", err,
s.logger.Warnw("failed to subscribe", err,
"attempt", s.numAttempts.Load(),
"trackID", s.trackID,
)
if s.durationSinceStart() > subscriptionTimeout {
s.maybeRecordError(m.params.Telemetry, m.params.Participant.ID(), err, false)
@@ -321,9 +318,7 @@ func (m *SubscriptionManager) reconcileSubscription(s *trackSubscription) {
if s.needsUnsubscribe() {
if err := m.unsubscribe(s); err != nil {
m.params.Logger.Errorw("failed to unsubscribe", err,
"trackID", s.trackID,
)
s.logger.Errorw("failed to unsubscribe", err)
} else {
// successfully unsubscribed, remove from map
m.lock.Lock()
@@ -338,11 +333,7 @@ func (m *SubscriptionManager) reconcileSubscription(s *trackSubscription) {
if s.needsBind() {
// check bound status, notify error callback if it's not bound
if s.durationSinceStart() > subscriptionTimeout {
m.params.Logger.Errorw("track not bound after timeout", nil,
"trackID", s.trackID,
"publisherID", s.getPublisherID(),
"publisherIdentity", s.getPublisherIdentity(),
)
s.logger.Errorw("track not bound after timeout", nil)
s.maybeRecordError(m.params.Telemetry, m.params.Participant.ID(), ErrTrackNotBound, false)
m.params.OnSubcriptionError(s.trackID)
}
@@ -383,6 +374,7 @@ func (m *SubscriptionManager) reconcileWorker() {
}
func (m *SubscriptionManager) subscribe(s *trackSubscription) error {
s.logger.Debugw("executing subscribe")
s.startAttempt()
if !m.params.Participant.CanSubscribe() {
@@ -394,6 +386,8 @@ func (m *SubscriptionManager) subscribe(s *trackSubscription) error {
return err
}
s.logger.Debugw("resolved track", "result", res)
if res.TrackChangeNotifier != nil && s.setChangeNotifier(res.TrackChangeNotifier) {
// set callback only when we haven't done it before
// we set the observer before checking for existence of track, so that we may get notified when track becomes
@@ -459,7 +453,8 @@ func (m *SubscriptionManager) subscribe(s *trackSubscription) error {
}
func (m *SubscriptionManager) unsubscribe(s *trackSubscription) error {
// remove from subscribedTo
s.logger.Debugw("executing unsubscribe")
subTrack := s.getSubscribedTrack()
if subTrack == nil {
// already unsubscribed
@@ -479,16 +474,14 @@ func (m *SubscriptionManager) unsubscribe(s *trackSubscription) error {
// - UpTrack was closed
// - publisher revoked permissions for the participant
func (m *SubscriptionManager) handleSubscribedTrackClose(s *trackSubscription, willBeResumed bool) {
m.params.Logger.Debugw("subscribed track closed",
"trackID", s.trackID,
"publisherID", s.getPublisherID(),
"publisherIdentity", s.getPublisherIdentity(),
s.logger.Debugw("subscribed track closed",
"willBeResumed", willBeResumed,
)
subTrack := s.getSubscribedTrack()
if subTrack == nil {
return
}
s.setSubscribedTrack(nil)
// remove from subscribedTo
publisherID := s.getPublisherID()
@@ -509,7 +502,6 @@ func (m *SubscriptionManager) handleSubscribedTrackClose(s *trackSubscription, w
}
subTrack.OnClose(nil)
s.setSubscribedTrack(nil)
go m.params.OnTrackUnsubscribed(subTrack)
// always trigger to decrement unsubscribed counter. However, only log an analytics event when
@@ -524,9 +516,7 @@ func (m *SubscriptionManager) handleSubscribedTrackClose(s *trackSubscription, w
if !willBeResumed {
sender := subTrack.RTPSender()
if sender != nil {
m.params.Logger.Debugw("removing PeerConnection track",
"publisher", subTrack.PublisherIdentity(),
"publisherID", subTrack.PublisherID(),
s.logger.Debugw("removing PeerConnection track",
"kind", subTrack.MediaTrack().Kind(),
)
@@ -553,6 +543,7 @@ func (m *SubscriptionManager) handleSubscribedTrackClose(s *trackSubscription, w
type trackSubscription struct {
subscriberID livekit.ParticipantID
trackID livekit.TrackID
logger logger.Logger
lock sync.RWMutex
desired bool
@@ -568,10 +559,11 @@ type trackSubscription struct {
subStartedAt atomic.Pointer[time.Time]
}
func newTrackSubscription(subscriberID livekit.ParticipantID, trackID livekit.TrackID) *trackSubscription {
func newTrackSubscription(subscriberID livekit.ParticipantID, trackID livekit.TrackID, l logger.Logger) *trackSubscription {
return &trackSubscription{
subscriberID: subscriberID,
trackID: trackID,
logger: l,
// default allow
hasPermission: true,
}
@@ -609,6 +601,10 @@ func (s *trackSubscription) setDesired(desired bool) bool {
s.changeNotifier.RemoveObserver(string(s.subscriberID))
s.changeNotifier = nil
}
if desired {
// reset attempt
s.numAttempts.Store(0)
}
return true
}
@@ -643,6 +639,7 @@ func (s *trackSubscription) setSubscribedTrack(track types.SubscribedTrack) {
s.lock.Unlock()
if settings != nil && track != nil {
s.logger.Debugw("restoring subscriber settings", "settings", settings)
track.UpdateSubscriberSettings(settings)
}
}
+1
View File
@@ -198,6 +198,7 @@ func TestUnsubscribe(t *testing.T) {
publisherIdentity: "pub",
hasPermission: true,
bound: true,
logger: logger.GetLogger(),
}
// a bunch of unfortunate manual wiring
res, err := resolver.Resolve("sub", s.publisherID, s.trackID)