From 4f9467040e9110cf7bda23c38d62510f96c732c5 Mon Sep 17 00:00:00 2001 From: Pingos Date: Tue, 26 Sep 2023 13:36:54 +0800 Subject: [PATCH 01/12] Bind() function fails when mime == "audio/red" (#2104) --- pkg/sfu/buffer/buffer.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/sfu/buffer/buffer.go b/pkg/sfu/buffer/buffer.go index f27e449a2..f01c8f469 100644 --- a/pkg/sfu/buffer/buffer.go +++ b/pkg/sfu/buffer/buffer.go @@ -252,7 +252,7 @@ func (b *Buffer) Bind(params webrtc.RTPParameters, codec webrtc.RTPCodecCapabili // codec settings at track level for same codec type, so enable nack for all audio receivers but don't create nack queue // for red codec. if strings.EqualFold(b.mime, "audio/red") { - return + break } b.logger.Debugw("Setting feedback", "type", webrtc.TypeRTCPFBNACK) b.nacker = nack.NewNACKQueue(nack.NackQueueParamsDefault) From 789a2bd39013a23ff30fcd330ea2c87de0d8c9e2 Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Tue, 26 Sep 2023 13:21:08 +0530 Subject: [PATCH 02/12] Log ICE candidates when closing the peer connection. (#2105) There are cases that experience the signalling channel timeout and disconnect and there are no logs of what the state of ICE at all. Log ICE candidates when closing transport so that there is some visibility in those cases. n --- pkg/rtc/transport.go | 1 + 1 file changed, 1 insertion(+) diff --git a/pkg/rtc/transport.go b/pkg/rtc/transport.go index 517c90c1f..47338e812 100644 --- a/pkg/rtc/transport.go +++ b/pkg/rtc/transport.go @@ -1443,6 +1443,7 @@ func (t *PCTransport) processEvents() { t.clearSignalStateCheckTimer() t.params.Logger.Debugw("leaving events processor") + t.handleLogICECandidates(nil) } func (t *PCTransport) handleEvent(e *event) error { From 4c9d95149d58836baa621c169610a7824826b18e Mon Sep 17 00:00:00 2001 From: David Zhao Date: Wed, 27 Sep 2023 21:16:13 -0700 Subject: [PATCH 03/12] Enable SVC codecs by default (#2109) --- README.md | 2 ++ config-sample.yaml | 4 ++-- pkg/config/config.go | 4 ++-- 3 files changed, 6 insertions(+), 4 deletions(-) diff --git a/README.md b/README.md index 3692b94e5..524d8d98b 100644 --- a/README.md +++ b/README.md @@ -33,6 +33,8 @@ LiveKit's server is written in Go, using the awesome [Pion WebRTC](https://githu - [end-to-end optimizations](https://blog.livekit.io/livekit-one-dot-zero/) - [selective subscription](https://docs.livekit.io/guides/room/receive/#selective-subscription) - [moderation APIs](https://docs.livekit.io/guides/server-api/) + - end-to-end encryption + - SVC codecs (VP9, AV1) - [webhooks](https://docs.livekit.io/guides/webhooks/) - [distributed and multi-region](https://docs.livekit.io/deploy/distributed/) diff --git a/config-sample.yaml b/config-sample.yaml index 15c7ada18..d7a4fe773 100644 --- a/config-sample.yaml +++ b/config-sample.yaml @@ -43,7 +43,7 @@ rtc: # that maps to an external one use_external_ip: true # # when set, LiveKit will attempt to use a UDP mux so all UDP traffic goes through - # # listed port(s). To maximize system performance, we recommend using a range of ports + # # listed port(s). To maximize system performance, we recommend using a range of ports # # greater or equal to the number of vCPUs on the machine. # # port_range_start & end must not be set for this config to take effect # udp_port: 7882-7892 @@ -148,7 +148,7 @@ keys: # max_participants: 0 # # only accept specific codecs for clients publishing to this room # # this is useful to standardize codecs across clients -# # other supported codecs are video/h264 +# # other supported codecs are video/h264, video/vp9, video/av1, audio/red # enabled_codecs: # - mime: audio/opus # - mime: video/vp8 diff --git a/pkg/config/config.go b/pkg/config/config.go index 62142bc41..b99be109e 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -459,8 +459,8 @@ var DefaultConfig = Config{ {Mime: "audio/red"}, {Mime: webrtc.MimeTypeVP8}, {Mime: webrtc.MimeTypeH264}, - // {Mime: webrtc.MimeTypeAV1}, - // {Mime: webrtc.MimeTypeVP9}, + {Mime: webrtc.MimeTypeVP9}, + {Mime: webrtc.MimeTypeAV1}, }, EmptyTimeout: 5 * 60, }, From 92a355e1f3f7b7e8ddff43766ee4fa338922d74e Mon Sep 17 00:00:00 2001 From: cnderrauber Date: Thu, 28 Sep 2023 15:41:44 +0800 Subject: [PATCH 04/12] Add SyncStreams flag to Room (#2110) * Add SyncStreams flag to Room * Increase protocol version * Revert version change * Move flags to internal & solve comment --- config-sample.yaml | 3 + go.mod | 2 +- go.sum | 6 +- pkg/config/config.go | 1 + pkg/rtc/mediatracksubscriptions.go | 4 +- pkg/rtc/participant.go | 11 +- pkg/rtc/types/interfaces.go | 3 +- .../typesfakes/fake_local_participant.go | 135 +++++++++++++----- pkg/service/roomallocator.go | 15 +- pkg/service/roommanager.go | 5 +- 10 files changed, 132 insertions(+), 53 deletions(-) diff --git a/config-sample.yaml b/config-sample.yaml index d7a4fe773..6025a3a9b 100644 --- a/config-sample.yaml +++ b/config-sample.yaml @@ -162,6 +162,9 @@ keys: # enabled: true # min: 100 # max: 2000 +# # improves A/V sync when playout_delay set to a value larger than 200ms. It will disables transceiver re-use +# # so not recommended for rooms with frequent subscription changes +# sync_streams: true # Webhooks # when configured, LiveKit notifies your URL handler with room events diff --git a/go.mod b/go.mod index bb0795cb2..d31a29bdd 100644 --- a/go.mod +++ b/go.mod @@ -18,7 +18,7 @@ require ( github.com/jxskiss/base62 v1.1.0 github.com/livekit/mageutil v0.0.0-20230125210925-54e8a70427c1 github.com/livekit/mediatransportutil v0.0.0-20230919184714-b8f0fa0133c5 - github.com/livekit/protocol v1.7.3-0.20230920084913-821c244d8ce2 + github.com/livekit/protocol v1.7.3-0.20230928065809-281e00a4a67d github.com/livekit/psrpc v0.3.3 github.com/mackerelio/go-osstat v0.2.4 github.com/magefile/mage v1.15.0 diff --git a/go.sum b/go.sum index 757166a60..b08db3a89 100644 --- a/go.sum +++ b/go.sum @@ -127,10 +127,8 @@ github.com/livekit/mageutil v0.0.0-20230125210925-54e8a70427c1 h1:jm09419p0lqTkD github.com/livekit/mageutil v0.0.0-20230125210925-54e8a70427c1/go.mod h1:Rs3MhFwutWhGwmY1VQsygw28z5bWcnEYmS1OG9OxjOQ= github.com/livekit/mediatransportutil v0.0.0-20230919184714-b8f0fa0133c5 h1:CjXYkNKSrdIn7GMD1ySoXrURhL5U9d6vG32vxcUhzIU= github.com/livekit/mediatransportutil v0.0.0-20230919184714-b8f0fa0133c5/go.mod h1:+WIOYwiBMive5T81V8B2wdAc2zQNRjNQiJIcPxMTILY= -github.com/livekit/protocol v1.7.3-0.20230919182418-0708b5a5bb84 h1:4WOaspDesbbrjGPsu6Vp1VjcWxTXVjyjdtQAzIUXn5s= -github.com/livekit/protocol v1.7.3-0.20230919182418-0708b5a5bb84/go.mod h1:zbh0QPUcLGOeZeIO/VeigwWWbudz4Lv+Px94FnVfQH0= -github.com/livekit/protocol v1.7.3-0.20230920084913-821c244d8ce2 h1:yIRqvyO3qDPO+4EdcHMjsINJYl6KE9AXJUgfChsw+0s= -github.com/livekit/protocol v1.7.3-0.20230920084913-821c244d8ce2/go.mod h1:zbh0QPUcLGOeZeIO/VeigwWWbudz4Lv+Px94FnVfQH0= +github.com/livekit/protocol v1.7.3-0.20230928065809-281e00a4a67d h1:JLc/seGGKdnv0JUDCkMprJYzud2E8ahQ3QZgP/Imb14= +github.com/livekit/protocol v1.7.3-0.20230928065809-281e00a4a67d/go.mod h1:zbh0QPUcLGOeZeIO/VeigwWWbudz4Lv+Px94FnVfQH0= github.com/livekit/psrpc v0.3.3 h1:+lltbuN39IdaynXhLLxRShgYqYsRMWeeXKzv60oqyWo= github.com/livekit/psrpc v0.3.3/go.mod h1:n6JntEg+zT6Ji8InoyTpV7wusPNwGqqtxmHlkNhDN0U= github.com/mackerelio/go-osstat v0.2.4 h1:qxGbdPkFo65PXOb/F/nhDKpF2nGmGaCFDLXoZjJTtUs= diff --git a/pkg/config/config.go b/pkg/config/config.go index b99be109e..7997e7f93 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -222,6 +222,7 @@ type RoomConfig struct { EnableRemoteUnmute bool `yaml:"enable_remote_unmute,omitempty"` MaxMetadataSize uint32 `yaml:"max_metadata_size,omitempty"` PlayoutDelay PlayoutDelayConfig `yaml:"playout_delay,omitempty"` + SyncStreams bool `yaml:"sync_streams,omitempty"` } type CodecSpec struct { diff --git a/pkg/rtc/mediatracksubscriptions.go b/pkg/rtc/mediatracksubscriptions.go index 552fc06ce..33ca4921f 100644 --- a/pkg/rtc/mediatracksubscriptions.go +++ b/pkg/rtc/mediatracksubscriptions.go @@ -115,7 +115,7 @@ func (t *MediaTrackSubscriptions) AddSubscriber(sub types.LocalParticipant, wr * } streamID := wr.StreamID() - if sub.SupportSyncStreamID() && t.params.MediaTrack.Stream() != "" { + if sub.SupportsSyncStreamID() && t.params.MediaTrack.Stream() != "" { streamID = PackSyncStreamID(t.params.MediaTrack.PublisherID(), t.params.MediaTrack.Stream()) } @@ -246,7 +246,7 @@ func (t *MediaTrackSubscriptions) AddSubscriber(sub types.LocalParticipant, wr * } sub.VerifySubscribeParticipantInfo(subTrack.PublisherID(), subTrack.PublisherVersion()) - if sub.ProtocolVersion().SupportsTransceiverReuse() { + if sub.SupportsTransceiverReuse() { // // AddTrack will create a new transceiver or re-use an unused one // if the attributes match. This prevents SDP from bloating diff --git a/pkg/rtc/participant.go b/pkg/rtc/participant.go index 30dfea7d0..2d9d21cb7 100644 --- a/pkg/rtc/participant.go +++ b/pkg/rtc/participant.go @@ -122,6 +122,7 @@ type ParticipantParams struct { SubscriptionLimitAudio int32 SubscriptionLimitVideo int32 PlayoutDelay *livekit.PlayoutDelay + SyncStreams bool } type ParticipantImpl struct { @@ -1113,7 +1114,7 @@ func (p *ParticipantImpl) setupTransportManager() error { TCPFallbackRTTThreshold: p.params.TCPFallbackRTTThreshold, AllowUDPUnstableFallback: p.params.AllowUDPUnstableFallback, TURNSEnabled: p.params.TURNSEnabled, - AllowPlayoutDelay: p.params.PlayoutDelay.GetEnabled() && p.SupportSyncStreamID(), + AllowPlayoutDelay: p.params.PlayoutDelay.GetEnabled() && p.SupportsSyncStreamID(), Logger: p.params.Logger.WithComponent(sutils.ComponentTransport), }) if err != nil { @@ -2250,8 +2251,12 @@ func (p *ParticipantImpl) GetPlayoutDelayConfig() *livekit.PlayoutDelay { return p.params.PlayoutDelay } -func (p *ParticipantImpl) SupportSyncStreamID() bool { - return p.ProtocolVersion().SupportSyncStreamID() && !p.params.ClientInfo.isFirefox() +func (p *ParticipantImpl) SupportsSyncStreamID() bool { + return p.ProtocolVersion().SupportSyncStreamID() && !p.params.ClientInfo.isFirefox() && p.params.SyncStreams +} + +func (p *ParticipantImpl) SupportsTransceiverReuse() bool { + return p.ProtocolVersion().SupportsTransceiverReuse() && !p.SupportsSyncStreamID() } func codecsFromMediaDescription(m *sdp.MediaDescription) (out []sdp.Codec, err error) { diff --git a/pkg/rtc/types/interfaces.go b/pkg/rtc/types/interfaces.go index 0b0e003ad..3a971fce4 100644 --- a/pkg/rtc/types/interfaces.go +++ b/pkg/rtc/types/interfaces.go @@ -303,7 +303,8 @@ type LocalParticipant interface { GetLogger() logger.Logger GetAdaptiveStream() bool ProtocolVersion() ProtocolVersion - SupportSyncStreamID() bool + SupportsSyncStreamID() bool + SupportsTransceiverReuse() bool ConnectedAt() time.Time IsClosed() bool IsReady() bool diff --git a/pkg/rtc/types/typesfakes/fake_local_participant.go b/pkg/rtc/types/typesfakes/fake_local_participant.go index fedebacf5..ac0b18932 100644 --- a/pkg/rtc/types/typesfakes/fake_local_participant.go +++ b/pkg/rtc/types/typesfakes/fake_local_participant.go @@ -793,14 +793,24 @@ type FakeLocalParticipant struct { arg2 livekit.TrackID arg3 bool } - SupportSyncStreamIDStub func() bool - supportSyncStreamIDMutex sync.RWMutex - supportSyncStreamIDArgsForCall []struct { + SupportsSyncStreamIDStub func() bool + supportsSyncStreamIDMutex sync.RWMutex + supportsSyncStreamIDArgsForCall []struct { } - supportSyncStreamIDReturns struct { + supportsSyncStreamIDReturns struct { result1 bool } - supportSyncStreamIDReturnsOnCall map[int]struct { + supportsSyncStreamIDReturnsOnCall map[int]struct { + result1 bool + } + SupportsTransceiverReuseStub func() bool + supportsTransceiverReuseMutex sync.RWMutex + supportsTransceiverReuseArgsForCall []struct { + } + supportsTransceiverReuseReturns struct { + result1 bool + } + supportsTransceiverReuseReturnsOnCall map[int]struct { result1 bool } ToProtoStub func() *livekit.ParticipantInfo @@ -5153,15 +5163,15 @@ func (fake *FakeLocalParticipant) SubscriptionPermissionUpdateArgsForCall(i int) return argsForCall.arg1, argsForCall.arg2, argsForCall.arg3 } -func (fake *FakeLocalParticipant) SupportSyncStreamID() bool { - fake.supportSyncStreamIDMutex.Lock() - ret, specificReturn := fake.supportSyncStreamIDReturnsOnCall[len(fake.supportSyncStreamIDArgsForCall)] - fake.supportSyncStreamIDArgsForCall = append(fake.supportSyncStreamIDArgsForCall, struct { +func (fake *FakeLocalParticipant) SupportsSyncStreamID() bool { + fake.supportsSyncStreamIDMutex.Lock() + ret, specificReturn := fake.supportsSyncStreamIDReturnsOnCall[len(fake.supportsSyncStreamIDArgsForCall)] + fake.supportsSyncStreamIDArgsForCall = append(fake.supportsSyncStreamIDArgsForCall, struct { }{}) - stub := fake.SupportSyncStreamIDStub - fakeReturns := fake.supportSyncStreamIDReturns - fake.recordInvocation("SupportSyncStreamID", []interface{}{}) - fake.supportSyncStreamIDMutex.Unlock() + stub := fake.SupportsSyncStreamIDStub + fakeReturns := fake.supportsSyncStreamIDReturns + fake.recordInvocation("SupportsSyncStreamID", []interface{}{}) + fake.supportsSyncStreamIDMutex.Unlock() if stub != nil { return stub() } @@ -5171,37 +5181,90 @@ func (fake *FakeLocalParticipant) SupportSyncStreamID() bool { return fakeReturns.result1 } -func (fake *FakeLocalParticipant) SupportSyncStreamIDCallCount() int { - fake.supportSyncStreamIDMutex.RLock() - defer fake.supportSyncStreamIDMutex.RUnlock() - return len(fake.supportSyncStreamIDArgsForCall) +func (fake *FakeLocalParticipant) SupportsSyncStreamIDCallCount() int { + fake.supportsSyncStreamIDMutex.RLock() + defer fake.supportsSyncStreamIDMutex.RUnlock() + return len(fake.supportsSyncStreamIDArgsForCall) } -func (fake *FakeLocalParticipant) SupportSyncStreamIDCalls(stub func() bool) { - fake.supportSyncStreamIDMutex.Lock() - defer fake.supportSyncStreamIDMutex.Unlock() - fake.SupportSyncStreamIDStub = stub +func (fake *FakeLocalParticipant) SupportsSyncStreamIDCalls(stub func() bool) { + fake.supportsSyncStreamIDMutex.Lock() + defer fake.supportsSyncStreamIDMutex.Unlock() + fake.SupportsSyncStreamIDStub = stub } -func (fake *FakeLocalParticipant) SupportSyncStreamIDReturns(result1 bool) { - fake.supportSyncStreamIDMutex.Lock() - defer fake.supportSyncStreamIDMutex.Unlock() - fake.SupportSyncStreamIDStub = nil - fake.supportSyncStreamIDReturns = struct { +func (fake *FakeLocalParticipant) SupportsSyncStreamIDReturns(result1 bool) { + fake.supportsSyncStreamIDMutex.Lock() + defer fake.supportsSyncStreamIDMutex.Unlock() + fake.SupportsSyncStreamIDStub = nil + fake.supportsSyncStreamIDReturns = struct { result1 bool }{result1} } -func (fake *FakeLocalParticipant) SupportSyncStreamIDReturnsOnCall(i int, result1 bool) { - fake.supportSyncStreamIDMutex.Lock() - defer fake.supportSyncStreamIDMutex.Unlock() - fake.SupportSyncStreamIDStub = nil - if fake.supportSyncStreamIDReturnsOnCall == nil { - fake.supportSyncStreamIDReturnsOnCall = make(map[int]struct { +func (fake *FakeLocalParticipant) SupportsSyncStreamIDReturnsOnCall(i int, result1 bool) { + fake.supportsSyncStreamIDMutex.Lock() + defer fake.supportsSyncStreamIDMutex.Unlock() + fake.SupportsSyncStreamIDStub = nil + if fake.supportsSyncStreamIDReturnsOnCall == nil { + fake.supportsSyncStreamIDReturnsOnCall = make(map[int]struct { result1 bool }) } - fake.supportSyncStreamIDReturnsOnCall[i] = struct { + fake.supportsSyncStreamIDReturnsOnCall[i] = struct { + result1 bool + }{result1} +} + +func (fake *FakeLocalParticipant) SupportsTransceiverReuse() bool { + fake.supportsTransceiverReuseMutex.Lock() + ret, specificReturn := fake.supportsTransceiverReuseReturnsOnCall[len(fake.supportsTransceiverReuseArgsForCall)] + fake.supportsTransceiverReuseArgsForCall = append(fake.supportsTransceiverReuseArgsForCall, struct { + }{}) + stub := fake.SupportsTransceiverReuseStub + fakeReturns := fake.supportsTransceiverReuseReturns + fake.recordInvocation("SupportsTransceiverReuse", []interface{}{}) + fake.supportsTransceiverReuseMutex.Unlock() + if stub != nil { + return stub() + } + if specificReturn { + return ret.result1 + } + return fakeReturns.result1 +} + +func (fake *FakeLocalParticipant) SupportsTransceiverReuseCallCount() int { + fake.supportsTransceiverReuseMutex.RLock() + defer fake.supportsTransceiverReuseMutex.RUnlock() + return len(fake.supportsTransceiverReuseArgsForCall) +} + +func (fake *FakeLocalParticipant) SupportsTransceiverReuseCalls(stub func() bool) { + fake.supportsTransceiverReuseMutex.Lock() + defer fake.supportsTransceiverReuseMutex.Unlock() + fake.SupportsTransceiverReuseStub = stub +} + +func (fake *FakeLocalParticipant) SupportsTransceiverReuseReturns(result1 bool) { + fake.supportsTransceiverReuseMutex.Lock() + defer fake.supportsTransceiverReuseMutex.Unlock() + fake.SupportsTransceiverReuseStub = nil + fake.supportsTransceiverReuseReturns = struct { + result1 bool + }{result1} +} + +func (fake *FakeLocalParticipant) SupportsTransceiverReuseReturnsOnCall(i int, result1 bool) { + fake.supportsTransceiverReuseMutex.Lock() + defer fake.supportsTransceiverReuseMutex.Unlock() + fake.SupportsTransceiverReuseStub = nil + if fake.supportsTransceiverReuseReturnsOnCall == nil { + fake.supportsTransceiverReuseReturnsOnCall = make(map[int]struct { + result1 bool + }) + } + fake.supportsTransceiverReuseReturnsOnCall[i] = struct { result1 bool }{result1} } @@ -6033,8 +6096,10 @@ func (fake *FakeLocalParticipant) Invocations() map[string][][]interface{} { defer fake.subscriptionPermissionMutex.RUnlock() fake.subscriptionPermissionUpdateMutex.RLock() defer fake.subscriptionPermissionUpdateMutex.RUnlock() - fake.supportSyncStreamIDMutex.RLock() - defer fake.supportSyncStreamIDMutex.RUnlock() + fake.supportsSyncStreamIDMutex.RLock() + defer fake.supportsSyncStreamIDMutex.RUnlock() + fake.supportsTransceiverReuseMutex.RLock() + defer fake.supportsTransceiverReuseMutex.RUnlock() fake.toProtoMutex.RLock() defer fake.toProtoMutex.RUnlock() fake.toProtoWithVersionMutex.RLock() diff --git a/pkg/service/roomallocator.go b/pkg/service/roomallocator.go index 6db962549..2d538cdd1 100644 --- a/pkg/service/roomallocator.go +++ b/pkg/service/roomallocator.go @@ -68,7 +68,8 @@ func (r *StandardRoomAllocator) CreateRoom(ctx context.Context, req *livekit.Cre CreationTime: time.Now().Unix(), TurnPassword: utils.RandomSecret(), } - applyDefaultRoomConfig(rm, &r.config.Room) + internal = &livekit.RoomInternal{} + applyDefaultRoomConfig(rm, internal, &r.config.Room) } else if err != nil { return nil, err } @@ -83,15 +84,18 @@ func (r *StandardRoomAllocator) CreateRoom(ctx context.Context, req *livekit.Cre rm.Metadata = req.Metadata } if req.Egress != nil && req.Egress.Tracks != nil { - internal = &livekit.RoomInternal{TrackEgress: req.Egress.Tracks} + internal.TrackEgress = req.Egress.Tracks } if req.MinPlayoutDelay > 0 || req.MaxPlayoutDelay > 0 { - rm.PlayoutDelay = &livekit.PlayoutDelay{ + internal.PlayoutDelay = &livekit.PlayoutDelay{ Enabled: true, Min: req.MinPlayoutDelay, Max: req.MaxPlayoutDelay, } } + if req.SyncStreams { + internal.SyncStreams = true + } if err = r.roomStore.StoreRoom(ctx, rm, internal); err != nil { return nil, err @@ -149,7 +153,7 @@ func (r *StandardRoomAllocator) ValidateCreateRoom(ctx context.Context, roomName return nil } -func applyDefaultRoomConfig(room *livekit.Room, conf *config.RoomConfig) { +func applyDefaultRoomConfig(room *livekit.Room, internal *livekit.RoomInternal, conf *config.RoomConfig) { room.EmptyTimeout = conf.EmptyTimeout room.MaxParticipants = conf.MaxParticipants for _, codec := range conf.EnabledCodecs { @@ -158,9 +162,10 @@ func applyDefaultRoomConfig(room *livekit.Room, conf *config.RoomConfig) { FmtpLine: codec.FmtpLine, }) } - room.PlayoutDelay = &livekit.PlayoutDelay{ + internal.PlayoutDelay = &livekit.PlayoutDelay{ Enabled: conf.PlayoutDelay.Enabled, Min: uint32(conf.PlayoutDelay.Min), Max: uint32(conf.PlayoutDelay.Max), } + internal.SyncStreams = conf.SyncStreams } diff --git a/pkg/service/roommanager.go b/pkg/service/roommanager.go index b436c7bca..5797142e3 100644 --- a/pkg/service/roommanager.go +++ b/pkg/service/roommanager.go @@ -240,7 +240,7 @@ func (r *RoomManager) StartSession( } defer room.Release() - protoRoom := room.ToProto() + protoRoom, roomInternal := room.ToProto(), room.Internal() // only create the room, but don't start a participant session if pi.Identity == "" { @@ -408,7 +408,8 @@ func (r *RoomManager) StartSession( SubscriberAllowPause: subscriberAllowPause, SubscriptionLimitAudio: r.config.Limit.SubscriptionLimitAudio, SubscriptionLimitVideo: r.config.Limit.SubscriptionLimitVideo, - PlayoutDelay: protoRoom.PlayoutDelay, + PlayoutDelay: roomInternal.GetPlayoutDelay(), + SyncStreams: roomInternal.GetSyncStreams(), }) if err != nil { return err From 39d0af43816d14b9266173d9b5a0d66bc2a32bc4 Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Thu, 28 Sep 2023 18:44:48 +0530 Subject: [PATCH 05/12] Start key frame requester on start. (#2111) Sending a single PLI on connected & bound meant that the upstream throttler may not have sent it and down stream does not have a key frame to lock onto. Caused some e2e test failures due to limited time of track. --- pkg/sfu/downtrack.go | 7 ++----- 1 file changed, 2 insertions(+), 5 deletions(-) diff --git a/pkg/sfu/downtrack.go b/pkg/sfu/downtrack.go index 0fe32361a..929386f02 100644 --- a/pkg/sfu/downtrack.go +++ b/pkg/sfu/downtrack.go @@ -1782,19 +1782,16 @@ func (d *DownTrack) GetAndResetBytesSent() (uint32, uint32) { } func (d *DownTrack) onBindAndConnectedChange() { + d.writable.Store(d.connected.Load() && d.bound.Load()) if d.connected.Load() && d.bound.Load() && !d.bindAndConnectedOnce.Swap(true) { if d.kind == webrtc.RTPCodecTypeVideo { - _, layer := d.forwarder.CheckSync() - if layer != buffer.InvalidLayerSpatial { - d.params.Receiver.SendPLI(layer, true) - } + d.maybeStartKeyFrameRequester() } if d.activePaddingOnMuteUpTrack.Load() { go d.sendPaddingOnMute() } } - d.writable.Store(d.connected.Load() && d.bound.Load()) } func (d *DownTrack) sendPaddingOnMute() { From ee3a7c01bcb864aad348c4b1b6734d0b0d5688da Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Sat, 30 Sep 2023 08:42:46 +0530 Subject: [PATCH 06/12] Log resync at Infow. (#2114) * Log resync at Infow. Seeing potentially large sequence number jumps on a resync. And it seems to happen on a lot of subscribe/unsubscribe. Logging at Infow to understand better. Probably need to find a way to avoid resync. But, logging for now to check if I can catch one. * Remove resync and log large sequence number jumps --- pkg/sfu/buffer/rtpstats_base.go | 10 +-- pkg/sfu/buffer/rtpstats_receiver.go | 100 ++-------------------------- pkg/sfu/buffer/rtpstats_sender.go | 4 ++ 3 files changed, 15 insertions(+), 99 deletions(-) diff --git a/pkg/sfu/buffer/rtpstats_base.go b/pkg/sfu/buffer/rtpstats_base.go index d25341c8b..a024d18c5 100644 --- a/pkg/sfu/buffer/rtpstats_base.go +++ b/pkg/sfu/buffer/rtpstats_base.go @@ -106,11 +106,11 @@ type snapshot struct { } type RTCPSenderReportData struct { - RTPTimestamp uint32 - RTPTimestampExt uint64 - NTPTimestamp mediatransportutil.NtpTime - PacketCount uint32 - PacketCountExt uint64 + RTPTimestamp uint32 + RTPTimestampExt uint64 + NTPTimestamp mediatransportutil.NtpTime + PacketCount uint32 + // RAJA-REMOVE PacketCountExt uint64 PaddingOnlyDrops uint64 At time.Time } diff --git a/pkg/sfu/buffer/rtpstats_receiver.go b/pkg/sfu/buffer/rtpstats_receiver.go index fff35b2f5..ae42dc426 100644 --- a/pkg/sfu/buffer/rtpstats_receiver.go +++ b/pkg/sfu/buffer/rtpstats_receiver.go @@ -46,9 +46,6 @@ type RTPFlowState struct { type RTPStatsReceiver struct { *rtpStatsBase - resyncOnNextPacket bool - shouldDiscountPaddingOnlyDrops bool - sequenceNumber *utils.WrapAround[uint16, uint64] timestamp *utils.WrapAround[uint32, uint64] @@ -89,11 +86,6 @@ func (r *RTPStatsReceiver) Update( return } - if r.resyncOnNextPacket { - r.resyncOnNextPacket = false - r.resync(packetTime, sequenceNumber, timestamp) - } - var resSN utils.WrapAroundUpdateResult[uint64] var resTS utils.WrapAroundUpdateResult[uint64] if !r.initialized { @@ -136,11 +128,11 @@ func (r *RTPStatsReceiver) Update( if payloadSize == 0 { // do not start on a padding only packet if resTS.IsRestart { - r.logger.Infow("rolling back timestamp restart", "tsBefore", r.timestamp.GetExtendedStart(), "tsAfter", resTS.PreExtendedStart) + r.logger.Infow("rolling back timestamp restart", "tsAfter", r.timestamp.GetExtendedStart(), "tsBefore", resTS.PreExtendedStart) r.timestamp.RollbackRestart(resTS.PreExtendedStart) } if resSN.IsRestart { - r.logger.Infow("rolling back sequence number restart", "snBefore", r.sequenceNumber.GetExtendedStart(), "snAfter", resSN.PreExtendedStart) + r.logger.Infow("rolling back sequence number restart", "snAfter", r.sequenceNumber.GetExtendedStart(), "snBefore", resSN.PreExtendedStart) r.sequenceNumber.RollbackRestart(resSN.PreExtendedStart) return } @@ -192,6 +184,10 @@ func (r *RTPStatsReceiver) Update( flowState.ExtSequenceNumber = resSN.ExtendedVal flowState.ExtTimestamp = resTS.ExtendedVal } else { // in-order + if gapSN >= cNumSequenceNumbers { + r.logger.Warnw("large sequence number gap", nil, "prev", resSN.PreExtendedHighest, "curr", resSN.ExtendedVal, "gap", gapSN) + } + // update gap histogram r.updateGapHistogram(int(gapSN)) @@ -235,83 +231,6 @@ func (r *RTPStatsReceiver) Update( return } -func (r *RTPStatsReceiver) ResyncOnNextPacket(shouldDiscountPaddingOnlyDrops bool) { - r.lock.Lock() - defer r.lock.Unlock() - - r.resyncOnNextPacket = true - r.shouldDiscountPaddingOnlyDrops = shouldDiscountPaddingOnlyDrops -} - -func (r *RTPStatsReceiver) resync(packetTime time.Time, sn uint16, ts uint32) { - if !r.initialized { - return - } - - extHighestSN := r.sequenceNumber.GetExtendedHighest() - var newestPacketCount uint64 - var paddingOnlyDrops uint64 - var extExpectedHighestSN uint64 - var expectedHighestSN uint16 - var snCycles uint64 - - extHighestTS := r.timestamp.GetExtendedHighest() - var newestTS uint64 - var extExpectedHighestTS uint64 - var expectedHighestTS uint32 - var tsCycles uint64 - if r.srNewest != nil { - newestPacketCount = r.srNewest.PacketCountExt - paddingOnlyDrops = r.srNewest.PaddingOnlyDrops - if newestPacketCount != 0 { - extExpectedHighestSN = r.sequenceNumber.GetExtendedStart() + newestPacketCount - if r.shouldDiscountPaddingOnlyDrops { - extExpectedHighestSN -= paddingOnlyDrops - } - expectedHighestSN = uint16(extExpectedHighestSN & 0xFFFF) - snCycles = extExpectedHighestSN & 0xFFFF_FFFF_FFFF_0000 - if sn-expectedHighestSN < (1<<15) && sn < expectedHighestSN { - snCycles += (1 << 16) - } - if snCycles != 0 && expectedHighestSN-sn < (1<<15) && expectedHighestSN < sn { - snCycles -= (1 << 16) - } - } - - newestTS = r.srNewest.RTPTimestampExt - extExpectedHighestTS = newestTS - expectedHighestTS = uint32(extExpectedHighestTS & 0xFFFF_FFFF) - tsCycles = extExpectedHighestTS & 0xFFFF_FFFF_0000_0000 - if ts-expectedHighestTS < (1<<31) && ts < expectedHighestTS { - tsCycles += (1 << 32) - } - if tsCycles != 0 && expectedHighestTS-ts < (1<<31) && expectedHighestTS < ts { - tsCycles -= (1 << 32) - } - } - r.sequenceNumber.ResetHighest(snCycles + uint64(sn) - 1) - r.timestamp.ResetHighest(tsCycles + uint64(ts)) - r.highestTime = packetTime - r.logger.Debugw( - "resync", - "newestPacketCount", newestPacketCount, - "paddingOnlyDrops", paddingOnlyDrops, - "extExpectedHighestSN", extExpectedHighestSN, - "expectedHighestSN", expectedHighestSN, - "snCycles", snCycles, - "rtpSN", sn, - "beforeExtHighestSN", extHighestSN, - "afterExtHighestSN", r.sequenceNumber.GetExtendedHighest(), - "newestTS", newestTS, - "extExpectedHighestTS", extExpectedHighestTS, - "expectedHighestTS", expectedHighestTS, - "tsCycles", tsCycles, - "rtpTS", ts, - "beforeExtHighestTS", extHighestTS, - "afterExtHighestTS", r.timestamp.GetExtendedHighest(), - ) -} - func (r *RTPStatsReceiver) SetRtcpSenderReportData(srData *RTCPSenderReportData) { r.lock.Lock() defer r.lock.Unlock() @@ -335,22 +254,15 @@ func (r *RTPStatsReceiver) SetRtcpSenderReportData(srData *RTCPSenderReportData) } tsCycles := uint64(0) - pcCycles := uint64(0) if r.srNewest != nil { tsCycles = r.srNewest.RTPTimestampExt & 0xFFFF_FFFF_0000_0000 if (srData.RTPTimestamp-r.srNewest.RTPTimestamp) < (1<<31) && srData.RTPTimestamp < r.srNewest.RTPTimestamp { tsCycles += (1 << 32) } - - pcCycles = r.srNewest.PacketCountExt & 0xFFFF_FFFF_0000_0000 - if (srData.PacketCount-r.srNewest.PacketCount) < (1<<31) && srData.PacketCount < r.srNewest.PacketCount { - pcCycles += (1 << 32) - } } srDataCopy := *srData srDataCopy.RTPTimestampExt = uint64(srDataCopy.RTPTimestamp) + tsCycles - srDataCopy.PacketCountExt = uint64(srDataCopy.PacketCount) + pcCycles r.maybeAdjustFirstPacketTime(srDataCopy.RTPTimestampExt, r.timestamp.GetExtendedStart()) diff --git a/pkg/sfu/buffer/rtpstats_sender.go b/pkg/sfu/buffer/rtpstats_sender.go index 5a66c65d4..bb8ca227d 100644 --- a/pkg/sfu/buffer/rtpstats_sender.go +++ b/pkg/sfu/buffer/rtpstats_sender.go @@ -302,6 +302,10 @@ func (r *RTPStatsSender) Update( r.setSnInfo(extSequenceNumber, r.extHighestSN, uint16(pktSize), uint8(hdrSize), uint16(payloadSize), marker, true) } } else { // in-order + if gapSN >= cNumSequenceNumbers { + r.logger.Warnw("large sequence number gap", nil, "prev", r.extHighestSN, "curr", extSequenceNumber, "gap", gapSN) + } + // update gap histogram r.updateGapHistogram(int(gapSN)) From 180ad541fc7fe451379ce91d2aa6e65ba3bc95f6 Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Sat, 30 Sep 2023 10:24:13 +0530 Subject: [PATCH 07/12] Mark packet not handled if restart is rejected (#2115) * Mark packet not handled if restart is rejected * log both sn and ts on restart/rollback --- pkg/sfu/buffer/rtpstats_receiver.go | 21 +++++++++++++++++++-- 1 file changed, 19 insertions(+), 2 deletions(-) diff --git a/pkg/sfu/buffer/rtpstats_receiver.go b/pkg/sfu/buffer/rtpstats_receiver.go index ae42dc426..0bf9bb49b 100644 --- a/pkg/sfu/buffer/rtpstats_receiver.go +++ b/pkg/sfu/buffer/rtpstats_receiver.go @@ -128,12 +128,25 @@ func (r *RTPStatsReceiver) Update( if payloadSize == 0 { // do not start on a padding only packet if resTS.IsRestart { - r.logger.Infow("rolling back timestamp restart", "tsAfter", r.timestamp.GetExtendedStart(), "tsBefore", resTS.PreExtendedStart) + r.logger.Infow( + "rolling back timestamp restart", + "tsBefore", resTS.PreExtendedStart, + "tsAfter", r.timestamp.GetExtendedStart(), + "snBefore", resSN.PreExtendedStart, + "snAfter", r.sequenceNumber.GetExtendedStart(), + ) r.timestamp.RollbackRestart(resTS.PreExtendedStart) } if resSN.IsRestart { - r.logger.Infow("rolling back sequence number restart", "snAfter", r.sequenceNumber.GetExtendedStart(), "snBefore", resSN.PreExtendedStart) + r.logger.Infow( + "rolling back sequence number restart", + "snBefore", resSN.PreExtendedStart, + "snAfter", r.sequenceNumber.GetExtendedStart(), + "tsBefore", resTS.PreExtendedStart, + "tsAfter", r.timestamp.GetExtendedStart(), + ) r.sequenceNumber.RollbackRestart(resSN.PreExtendedStart) + flowState.IsNotHandled = true return } } @@ -157,6 +170,8 @@ func (r *RTPStatsReceiver) Update( "adjusting start sequence number", "snBefore", resSN.PreExtendedStart, "snAfter", resSN.ExtendedVal, + "tsBefore", resTS.PreExtendedStart, + "tsAfter", resTS.ExtendedVal, ) } @@ -165,6 +180,8 @@ func (r *RTPStatsReceiver) Update( "adjusting start timestamp", "tsBefore", resTS.PreExtendedStart, "tsAfter", resTS.ExtendedVal, + "snBefore", resSN.PreExtendedStart, + "snAfter", resSN.ExtendedVal, ) } From 99a4eeb64231cae8c5099517ebd5a2be0d277af0 Mon Sep 17 00:00:00 2001 From: bharat <114856815+sonigeez@users.noreply.github.com> Date: Sun, 1 Oct 2023 11:53:01 +0530 Subject: [PATCH 08/12] Update X's username in README.md (#2116) --- README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/README.md b/README.md index 524d8d98b..7bdc9863e 100644 --- a/README.md +++ b/README.md @@ -15,7 +15,7 @@ LiveKit's server is written in Go, using the awesome [Pion WebRTC](https://githu [![GitHub stars](https://img.shields.io/github/stars/livekit/livekit?style=social&label=Star&maxAge=2592000)](https://github.com/livekit/livekit/stargazers/) [![Slack community](https://img.shields.io/endpoint?url=https%3A%2F%2Flivekit.io%2Fbadges%2Fslack)](https://livekit.io/join-slack) -[![Twitter Follow](https://img.shields.io/twitter/follow/livekitted)](https://twitter.com/livekitted) +[![Twitter Follow](https://img.shields.io/twitter/follow/livekit)](https://twitter.com/livekit) [![GitHub release (latest SemVer)](https://img.shields.io/github/v/release/livekit/livekit)](https://github.com/livekit/livekit/releases/latest) [![GitHub Workflow Status](https://img.shields.io/github/actions/workflow/status/livekit/livekit/buildtest.yaml?branch=master)](https://github.com/livekit/livekit/actions/workflows/buildtest.yaml) [![License](https://img.shields.io/github/license/livekit/livekit)](https://github.com/livekit/livekit/blob/master/LICENSE) From 989d621c9746f1912fc44f8f8021a0388fdd773f Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Mon, 2 Oct 2023 11:51:46 +0530 Subject: [PATCH 09/12] A small change to move out of order check before RTT calc. (#2117) --- pkg/sfu/buffer/rtpstats_sender.go | 24 ++++++++++++------------ 1 file changed, 12 insertions(+), 12 deletions(-) diff --git a/pkg/sfu/buffer/rtpstats_sender.go b/pkg/sfu/buffer/rtpstats_sender.go index bb8ca227d..df4f7381d 100644 --- a/pkg/sfu/buffer/rtpstats_sender.go +++ b/pkg/sfu/buffer/rtpstats_sender.go @@ -375,18 +375,6 @@ func (r *RTPStatsSender) UpdateFromReceiverReport(rr rtcp.ReceptionReport) (rtt return } - var err error - if r.srNewest != nil { - rtt, err = mediatransportutil.GetRttMs(&rr, r.srNewest.NTPTimestamp, r.srNewest.At) - if err == nil { - isRttChanged = rtt != r.rtt - } else { - if !errors.Is(err, mediatransportutil.ErrRttNotLastSenderReport) && !errors.Is(err, mediatransportutil.ErrRttNoLastSenderReport) { - r.logger.Warnw("error getting rtt", err) - } - } - } - if !r.lastRRTime.IsZero() && r.extHighestSNFromRR > extHighestSNFromRR { r.logger.Debugw( fmt.Sprintf("receiver report potentially out of order, highestSN: existing: %d, received: %d", r.extHighestSNFromRR, extHighestSNFromRR), @@ -400,6 +388,18 @@ func (r *RTPStatsSender) UpdateFromReceiverReport(rr rtcp.ReceptionReport) (rtt r.extHighestSNFromRR = extHighestSNFromRR + if r.srNewest != nil { + var err error + rtt, err = mediatransportutil.GetRttMs(&rr, r.srNewest.NTPTimestamp, r.srNewest.At) + if err == nil { + isRttChanged = rtt != r.rtt + } else { + if !errors.Is(err, mediatransportutil.ErrRttNotLastSenderReport) && !errors.Is(err, mediatransportutil.ErrRttNoLastSenderReport) { + r.logger.Warnw("error getting rtt", err) + } + } + } + packetsLostFromRR := r.packetsLostFromRR&0xFFFF_FFFF_0000_0000 + uint64(rr.TotalLost) if (rr.TotalLost-r.lastRR.TotalLost) < (1<<31) && rr.TotalLost < r.lastRR.TotalLost { packetsLostFromRR += (1 << 32) From e6e3e2a729b3f38027b2d46c4a833589f73e8aa5 Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Mon, 2 Oct 2023 22:47:40 +0530 Subject: [PATCH 10/12] sligtly easier readability (#2121) --- pkg/sfu/buffer/buffer.go | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pkg/sfu/buffer/buffer.go b/pkg/sfu/buffer/buffer.go index f01c8f469..2b62b3427 100644 --- a/pkg/sfu/buffer/buffer.go +++ b/pkg/sfu/buffer/buffer.go @@ -426,9 +426,10 @@ func (b *Buffer) calc(pkt []byte, arrivalTime time.Time) { return } - flowState := b.updateStreamState(&rtpPacket, arrivalTime) // process header extensions always as padding packets could be used for probing b.processHeaderExtensions(&rtpPacket, arrivalTime) + + flowState := b.updateStreamState(&rtpPacket, arrivalTime) if flowState.IsNotHandled { return } From 96ccf696d3bcb396d8cd82b8e4e520d8c1701306 Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Tue, 3 Oct 2023 12:36:19 +0530 Subject: [PATCH 11/12] Cap expected packets to padding diff. (#2122) * Cap expected packets to padding diff. On the receiver, no longer using packet metadata cache to calculate interval stats. An optimisation to get rid of packet metadata cache on receiver side. Because of that, padding packets in an interval could be more than expected packets. As padding packets is just a counter, out-of-order padding packets will make the diff look larger than expected packets in a window. Cap the expected to 0. NOTE: This makes it so that the count is not accurate in a window, but that is okay occasionally. It will affect reported stats and quality calculations, but it should be rare. For example, if 30 packets were received in a window and 60 out-of-order padding packets were received, it would reported as 0 packets were received. One option is to not increment padding packets when they are out-of-order, but that will mess up overall stats. Will make that change if we see this happen a lot. * log unexpected padding packets --- go.mod | 2 +- go.sum | 4 ++-- pkg/sfu/buffer/rtpstats_base.go | 14 ++++++++++++-- 3 files changed, 15 insertions(+), 5 deletions(-) diff --git a/go.mod b/go.mod index d31a29bdd..39c2d94cc 100644 --- a/go.mod +++ b/go.mod @@ -17,7 +17,7 @@ require ( github.com/hashicorp/golang-lru/v2 v2.0.6 github.com/jxskiss/base62 v1.1.0 github.com/livekit/mageutil v0.0.0-20230125210925-54e8a70427c1 - github.com/livekit/mediatransportutil v0.0.0-20230919184714-b8f0fa0133c5 + github.com/livekit/mediatransportutil v0.0.0-20231003064835-a51e0ca1b1fd github.com/livekit/protocol v1.7.3-0.20230928065809-281e00a4a67d github.com/livekit/psrpc v0.3.3 github.com/mackerelio/go-osstat v0.2.4 diff --git a/go.sum b/go.sum index b08db3a89..c25e8fc89 100644 --- a/go.sum +++ b/go.sum @@ -125,8 +125,8 @@ github.com/lithammer/shortuuid/v4 v4.0.0 h1:QRbbVkfgNippHOS8PXDkti4NaWeyYfcBTHtw github.com/lithammer/shortuuid/v4 v4.0.0/go.mod h1:Zs8puNcrvf2rV9rTH51ZLLcj7ZXqQI3lv67aw4KiB1Y= github.com/livekit/mageutil v0.0.0-20230125210925-54e8a70427c1 h1:jm09419p0lqTkDaKb5iXdynYrzB84ErPPO4LbRASk58= github.com/livekit/mageutil v0.0.0-20230125210925-54e8a70427c1/go.mod h1:Rs3MhFwutWhGwmY1VQsygw28z5bWcnEYmS1OG9OxjOQ= -github.com/livekit/mediatransportutil v0.0.0-20230919184714-b8f0fa0133c5 h1:CjXYkNKSrdIn7GMD1ySoXrURhL5U9d6vG32vxcUhzIU= -github.com/livekit/mediatransportutil v0.0.0-20230919184714-b8f0fa0133c5/go.mod h1:+WIOYwiBMive5T81V8B2wdAc2zQNRjNQiJIcPxMTILY= +github.com/livekit/mediatransportutil v0.0.0-20231003064835-a51e0ca1b1fd h1:Dqfr8ehKCukEYE9zr0C035VAOmMq8lntE4iD/DE7504= +github.com/livekit/mediatransportutil v0.0.0-20231003064835-a51e0ca1b1fd/go.mod h1:+WIOYwiBMive5T81V8B2wdAc2zQNRjNQiJIcPxMTILY= github.com/livekit/protocol v1.7.3-0.20230928065809-281e00a4a67d h1:JLc/seGGKdnv0JUDCkMprJYzud2E8ahQ3QZgP/Imb14= github.com/livekit/protocol v1.7.3-0.20230928065809-281e00a4a67d/go.mod h1:zbh0QPUcLGOeZeIO/VeigwWWbudz4Lv+Px94FnVfQH0= github.com/livekit/psrpc v0.3.3 h1:+lltbuN39IdaynXhLLxRShgYqYsRMWeeXKzv60oqyWo= diff --git a/pkg/sfu/buffer/rtpstats_base.go b/pkg/sfu/buffer/rtpstats_base.go index a024d18c5..86b0f82c5 100644 --- a/pkg/sfu/buffer/rtpstats_base.go +++ b/pkg/sfu/buffer/rtpstats_base.go @@ -546,16 +546,26 @@ func (r *rtpStatsBase) deltaInfo(snapshotID uint32, extStartSN uint64, extHighes if int32(packetsLost) < 0 { packetsLost = 0 } + + // padding packets delta could be higher than expected due to out-of-order padding packets + packetsPadding := now.packetsPadding - then.packetsPadding + if packetsExpected < packetsPadding { + r.logger.Infow("padding packets more than expected", "packetsExpected", packetsExpected, "packetsPadding", packetsPadding) + packetsExpected = 0 + } else { + packetsExpected -= packetsPadding + } + return &RTPDeltaInfo{ StartTime: startTime, Duration: endTime.Sub(startTime), - Packets: uint32(packetsExpected - (now.packetsPadding - then.packetsPadding)), + Packets: uint32(packetsExpected), Bytes: now.bytes - then.bytes, HeaderBytes: now.headerBytes - then.headerBytes, PacketsDuplicate: uint32(now.packetsDuplicate - then.packetsDuplicate), BytesDuplicate: now.bytesDuplicate - then.bytesDuplicate, HeaderBytesDuplicate: now.headerBytesDuplicate - then.headerBytesDuplicate, - PacketsPadding: uint32(now.packetsPadding - then.packetsPadding), + PacketsPadding: uint32(packetsPadding), BytesPadding: now.bytesPadding - then.bytesPadding, HeaderBytesPadding: now.headerBytesPadding - then.headerBytesPadding, PacketsLost: packetsLost, From cb0a48c12d761ebd3a801d8d7bc996b9c44cfbe8 Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Tue, 3 Oct 2023 19:38:55 +0530 Subject: [PATCH 12/12] Handle RED extended sequence number. (#2123) When converting from RED -> Opus, if there is a loss, SFU recovers that loss if it can using a subsequent redundant packet. That path was not setting the extended sequence number properly. Also, ensuring use of monotonic clock for first packet time adjustment also. --- pkg/sfu/buffer/rtpstats_base.go | 3 ++- pkg/sfu/buffer/rtpstats_sender.go | 14 ++++++++++++ pkg/sfu/redprimaryreceiver.go | 38 +++++++++++++++++++++++++++---- 3 files changed, 49 insertions(+), 6 deletions(-) diff --git a/pkg/sfu/buffer/rtpstats_base.go b/pkg/sfu/buffer/rtpstats_base.go index 86b0f82c5..7438c9c72 100644 --- a/pkg/sfu/buffer/rtpstats_base.go +++ b/pkg/sfu/buffer/rtpstats_base.go @@ -474,7 +474,8 @@ func (r *rtpStatsBase) maybeAdjustFirstPacketTime(ets uint64, extStartTS uint64) } samplesDuration := time.Duration(float64(samplesDiff) / float64(r.params.ClockRate) * float64(time.Second)) - now := time.Now() + timeSinceFirst := time.Since(r.firstTime) + now := r.firstTime.Add(timeSinceFirst) firstTime := now.Add(-samplesDuration) if firstTime.Before(r.firstTime) { r.logger.Debugw( diff --git a/pkg/sfu/buffer/rtpstats_sender.go b/pkg/sfu/buffer/rtpstats_sender.go index df4f7381d..22392ec8a 100644 --- a/pkg/sfu/buffer/rtpstats_sender.go +++ b/pkg/sfu/buffer/rtpstats_sender.go @@ -281,10 +281,24 @@ func (r *RTPStatsSender) Update( } } + r.logger.Infow( + "adjusting start sequence number", + "snBefore", r.extStartSN, + "snAfter", extSequenceNumber, + "tsBefore", r.extStartTS, + "tsAfter", extTimestamp, + ) r.extStartSN = extSequenceNumber } if extTimestamp < r.extStartTS { + r.logger.Infow( + "adjusting start timestamp", + "snBefore", r.extStartSN, + "snAfter", extSequenceNumber, + "tsBefore", r.extStartTS, + "tsAfter", extTimestamp, + ) r.extStartTS = extTimestamp } diff --git a/pkg/sfu/redprimaryreceiver.go b/pkg/sfu/redprimaryreceiver.go index eb0965627..e7d33099f 100644 --- a/pkg/sfu/redprimaryreceiver.go +++ b/pkg/sfu/redprimaryreceiver.go @@ -69,8 +69,14 @@ func (r *RedPrimaryReceiver) ForwardRTP(pkt *buffer.ExtPacket, spatialLayer int3 return } - for _, sendPkt := range pkts { + for i, sendPkt := range pkts { pPkt := *pkt + if i != len(pkts)-1 { + // patch extended sequence number and time stmap for all but the last packet, + // last packet is the primary payload + pPkt.ExtSequenceNumber -= uint64(pkts[len(pkts)-1].SequenceNumber - pkts[i].SequenceNumber) + pPkt.ExtTimestamp -= uint64(pkts[len(pkts)-1].Timestamp - pkts[i].Timestamp) + } pPkt.Packet = sendPkt // not modify the ExtPacket.RawPacket here for performance since it is not used by the DownTrack, @@ -143,6 +149,7 @@ func (r *RedPrimaryReceiver) getSendPktsFromRed(rtp *rtp.Packet) ([]*rtp.Packet, switch { case diff == 0: // duplicate break + case diff > 0x8000: // unorder // in history if 65535-diff < 8 { @@ -191,25 +198,36 @@ func extractPktsFromRed(redPkt *rtp.Packet, recoverBits byte) ([]*rtp.Packet, er var blocks []block var blockLength int for { + if len(payload) < 1 { + // illegal data, need at least one byte for primary encoding + return nil, ErrIncompleteRedHeader + } + if payload[0]&0x80 == 0 { // last block is primary encoding data + pt := uint8(payload[0] & 0x7F) + + blocks = append(blocks, block{pt: pt, primary: true}) + payload = payload[1:] - blocks = append(blocks, block{primary: true}) break } else { if len(payload) < 4 { // illegal data return nil, ErrIncompleteRedHeader } + blockHead := binary.BigEndian.Uint32(payload[0:]) length := int(blockHead & 0x03FF) blockHead >>= 10 tsOffset := blockHead & 0x3FFF blockHead >>= 14 pt := uint8(blockHead & 0x7F) - payload = payload[4:] - blockLength += length + blocks = append(blocks, block{pt: pt, length: length, tsOffset: tsOffset}) + + blockLength += length + payload = payload[4:] } } @@ -220,22 +238,26 @@ func extractPktsFromRed(redPkt *rtp.Packet, recoverBits byte) ([]*rtp.Packet, er pkts := make([]*rtp.Packet, 0, len(blocks)) for i, b := range blocks { if b.primary { + header := redPkt.Header + header.PayloadType = b.pt pkts = append(pkts, &rtp.Packet{Header: redPkt.Header, Payload: payload}) break } - // last block is primary encoding recoverIndex := len(blocks) - i - 1 if recoverIndex < 1 || recoverBits&(1<<(recoverIndex-1)) == 0 { + // skip past packet/block that does not need recovery payload = payload[b.length:] continue } + // recover missing packet header := redPkt.Header header.SequenceNumber -= uint16(recoverIndex) header.Timestamp -= b.tsOffset header.PayloadType = b.pt pkts = append(pkts, &rtp.Packet{Header: header, Payload: payload[:b.length]}) + payload = payload[b.length:] } @@ -257,6 +279,11 @@ func extractPrimaryEncodingForRED(payload []byte) ([]byte, error) { var blockLength int for { + if len(payload) < 1 { + // illegal data, need at least one byte for primary encoding + return nil, ErrIncompleteRedHeader + } + if payload[0]&0x80 == 0 { // last block is primary encoding data payload = payload[1:] @@ -266,6 +293,7 @@ func extractPrimaryEncodingForRED(payload []byte) ([]byte, error) { // illegal data return nil, ErrIncompleteRedHeader } + blockLength += int(binary.BigEndian.Uint16(payload[2:]) & 0x03FF) payload = payload[4:] }