From 3e43f75143fa3d64573bb9163f1673eb6e4d759b Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Wed, 13 Mar 2024 14:31:39 +0530 Subject: [PATCH] Forward publisher sender report. (#2572) * Forward publisher sender report. Publisher side RTCP sernfer report is rebased to SFU time base and used to send sender rerport to subscriber. Will wait to merge till previous versions are out as this will require a bunch of testing. * - Add rebased report drift - update protocol dep - fix path change check, it has to check against delta of propagation delay and not propagation delay as the two side clocks could be way off. --- go.mod | 21 ++--- go.sum | 40 +++++----- pkg/rtc/wrappedreceiver.go | 4 +- pkg/sfu/buffer/buffer.go | 4 +- pkg/sfu/buffer/rtpstats_base.go | 70 ++++++++++------- pkg/sfu/buffer/rtpstats_receiver.go | 115 +++++++++++++++++++++++++--- pkg/sfu/buffer/rtpstats_sender.go | 96 ++++++++--------------- pkg/sfu/downtrack.go | 35 +++------ pkg/sfu/forwarder.go | 13 ++-- pkg/sfu/receiver.go | 10 +-- pkg/sfu/rtpmunger.go | 4 + pkg/sfu/streamtrackermanager.go | 24 ++---- 12 files changed, 248 insertions(+), 188 deletions(-) diff --git a/go.mod b/go.mod index 317336199..63a5c4fcb 100644 --- a/go.mod +++ b/go.mod @@ -19,8 +19,8 @@ require ( github.com/jxskiss/base62 v1.1.0 github.com/livekit/mageutil v0.0.0-20230125210925-54e8a70427c1 github.com/livekit/mediatransportutil v0.0.0-20240302142739-1c3dd691a1b8 - github.com/livekit/protocol v1.11.1-0.20240311174744-00c977ffbb49 - github.com/livekit/psrpc v0.5.3-0.20240228172457-3724cb4adbc4 + github.com/livekit/protocol v1.11.1-0.20240313083005-cf54792d0626 + github.com/livekit/psrpc v0.5.3-0.20240312110212-61ab09477c30 github.com/mackerelio/go-osstat v0.2.4 github.com/magefile/mage v1.15.0 github.com/maxbrunsfeld/counterfeiter/v6 v6.8.1 @@ -50,7 +50,7 @@ require ( go.uber.org/zap v1.27.0 golang.org/x/exp v0.0.0-20240222234643-814bf88cf225 golang.org/x/sync v0.6.0 - google.golang.org/protobuf v1.32.0 + google.golang.org/protobuf v1.33.0 gopkg.in/yaml.v3 v3.0.1 ) @@ -62,9 +62,10 @@ require ( github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f // indirect github.com/eapache/channels v1.1.0 // indirect github.com/eapache/queue v1.1.0 // indirect + github.com/fsnotify/fsnotify v1.7.0 // indirect github.com/go-jose/go-jose/v3 v3.0.3 // indirect github.com/go-logr/logr v1.4.1 // indirect - github.com/golang/protobuf v1.5.3 // indirect + github.com/golang/protobuf v1.5.4 // indirect github.com/google/go-cmp v0.6.0 // indirect github.com/google/subcommands v1.2.0 // indirect github.com/google/uuid v1.6.0 // indirect @@ -72,13 +73,13 @@ require ( github.com/hashicorp/go-retryablehttp v0.7.5 // indirect github.com/hashicorp/golang-lru v0.5.4 // indirect github.com/josharian/native v1.1.0 // indirect - github.com/klauspost/compress v1.17.6 // indirect + github.com/klauspost/compress v1.17.7 // indirect github.com/klauspost/cpuid/v2 v2.2.6 // indirect github.com/lithammer/shortuuid/v4 v4.0.0 // indirect github.com/mattn/go-runewidth v0.0.9 // indirect github.com/mdlayher/netlink v1.7.1 // indirect github.com/mdlayher/socket v0.4.0 // indirect - github.com/nats-io/nats.go v1.32.0 // indirect + github.com/nats-io/nats.go v1.33.1 // indirect github.com/nats-io/nkeys v0.4.7 // indirect github.com/nats-io/nuid v1.0.1 // indirect github.com/pion/datachannel v1.5.5 // indirect @@ -96,13 +97,13 @@ require ( github.com/xrash/smetrics v0.0.0-20201216005158-039620a65673 // indirect github.com/zeebo/xxh3 v1.0.2 // indirect go.uber.org/multierr v1.11.0 // indirect - golang.org/x/crypto v0.19.0 // indirect + golang.org/x/crypto v0.21.0 // indirect golang.org/x/mod v0.16.0 // indirect golang.org/x/net v0.21.0 // indirect - golang.org/x/sys v0.17.0 // indirect + golang.org/x/sys v0.18.0 // indirect golang.org/x/text v0.14.0 // indirect golang.org/x/tools v0.18.0 // indirect - google.golang.org/genproto/googleapis/rpc v0.0.0-20240221002015-b0ce06bbee7c // indirect - google.golang.org/grpc v1.62.0 // indirect + google.golang.org/genproto/googleapis/rpc v0.0.0-20240311173647-c811ad7063a7 // indirect + google.golang.org/grpc v1.62.1 // indirect gopkg.in/yaml.v2 v2.4.0 // indirect ) diff --git a/go.sum b/go.sum index fc8764a1a..8b2635d44 100644 --- a/go.sum +++ b/go.sum @@ -40,6 +40,8 @@ github.com/frostbyte73/core v0.0.10 h1:D4DQXdPb8ICayz0n75rs4UYTXrUSdxzUfeleuNJOR github.com/frostbyte73/core v0.0.10/go.mod h1:XsOGqrqe/VEV7+8vJ+3a8qnCIXNbKsoEiu/czs7nrcU= github.com/fsnotify/fsnotify v1.4.7/go.mod h1:jwhsz4b93w/PPRr/qN1Yymfu8t87LnFCMoQvtojpjFo= github.com/fsnotify/fsnotify v1.4.9/go.mod h1:znqG4EE+3YCdAaPaxE2ZRY/06pZUdp0tY4IgpuI1SZQ= +github.com/fsnotify/fsnotify v1.7.0 h1:8JEhPFa5W2WU7YfeZzPNqzMP6Lwt7L2715Ggo0nosvA= +github.com/fsnotify/fsnotify v1.7.0/go.mod h1:40Bi/Hjc2AVfZrqy+aj+yEI+/bRxZnMJyTJwOpGvigM= github.com/gammazero/deque v0.2.1 h1:qSdsbG6pgp6nL7A0+K/B7s12mcCY/5l5SIUpMOl+dC0= github.com/gammazero/deque v0.2.1/go.mod h1:LFroj8x4cMYCukHJDbxFCkT+r9AndaJnFMuZDV34tuU= github.com/gammazero/workerpool v1.1.3 h1:WixN4xzukFoN0XSeXF6puqEqFTl2mECI9S6W44HWy9Q= @@ -58,8 +60,8 @@ github.com/golang/protobuf v1.4.0/go.mod h1:jodUvKwWbYaEsadDk5Fwe5c77LiNKVO9IDvq github.com/golang/protobuf v1.4.2/go.mod h1:oDoupMAO8OvCJWAcko0GGGIgR6R6ocIYbsSw735rRwI= github.com/golang/protobuf v1.5.0/go.mod h1:FsONVRAS9T7sI+LIUmWTfcYkHO4aIWwzhcaSAoJOfIk= github.com/golang/protobuf v1.5.2/go.mod h1:XVQd3VNwM+JqD3oG2Ue2ip4fOMUkwXdXDdiuN0vRsmY= -github.com/golang/protobuf v1.5.3 h1:KhyjKVUg7Usr/dYsdSqoFveMYd5ko72D+zANwlG1mmg= -github.com/golang/protobuf v1.5.3/go.mod h1:XVQd3VNwM+JqD3oG2Ue2ip4fOMUkwXdXDdiuN0vRsmY= +github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek= +github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps= github.com/google/go-cmp v0.2.0/go.mod h1:oXzfMopK8JAjlY9xF4vHSVASa0yLyX7SntLO5aqRK0M= github.com/google/go-cmp v0.3.0/go.mod h1:8QqcDgzrUqlUb/G2PQTWiueGozuR1884gddMywk6iLU= github.com/google/go-cmp v0.3.1/go.mod h1:8QqcDgzrUqlUb/G2PQTWiueGozuR1884gddMywk6iLU= @@ -111,8 +113,8 @@ github.com/jsimonetti/rtnetlink v0.0.0-20211022192332-93da33804786 h1:N527AHMa79 github.com/jsimonetti/rtnetlink v0.0.0-20211022192332-93da33804786/go.mod h1:v4hqbTdfQngbVSZJVWUhGE/lbTFf9jb+ygmNUDQMuOs= github.com/jxskiss/base62 v1.1.0 h1:A5zbF8v8WXx2xixnAKD2w+abC+sIzYJX+nxmhA6HWFw= github.com/jxskiss/base62 v1.1.0/go.mod h1:HhWAlUXvxKThfOlZbcuFzsqwtF5TcqS9ru3y5GfjWAc= -github.com/klauspost/compress v1.17.6 h1:60eq2E/jlfwQXtvZEeBUYADs+BwKBWURIY+Gj2eRGjI= -github.com/klauspost/compress v1.17.6/go.mod h1:/dCuZOvVtNoHsyb+cuJD3itjs3NbnF6KH9zAO4BDxPM= +github.com/klauspost/compress v1.17.7 h1:ehO88t2UGzQK66LMdE8tibEd1ErmzZjNEqWkjLAKQQg= +github.com/klauspost/compress v1.17.7/go.mod h1:Di0epgTjJY877eYKx5yC51cX2A2Vl2ibi7bDH9ttBbw= github.com/klauspost/cpuid/v2 v2.2.6 h1:ndNyv040zDGIDh8thGkXYjnFtiN02M1PVVF+JE/48xc= github.com/klauspost/cpuid/v2 v2.2.6/go.mod h1:Lcz8mBdAVJIBVzewtcLocK12l3Y+JytZYpaMropDUws= github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo= @@ -130,10 +132,10 @@ 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-20240302142739-1c3dd691a1b8 h1:xawydPEACNO5Ncs2LgioTjWghXQ0eUN1q1RnVUUyVnI= github.com/livekit/mediatransportutil v0.0.0-20240302142739-1c3dd691a1b8/go.mod h1:jwKUCmObuiEDH0iiuJHaGMXwRs3RjrB4G6qqgkr/5oE= -github.com/livekit/protocol v1.11.1-0.20240311174744-00c977ffbb49 h1:eI3blFvnnFvAaT5oTbVZBTQCufv9OEplW8ejYWDR6tY= -github.com/livekit/protocol v1.11.1-0.20240311174744-00c977ffbb49/go.mod h1:x+QyergF26N374J0HxxDw6AiMEtf8wNtJS7yJrWGufA= -github.com/livekit/psrpc v0.5.3-0.20240228172457-3724cb4adbc4 h1:253WtQ2VGVHzIIzW9MUZj7vUDDILESU3zsEbiRdxYF0= -github.com/livekit/psrpc v0.5.3-0.20240228172457-3724cb4adbc4/go.mod h1:CQUBSPfYYAaevg1TNCc6/aYsa8DJH4jSRFdCeSZk5u0= +github.com/livekit/protocol v1.11.1-0.20240313083005-cf54792d0626 h1:FMN/mD/7e06gDi51sT0BdaTGIWIcdK140YTQNJI4HrY= +github.com/livekit/protocol v1.11.1-0.20240313083005-cf54792d0626/go.mod h1:znZpBU024XvKIpC9jkFzfrJMDuucDb4B1ON9lbUJMGA= +github.com/livekit/psrpc v0.5.3-0.20240312110212-61ab09477c30 h1:3GEU6vP+KLTTOEqsFKW+PgIUp+i+s0jaUqogQc/hb7M= +github.com/livekit/psrpc v0.5.3-0.20240312110212-61ab09477c30/go.mod h1:CQUBSPfYYAaevg1TNCc6/aYsa8DJH4jSRFdCeSZk5u0= github.com/mackerelio/go-osstat v0.2.4 h1:qxGbdPkFo65PXOb/F/nhDKpF2nGmGaCFDLXoZjJTtUs= github.com/mackerelio/go-osstat v0.2.4/go.mod h1:Zy+qzGdZs3A9cuIqmgbJvwbmLQH9dJvtio5ZjJTbdlQ= github.com/magefile/mage v1.15.0 h1:BvGheCMAsG3bWUDbZ8AyXXpCNwU9u5CB6sM+HNb9HYg= @@ -163,8 +165,8 @@ github.com/mdlayher/socket v0.4.0 h1:280wsy40IC9M9q1uPGcLBwXpcTQDtoGwVt+BNoITxIw github.com/mdlayher/socket v0.4.0/go.mod h1:xxFqz5GRCUN3UEOm9CZqEJsAbe1C8OwSK46NlmWuVoc= github.com/mitchellh/go-homedir v1.1.0 h1:lukF9ziXFxDFPkA1vsr5zpc1XuPDn/wFntq5mG+4E0Y= github.com/mitchellh/go-homedir v1.1.0/go.mod h1:SfyaCUpYCn1Vlf4IUYiD9fPX4A5wJrkLzIz1N1q0pr0= -github.com/nats-io/nats.go v1.32.0 h1:Bx9BZS+aXYlxW08k8Gd3yR2s73pV5XSoAQUyp1Kwvp0= -github.com/nats-io/nats.go v1.32.0/go.mod h1:Ubdu4Nh9exXdSz0RVWRFBbRfrbSxOYd26oF0wkWclB8= +github.com/nats-io/nats.go v1.33.1 h1:8TxLZZ/seeEfR97qV0/Bl939tpDnt2Z2fK3HkPypj70= +github.com/nats-io/nats.go v1.33.1/go.mod h1:Ubdu4Nh9exXdSz0RVWRFBbRfrbSxOYd26oF0wkWclB8= github.com/nats-io/nkeys v0.4.7 h1:RwNJbbIdYCoClSDNY7QVKZlyb/wfT6ugvFCiKy6vDvI= github.com/nats-io/nkeys v0.4.7/go.mod h1:kqXRgRDPlGy7nGaEDMuYzmiJCIAAWDK0IMBtDmGD0nc= github.com/nats-io/nuid v1.0.1 h1:5iA8DT8V7q8WK2EScv2padNa/rTESc1KdnPw4TC2paw= @@ -304,8 +306,9 @@ golang.org/x/crypto v0.11.0/go.mod h1:xgJhtzW8F9jGdVFWZESrid1U1bjeNy4zgy5cRr/CIi golang.org/x/crypto v0.12.0/go.mod h1:NF0Gs7EO5K4qLn+Ylc+fih8BSTeIjAP05siRnAh98yw= golang.org/x/crypto v0.13.0/go.mod h1:y6Z2r+Rw4iayiXXAIxJIDAJ1zMW4yaTpebo8fPOliYc= golang.org/x/crypto v0.18.0/go.mod h1:R0j02AL6hcrfOiy9T4ZYp/rcWeMxM3L6QYxlOuEG1mg= -golang.org/x/crypto v0.19.0 h1:ENy+Az/9Y1vSrlrvBSyna3PITt4tiZLf7sgCjZBX7Wo= golang.org/x/crypto v0.19.0/go.mod h1:Iy9bg/ha4yyC70EfRS8jz+B6ybOBKMaSxLj6P6oBDfU= +golang.org/x/crypto v0.21.0 h1:X31++rzVUdKhX5sWmSOFZxx8UW/ldWx55cbf08iNAMA= +golang.org/x/crypto v0.21.0/go.mod h1:0BP7YvVV9gBbVKyeTG0Gyn+gZm94bibOW5BjDEYAOMs= golang.org/x/exp v0.0.0-20240222234643-814bf88cf225 h1:LfspQV/FYTatPTr/3HzIcmiUFH7PGP+OQ6mgDYo3yuQ= golang.org/x/exp v0.0.0-20240222234643-814bf88cf225/go.mod h1:CxmFvTBINI24O/j8iY7H1xHzx2i4OsyguNBmN/uPtqc= golang.org/x/mod v0.3.0/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA= @@ -396,8 +399,9 @@ golang.org/x/sys v0.10.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.11.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.12.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.16.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= -golang.org/x/sys v0.17.0 h1:25cE3gD+tdBA7lp7QfhuV+rJiE9YXTcS3VG1SqssI/Y= golang.org/x/sys v0.17.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= +golang.org/x/sys v0.18.0 h1:DBdB3niSjOA/O0blCZBqDefyWNYveAYMNF1Wum0DYQ4= +golang.org/x/sys v0.18.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= golang.org/x/term v0.0.0-20210927222741-03fcf44c2211/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8= golang.org/x/term v0.1.0/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8= @@ -434,10 +438,10 @@ golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8T golang.org/x/xerrors v0.0.0-20191011141410-1b5146add898/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20200804184101-5ec99f83aff1/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= -google.golang.org/genproto/googleapis/rpc v0.0.0-20240221002015-b0ce06bbee7c h1:NUsgEN92SQQqzfA+YtqYNqYmB3DMMYLlIwUZAQFVFbo= -google.golang.org/genproto/googleapis/rpc v0.0.0-20240221002015-b0ce06bbee7c/go.mod h1:H4O17MA/PE9BsGx3w+a+W2VOLLD1Qf7oJneAoU6WktY= -google.golang.org/grpc v1.62.0 h1:HQKZ/fa1bXkX1oFOvSjmZEUL8wLSaZTjCcLAlmZRtdk= -google.golang.org/grpc v1.62.0/go.mod h1:IWTG0VlJLCh1SkC58F7np9ka9mx/WNkjl4PGJaiq+QE= +google.golang.org/genproto/googleapis/rpc v0.0.0-20240311173647-c811ad7063a7 h1:8EeVk1VKMD+GD/neyEHGmz7pFblqPjHoi+PGQIlLx2s= +google.golang.org/genproto/googleapis/rpc v0.0.0-20240311173647-c811ad7063a7/go.mod h1:WtryC6hu0hhx87FDGxWCDptyssuo68sk10vYjF+T9fY= +google.golang.org/grpc v1.62.1 h1:B4n+nfKzOICUXMgyrNd19h/I9oH0L1pizfk1d4zSgTk= +google.golang.org/grpc v1.62.1/go.mod h1:IWTG0VlJLCh1SkC58F7np9ka9mx/WNkjl4PGJaiq+QE= google.golang.org/protobuf v0.0.0-20200109180630-ec00e32a8dfd/go.mod h1:DFci5gLYBciE7Vtevhsrf46CRTquxDuWsQurQQe4oz8= google.golang.org/protobuf v0.0.0-20200221191635-4d8936d0db64/go.mod h1:kwYJMbMJ01Woi6D6+Kah6886xMZcty6N08ah7+eCXa0= google.golang.org/protobuf v0.0.0-20200228230310-ab0ca4ff8a60/go.mod h1:cfTl7dwQJ+fmap5saPgwCLgHXTUD7jkjRqWcaiX5VyM= @@ -446,8 +450,8 @@ google.golang.org/protobuf v1.21.0/go.mod h1:47Nbq4nVaFHyn7ilMalzfO3qCViNmqZ2kzi google.golang.org/protobuf v1.23.0/go.mod h1:EGpADcykh3NcUnDUJcl1+ZksZNG86OlYog2l/sGQquU= google.golang.org/protobuf v1.26.0-rc.1/go.mod h1:jlhhOSvTdKEhbULTjvd4ARK9grFBp09yW+WbY/TyQbw= google.golang.org/protobuf v1.26.0/go.mod h1:9q0QmTI4eRPtz6boOQmLYwt+qCgq0jsYwAQnmE0givc= -google.golang.org/protobuf v1.32.0 h1:pPC6BG5ex8PDFnkbrGU3EixyhKcQ2aDuBS36lqK/C7I= -google.golang.org/protobuf v1.32.0/go.mod h1:c6P6GXX6sHbq/GpV6MGZEdwhWPcYBgnhAHhKbcUYpos= +google.golang.org/protobuf v1.33.0 h1:uNO2rsAINq/JlFpSdYEKIZ0uKD/R9cpdv0T+yoGwGmI= +google.golang.org/protobuf v1.33.0/go.mod h1:c6P6GXX6sHbq/GpV6MGZEdwhWPcYBgnhAHhKbcUYpos= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v1.0.0-20180628173108-788fd7840127/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v1.0.0-20190902080502-41f04d3bba15/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= diff --git a/pkg/rtc/wrappedreceiver.go b/pkg/rtc/wrappedreceiver.go index e088c9eb3..442a8df83 100644 --- a/pkg/rtc/wrappedreceiver.go +++ b/pkg/rtc/wrappedreceiver.go @@ -324,11 +324,11 @@ func (d *DummyReceiver) GetReferenceLayerRTPTimestamp(ts uint32, layer int32, re return 0, errors.New("receiver not available") } -func (d *DummyReceiver) GetRTCPSenderReportData(layer int32) (*buffer.RTCPSenderReportData, *buffer.RTCPSenderReportData) { +func (d *DummyReceiver) GetRTCPSenderReportData(layer int32) *buffer.RTCPSenderReportData { if r, ok := d.receiver.Load().(sfu.TrackReceiver); ok { return r.GetRTCPSenderReportData(layer) } - return nil, nil + return nil } func (d *DummyReceiver) GetTrackStats() *livekit.RTPStats { diff --git a/pkg/sfu/buffer/buffer.go b/pkg/sfu/buffer/buffer.go index b2e6b75a7..123472ad6 100644 --- a/pkg/sfu/buffer/buffer.go +++ b/pkg/sfu/buffer/buffer.go @@ -840,7 +840,7 @@ func (b *Buffer) SetSenderReportData(rtpTime uint32, ntpTime uint64) { } } -func (b *Buffer) GetSenderReportData() (*RTCPSenderReportData, *RTCPSenderReportData) { +func (b *Buffer) GetSenderReportData() *RTCPSenderReportData { b.RLock() defer b.RUnlock() @@ -848,7 +848,7 @@ func (b *Buffer) GetSenderReportData() (*RTCPSenderReportData, *RTCPSenderReport return b.rtpStats.GetRtcpSenderReportData() } - return nil, nil + return nil } func (b *Buffer) SetLastFractionLostReport(lost uint8) { diff --git a/pkg/sfu/buffer/rtpstats_base.go b/pkg/sfu/buffer/rtpstats_base.go index 7af312c11..11c9517c6 100644 --- a/pkg/sfu/buffer/rtpstats_base.go +++ b/pkg/sfu/buffer/rtpstats_base.go @@ -484,7 +484,7 @@ func (r *rtpStatsBase) GetRtt() uint32 { return r.rtt } -func (r *rtpStatsBase) maybeAdjustFirstPacketTime(ts uint32, startTS uint32) { +func (r *rtpStatsBase) maybeAdjustFirstPacketTime(srData *RTCPSenderReportData, tsOffset uint64, extStartTS uint64) { if time.Since(r.startTime) > cFirstPacketTimeAdjustWindow { return } @@ -495,7 +495,9 @@ func (r *rtpStatsBase) maybeAdjustFirstPacketTime(ts uint32, startTS uint32) { // abnormal delay (maybe due to pacing or maybe due to queuing // in some network element along the way), push back first time // to an earlier instance. - samplesDiff := int32(ts - startTS) + timeSinceReceive := time.Since(srData.At) + extNowTS := srData.RTPTimestampExt - tsOffset + uint64(timeSinceReceive.Nanoseconds()*int64(r.params.ClockRate)/1e9) + samplesDiff := int64(extNowTS - extStartTS) if samplesDiff < 0 { // out-of-order, skip return @@ -505,28 +507,24 @@ func (r *rtpStatsBase) maybeAdjustFirstPacketTime(ts uint32, startTS uint32) { timeSinceFirst := time.Since(r.firstTime) now := r.firstTime.Add(timeSinceFirst) firstTime := now.Add(-samplesDuration) + + getFields := func() []interface{} { + return []interface{}{ + "startTime", r.startTime.String(), + "nowTime", now.String(), + "before", r.firstTime.String(), + "after", firstTime.String(), + "adjustment", r.firstTime.Sub(firstTime).String(), + "extNowTS", extNowTS, + "extStartTS", extStartTS, + } + } + if firstTime.Before(r.firstTime) { if r.firstTime.Sub(firstTime) > cFirstPacketTimeAdjustThreshold { - r.logger.Infow("adjusting first packet time, too big, ignoring", - "startTime", r.startTime.String(), - "nowTime", now.String(), - "before", r.firstTime.String(), - "after", firstTime.String(), - "adjustment", r.firstTime.Sub(firstTime).String(), - "nowTS", ts, - "startTS", startTS, - ) + r.logger.Infow("adjusting first packet time, too big, ignoring", getFields()...) } else { - r.logger.Debugw( - "adjusting first packet time", - "startTime", r.startTime.String(), - "nowTime", now.String(), - "before", r.firstTime.String(), - "after", firstTime.String(), - "adjustment", r.firstTime.Sub(firstTime).String(), - "nowTS", ts, - "startTS", startTS, - ) + r.logger.Debugw("adjusting first packet time", getFields()...) r.firstTime = firstTime } } @@ -678,7 +676,7 @@ func (r *rtpStatsBase) toString( str += ", rtt(ms):" str += fmt.Sprintf("%d|%d", p.RttCurrent, p.RttMax) - str += fmt.Sprintf(", pd: %s, rd: %s", RTPDriftToString(p.PacketDrift), RTPDriftToString(p.ReportDrift)) + str += fmt.Sprintf(", pd: %s, nrd: %s, rrd: %s", RTPDriftToString(p.PacketDrift), RTPDriftToString(p.ReportDrift), RTPDriftToString(p.RebasedReportDrift)) return str } @@ -719,7 +717,7 @@ func (r *rtpStatsBase) toProto( jitterTime := jitter / float64(r.params.ClockRate) * 1e6 maxJitterTime := maxJitter / float64(r.params.ClockRate) * 1e6 - packetDrift, reportDrift := r.getDrift(extStartTS, extHighestTS) + packetDrift, ntpReportDrift, rebasedReportDrift := r.getDrift(extStartTS, extHighestTS) p := &livekit.RTPStats{ StartTime: timestamppb.New(r.startTime), @@ -763,7 +761,8 @@ func (r *rtpStatsBase) toProto( RttCurrent: r.rtt, RttMax: r.maxRtt, PacketDrift: packetDrift, - ReportDrift: reportDrift, + ReportDrift: ntpReportDrift, + RebasedReportDrift: rebasedReportDrift, } gapsPresent := false @@ -845,7 +844,7 @@ func (r *rtpStatsBase) getAndResetSnapshot(snapshotID uint32, extStartSN uint64, return &then, &now } -func (r *rtpStatsBase) getDrift(extStartTS, extHighestTS uint64) (packetDrift *livekit.RTPDrift, reportDrift *livekit.RTPDrift) { +func (r *rtpStatsBase) getDrift(extStartTS, extHighestTS uint64) (packetDrift *livekit.RTPDrift, ntpReportDrift *livekit.RTPDrift, rebasedReportDrift *livekit.RTPDrift) { if !r.firstTime.IsZero() { elapsed := r.highestTime.Sub(r.firstTime) rtpClockTicks := extHighestTS - extStartTS @@ -866,11 +865,12 @@ func (r *rtpStatsBase) getDrift(extStartTS, extHighestTS uint64) (packetDrift *l } if r.srFirst != nil && r.srNewest != nil && r.srFirst.RTPTimestamp != r.srNewest.RTPTimestamp { - elapsed := r.srNewest.NTPTimestamp.Time().Sub(r.srFirst.NTPTimestamp.Time()) rtpClockTicks := r.srNewest.RTPTimestampExt - r.srFirst.RTPTimestampExt + + elapsed := r.srNewest.NTPTimestamp.Time().Sub(r.srFirst.NTPTimestamp.Time()) driftSamples := int64(rtpClockTicks - uint64(elapsed.Nanoseconds()*int64(r.params.ClockRate)/1e9)) if elapsed.Seconds() > 0.0 { - reportDrift = &livekit.RTPDrift{ + ntpReportDrift = &livekit.RTPDrift{ StartTime: timestamppb.New(r.srFirst.NTPTimestamp.Time()), EndTime: timestamppb.New(r.srNewest.NTPTimestamp.Time()), Duration: elapsed.Seconds(), @@ -882,6 +882,22 @@ func (r *rtpStatsBase) getDrift(extStartTS, extHighestTS uint64) (packetDrift *l ClockRate: float64(rtpClockTicks) / elapsed.Seconds(), } } + + elapsed = r.srNewest.At.Sub(r.srFirst.At) + driftSamples = int64(rtpClockTicks - uint64(elapsed.Nanoseconds()*int64(r.params.ClockRate)/1e9)) + if elapsed.Seconds() > 0.0 { + rebasedReportDrift = &livekit.RTPDrift{ + StartTime: timestamppb.New(r.srFirst.At), + EndTime: timestamppb.New(r.srNewest.At), + Duration: elapsed.Seconds(), + StartTimestamp: r.srFirst.RTPTimestampExt, + EndTimestamp: r.srNewest.RTPTimestampExt, + RtpClockTicks: rtpClockTicks, + DriftSamples: driftSamples, + DriftMs: (float64(driftSamples) * 1000) / float64(r.params.ClockRate), + ClockRate: float64(rtpClockTicks) / elapsed.Seconds(), + } + } } return } diff --git a/pkg/sfu/buffer/rtpstats_receiver.go b/pkg/sfu/buffer/rtpstats_receiver.go index 2a789c443..c54dd53c8 100644 --- a/pkg/sfu/buffer/rtpstats_receiver.go +++ b/pkg/sfu/buffer/rtpstats_receiver.go @@ -28,6 +28,32 @@ import ( const ( cHistorySize = 4096 + + // RTCP Sender Reports are re-based to SFU time base so that all subscriber side + // can have the same time base (i. e. SFU time base). To convert publisher side + // RTCP Sender Reports to SFU timebase, a propagation delay is maintained. + // propagation_delay = time_of_report_reception - ntp_timestamp_in_report + // + // Propagation delay is adapted continuously. If it falls, adapt quickly to the + // lower value as that could be the real propagation delay. If it rises, adapt slowly + // as it might be a temporary change or slow drift. See below for handling of high deltas + // which could be a result of a path change. + cPropagationDelayFallFactor = float64(0.95) + cPropagationDelayRiseFactor = float64(0.05) + + // do not adapt to small OR large (outlier) changes + cPropagationDelayDeltaThresholdMin = 5 * time.Millisecond + cPropagationDelayDeltaThresholdMaxFactor = 2 + + // To account for path changes mid-stream, if the delta of the propagation delay is consistently higher, reset. + // Reset at whichever of the below happens later. + // + // A smoothed version of delta of propagation delay is maintained and delta propagation delay exceeding + // a factor of the smoothed version is considered a sharp increase. That will trigger the start of the + // path change condition and if it persists, propagation delay will be reset. + cPropagationDelayDeltaAdaptationFactor = float64(0.1) + cPropagationDelayDeltaHighResetNumReports = 3 + cPropagationDelayDeltaHighResetWait = 10 * time.Second ) type RTPFlowState struct { @@ -53,6 +79,11 @@ type RTPStatsReceiver struct { history *protoutils.Bitmap[uint64] + propagationDelay time.Duration + smoothedDeltaPropagationDelay time.Duration + propagationDelayDeltaHighCount int + propagationDelayDeltaHighStartTime time.Time + clockSkewCount int outOfOrderSsenderReportCount int } @@ -294,8 +325,6 @@ func (r *RTPStatsReceiver) SetRtcpSenderReportData(srData *RTCPSenderReportData) return } - r.maybeAdjustFirstPacketTime(srDataCopy.RTPTimestamp, r.timestamp.GetStart()) - if r.srNewest != nil { timeSinceLast := srData.NTPTimestamp.Time().Sub(r.srNewest.NTPTimestamp.Time()).Seconds() rtpDiffSinceLast := srDataCopy.RTPTimestampExt - r.srNewest.RTPTimestampExt @@ -326,26 +355,88 @@ func (r *RTPStatsReceiver) SetRtcpSenderReportData(srData *RTCPSenderReportData) } } - r.srNewest = &srDataCopy + var propagationDelay time.Duration + var deltaPropagationDelay time.Duration + getPropagationFields := func() []interface{} { + return []interface{}{ + "propagationDelay", r.propagationDelay.String(), + "receivedPropagationDelay", propagationDelay.String(), + "smoothedDeltaPropagationDelay", r.smoothedDeltaPropagationDelay.String(), + "receivedDeltaPropagationDelay", deltaPropagationDelay.String(), + "deltaHighCount", r.propagationDelayDeltaHighCount, + "sinceDeltaHighStart", time.Since(r.propagationDelayDeltaHighStartTime).String(), + "first", r.srFirst, + "last", r.srNewest, + "current", &srDataCopy, + } + } + initPropagationDelay := func(pd time.Duration) { + r.propagationDelay = pd + r.smoothedDeltaPropagationDelay = 0 + r.propagationDelayDeltaHighCount = 0 + r.propagationDelayDeltaHighStartTime = time.Time{} + } + + ntpTime := srDataCopy.NTPTimestamp.Time() + propagationDelay = srDataCopy.At.Sub(ntpTime) if r.srFirst == nil { r.srFirst = &srDataCopy + initPropagationDelay(propagationDelay) + r.logger.Debugw("initializing propagation delay", getPropagationFields()...) + } else { + deltaPropagationDelay = propagationDelay - r.propagationDelay + r.logger.Debugw("RAJA pd", getPropagationFields()...) // REMOVE + if r.smoothedDeltaPropagationDelay != 0 && deltaPropagationDelay > 0 && deltaPropagationDelay > r.smoothedDeltaPropagationDelay*time.Duration(cPropagationDelayDeltaThresholdMaxFactor) { + r.logger.Debugw("sharp increase in propagation delay, skipping", getPropagationFields()...) // TODO-REMOVE + r.propagationDelayDeltaHighCount++ + if r.propagationDelayDeltaHighStartTime.IsZero() { + r.propagationDelayDeltaHighStartTime = time.Now() + } + + if r.propagationDelayDeltaHighCount >= cPropagationDelayDeltaHighResetNumReports && time.Since(r.propagationDelayDeltaHighStartTime) >= cPropagationDelayDeltaHighResetWait { + r.logger.Infow("re-initializing propagation delay", append(getPropagationFields(), "newPropagationDelay", propagationDelay)...) + initPropagationDelay(propagationDelay) + } + } else { + if r.smoothedDeltaPropagationDelay == 0 { + r.smoothedDeltaPropagationDelay = deltaPropagationDelay + } else { + r.smoothedDeltaPropagationDelay += time.Duration(cPropagationDelayDeltaAdaptationFactor * float64(deltaPropagationDelay-r.smoothedDeltaPropagationDelay)) + } + r.propagationDelayDeltaHighCount = 0 + r.propagationDelayDeltaHighStartTime = time.Time{} + + if deltaPropagationDelay.Abs() > cPropagationDelayDeltaThresholdMin { + factor := cPropagationDelayFallFactor + if propagationDelay > r.propagationDelay { + factor = cPropagationDelayRiseFactor + } + fields := append( + getPropagationFields(), + "adjustedPropagationDelay", r.propagationDelay+time.Duration(factor*float64(propagationDelay-r.propagationDelay)), + ) // TODO-REMOVE + r.logger.Debugw("adapting propagation delay", fields...) // TODO-REMOVE + r.propagationDelay += time.Duration(factor * float64(propagationDelay-r.propagationDelay)) + } + } } + // adjust receive time to estimated propagation delay + srDataCopy.At = ntpTime.Add(r.propagationDelay) + r.srNewest = &srDataCopy + + r.maybeAdjustFirstPacketTime(r.srNewest, 0, r.timestamp.GetExtendedStart()) } -func (r *RTPStatsReceiver) GetRtcpSenderReportData() (srFirst *RTCPSenderReportData, srNewest *RTCPSenderReportData) { +func (r *RTPStatsReceiver) GetRtcpSenderReportData() *RTCPSenderReportData { r.lock.RLock() defer r.lock.RUnlock() - if r.srFirst != nil { - srFirstCopy := *r.srFirst - srFirst = &srFirstCopy + if r.srNewest == nil { + return nil } - if r.srNewest != nil { - srNewestCopy := *r.srNewest - srNewest = &srNewestCopy - } - return + srNewestCopy := *r.srNewest + return &srNewestCopy } func (r *RTPStatsReceiver) GetRtcpReceptionReport(ssrc uint32, proxyFracLost uint8, snapshotID uint32) *rtcp.ReceptionReport { diff --git a/pkg/sfu/buffer/rtpstats_sender.go b/pkg/sfu/buffer/rtpstats_sender.go index b78dc6b2a..4fe39c0c5 100644 --- a/pkg/sfu/buffer/rtpstats_sender.go +++ b/pkg/sfu/buffer/rtpstats_sender.go @@ -161,9 +161,6 @@ type RTPStatsSender struct { clockSkewCount int metadataCacheOverflowCount int - - srFeedFirst *RTCPSenderReportData - srFeedNewest *RTCPSenderReportData } func NewRTPStatsSender(params RTPStatsParams) *RTPStatsSender { @@ -202,15 +199,6 @@ func (r *RTPStatsSender) Seed(from *RTPStatsSender) { r.nextSenderSnapshotID = from.nextSenderSnapshotID r.senderSnapshots = make([]senderSnapshot, cap(from.senderSnapshots)) copy(r.senderSnapshots, from.senderSnapshots) - - if from.srFeedFirst != nil { - srFeedFirst := *from.srFeedFirst - r.srFeedFirst = &srFeedFirst - } - if from.srFeedNewest != nil { - srFeedNewest := *from.srFeedNewest - r.srFeedNewest = &srFeedNewest - } } func (r *RTPStatsSender) NewSnapshotId() uint32 { @@ -611,25 +599,15 @@ func (r *RTPStatsSender) LastReceiverReportTime() time.Time { return r.lastRRTime } -func (r *RTPStatsSender) MaybeAdjustFirstPacketTime(srFirst *RTCPSenderReportData, srNewest *RTCPSenderReportData, ts uint32) { +func (r *RTPStatsSender) MaybeAdjustFirstPacketTime(publisherSRData *RTCPSenderReportData, tsOffset uint64) { r.lock.Lock() defer r.lock.Unlock() - if !r.initialized { + if !r.initialized || publisherSRData == nil { return } - if srFirst != nil { - srFirstCopy := *srFirst - r.srFeedFirst = &srFirstCopy - } - - if srNewest != nil { - srNewestCopy := *srNewest - r.srFeedNewest = &srNewestCopy - } - - r.maybeAdjustFirstPacketTime(ts, uint32(r.extStartTS)) + r.maybeAdjustFirstPacketTime(publisherSRData, tsOffset, r.extStartTS) } func (r *RTPStatsSender) GetExpectedRTPTimestamp(at time.Time) (expectedTSExt uint64, err error) { @@ -647,22 +625,18 @@ func (r *RTPStatsSender) GetExpectedRTPTimestamp(at time.Time) (expectedTSExt ui return } -func (r *RTPStatsSender) GetRtcpSenderReport(ssrc uint32) *rtcp.SenderReport { +func (r *RTPStatsSender) GetRtcpSenderReport(ssrc uint32, publisherSRData *RTCPSenderReportData, tsOffset uint64) *rtcp.SenderReport { r.lock.Lock() defer r.lock.Unlock() - if !r.initialized { + if !r.initialized || publisherSRData == nil { return nil } - // construct current time based on monotonic clock - timeSinceFirst := time.Since(r.firstTime) - if timeSinceFirst < cSenderReportInitialWait { - return nil - } - now := r.firstTime.Add(timeSinceFirst) + timeSincePublisherSR := time.Since(publisherSRData.At) + now := publisherSRData.At.Add(timeSincePublisherSR) nowNTP := mediatransportutil.ToNtpTime(now) - nowRTPExt := r.extStartTS + uint64(timeSinceFirst.Nanoseconds()*int64(r.params.ClockRate)/1e9) + nowRTPExt := publisherSRData.RTPTimestampExt - tsOffset + uint64(timeSincePublisherSR.Nanoseconds()*int64(r.params.ClockRate)/1e9) srData := &RTCPSenderReportData{ NTPTimestamp: nowNTP, @@ -670,32 +644,39 @@ func (r *RTPStatsSender) GetRtcpSenderReport(ssrc uint32) *rtcp.SenderReport { RTPTimestampExt: nowRTPExt, At: now, } + + getFields := func() []interface{} { + return []interface{}{ + "first", r.srFirst, + "last", r.srNewest, + "curr", srData, + "feed", publisherSRData, + "tsOffset", tsOffset, + "timeNow", time.Now().String(), + "extStartTS", r.extStartTS, + "extHighestTS", r.extHighestTS, + "highestTime", r.highestTime.String(), + "timeSinceHighest", now.Sub(r.highestTime).String(), + "firstTime", r.firstTime.String(), + "timeSinceFirst", now.Sub(r.firstTime).String(), + "timeSincePublisherSR", timeSincePublisherSR.String(), + "nowRTPExt", nowRTPExt, + } + } if r.srNewest != nil && nowRTPExt >= r.srNewest.RTPTimestampExt { timeSinceLastReport := nowNTP.Time().Sub(r.srNewest.NTPTimestamp.Time()) rtpDiffSinceLastReport := nowRTPExt - r.srNewest.RTPTimestampExt windowClockRate := float64(rtpDiffSinceLastReport) / timeSinceLastReport.Seconds() if timeSinceLastReport.Seconds() > 0.2 && math.Abs(float64(r.params.ClockRate)-windowClockRate) > 0.2*float64(r.params.ClockRate) { if r.clockSkewCount%10 == 0 { - r.logger.Infow( - "sending sender report, clock skew", - "first", r.srFirst, - "last", r.srNewest, - "curr", srData, - "firstFeed", r.srFeedFirst, - "lastFeed", r.srFeedNewest, - "timeNow", time.Now().String(), - "extStartTS", r.extStartTS, - "extHighestTS", r.extHighestTS, - "highestTime", r.highestTime.String(), - "timeSinceHighest", now.Sub(r.highestTime).String(), - "firstTime", r.firstTime.String(), - "timeSinceFirst", timeSinceFirst.String(), - "nowRTPExt", nowRTPExt, + fields := append( + getFields(), "timeSinceLastReport", timeSinceLastReport.String(), "rtpDiffSinceLastReport", rtpDiffSinceLastReport, "windowClockRate", windowClockRate, "count", r.clockSkewCount, ) + r.logger.Infow("sending sender report, clock skew", fields...) } r.clockSkewCount++ } @@ -704,22 +685,7 @@ func (r *RTPStatsSender) GetRtcpSenderReport(ssrc uint32) *rtcp.SenderReport { if r.srNewest != nil && nowRTPExt < r.srNewest.RTPTimestampExt { // If report being generated is behind the last report, skip it. // Should not happen. - r.logger.Infow( - "sending sender report, out-of-order, skipping", - "first", r.srFirst, - "last", r.srNewest, - "curr", srData, - "firstFeed", r.srFeedFirst, - "lastFeed", r.srFeedNewest, - "timeNow", time.Now().String(), - "extStartTS", r.extStartTS, - "extHighestTS", r.extHighestTS, - "highestTime", r.highestTime.String(), - "timeSinceHighest", now.Sub(r.highestTime).String(), - "firstTime", r.firstTime.String(), - "timeSinceFirst", timeSinceFirst.String(), - "nowRTPExt", nowRTPExt, - ) + r.logger.Infow("sending sender report, out-of-order, skipping", getFields()...) return nil } diff --git a/pkg/sfu/downtrack.go b/pkg/sfu/downtrack.go index 1bc07b8ff..41e6eb67f 100644 --- a/pkg/sfu/downtrack.go +++ b/pkg/sfu/downtrack.go @@ -59,8 +59,7 @@ type TrackSender interface { payloadType webrtc.PayloadType, isSVC bool, layer int32, - srFirst *buffer.RTCPSenderReportData, - srNewest *buffer.RTCPSenderReportData, + publisherSRData *buffer.RTCPSenderReportData, ) error } @@ -1303,11 +1302,8 @@ func (d *DownTrack) CreateSenderReport() *rtcp.SenderReport { return nil } - clockLayer := d.forwarder.CurrentLayer().Spatial - if clockLayer == buffer.InvalidLayerSpatial { - clockLayer = d.forwarder.GetReferenceLayerSpatial() - } - return d.rtpStats.GetRtcpSenderReport(d.ssrc) + layer, tsOffset := d.forwarder.GetCurrentSpatialAndTSOffset() + return d.rtpStats.GetRtcpSenderReport(d.ssrc, d.params.Receiver.GetRTCPSenderReportData(layer), tsOffset) } func (d *DownTrack) writeBlankFrameRTP(duration float32, generation uint32) chan struct{} { @@ -1948,25 +1944,17 @@ func (d *DownTrack) HandleRTCPSenderReportData( _payloadType webrtc.PayloadType, isSVC bool, layer int32, - srFirst *buffer.RTCPSenderReportData, - srNewest *buffer.RTCPSenderReportData, + publisherSRData *buffer.RTCPSenderReportData, ) error { - if layer == d.forwarder.GetReferenceLayerSpatial() || (layer == 0 && isSVC) { - d.handleRTCPSenderReportData(srFirst, srNewest) + currentLayer, tsOffset := d.forwarder.GetCurrentSpatialAndTSOffset() + if layer == currentLayer || (layer == 0 && isSVC) { + d.handleRTCPSenderReportData(publisherSRData, tsOffset) } return nil } -func (d *DownTrack) handleRTCPSenderReportData(srFirst *buffer.RTCPSenderReportData, srNewest *buffer.RTCPSenderReportData) { - if srNewest == nil { - return - } - - d.rtpStats.MaybeAdjustFirstPacketTime( - srFirst, - srNewest, - srNewest.RTPTimestamp+uint32(d.forwarder.GetReferenceTimestampOffset()), - ) +func (d *DownTrack) handleRTCPSenderReportData(publisherSRData *buffer.RTCPSenderReportData, tsOffset uint64) { + d.rtpStats.MaybeAdjustFirstPacketTime(publisherSRData, tsOffset) } type sendPacketMetadata struct { @@ -2018,8 +2006,9 @@ func (d *DownTrack) sendingPacket(hdr *rtp.Header, payloadSize int, spmd *sendPa } if spmd.tp.isResuming { - // adjust first packet time on a resumption so that sender reports can lock in quicker - d.handleRTCPSenderReportData(d.params.Receiver.GetRTCPSenderReportData(d.forwarder.GetReferenceLayerSpatial())) + // adjust first packet time on a resumption so that subsequent switches get a more accurate expected time stamp + currentLayer, tsOffset := d.forwarder.GetCurrentSpatialAndTSOffset() + d.handleRTCPSenderReportData(d.params.Receiver.GetRTCPSenderReportData(currentLayer), tsOffset) if sal := d.getStreamAllocatorListener(); sal != nil { sal.OnResume(d) diff --git a/pkg/sfu/forwarder.go b/pkg/sfu/forwarder.go index 3e684e39d..e2bc0f1c9 100644 --- a/pkg/sfu/forwarder.go +++ b/pkg/sfu/forwarder.go @@ -556,18 +556,15 @@ func (f *Forwarder) GetMaxSubscribedSpatial() int32 { return layer } -func (f *Forwarder) GetReferenceLayerSpatial() int32 { +func (f *Forwarder) GetCurrentSpatialAndTSOffset() (int32, uint64) { f.lock.RLock() defer f.lock.RUnlock() - return f.referenceLayerSpatial -} + if f.kind == webrtc.RTPCodecTypeAudio { + return 0, f.rtpMunger.GetTSOffset() + } -func (f *Forwarder) GetReferenceTimestampOffset() uint64 { - f.lock.RLock() - defer f.lock.RUnlock() - - return f.refTSOffset + return f.vls.GetCurrent().Spatial, f.rtpMunger.GetTSOffset() } func (f *Forwarder) isDeficientLocked() bool { diff --git a/pkg/sfu/receiver.go b/pkg/sfu/receiver.go index eec709dba..3028e98e3 100644 --- a/pkg/sfu/receiver.go +++ b/pkg/sfu/receiver.go @@ -84,7 +84,7 @@ type TrackReceiver interface { GetTemporalLayerFpsForSpatial(layer int32) []float32 GetReferenceLayerRTPTimestamp(ts uint32, layer int32, referenceLayer int32) (uint32, error) - GetRTCPSenderReportData(layer int32) (*buffer.RTCPSenderReportData, *buffer.RTCPSenderReportData) + GetRTCPSenderReportData(layer int32) *buffer.RTCPSenderReportData GetTrackStats() *livekit.RTPStats } @@ -349,11 +349,11 @@ func (w *WebRTCReceiver) AddUpTrack(track *webrtc.TrackRemote, buff *buffer.Buff }) buff.OnRtcpFeedback(w.sendRTCP) buff.OnRtcpSenderReport(func() { - srFirst, srNewest := buff.GetSenderReportData() - w.streamTrackerManager.SetRTCPSenderReportData(layer, srFirst, srNewest) + srData := buff.GetSenderReportData() + w.streamTrackerManager.SetRTCPSenderReportData(layer, srData) w.downTrackSpreader.Broadcast(func(dt TrackSender) { - _ = dt.HandleRTCPSenderReportData(w.codec.PayloadType, w.isSVC, layer, srFirst, srNewest) + _ = dt.HandleRTCPSenderReportData(w.codec.PayloadType, w.isSVC, layer, srData) }) }) @@ -786,7 +786,7 @@ func (w *WebRTCReceiver) GetReferenceLayerRTPTimestamp(ts uint32, layer int32, r return w.streamTrackerManager.GetReferenceLayerRTPTimestamp(ts, layer, referenceLayer) } -func (w *WebRTCReceiver) GetRTCPSenderReportData(layer int32) (*buffer.RTCPSenderReportData, *buffer.RTCPSenderReportData) { +func (w *WebRTCReceiver) GetRTCPSenderReportData(layer int32) *buffer.RTCPSenderReportData { return w.streamTrackerManager.GetRTCPSenderReportData(layer) } diff --git a/pkg/sfu/rtpmunger.go b/pkg/sfu/rtpmunger.go index 9ff94ea1a..49c918037 100644 --- a/pkg/sfu/rtpmunger.go +++ b/pkg/sfu/rtpmunger.go @@ -123,6 +123,10 @@ func (r *RTPMunger) GetLast() RTPMungerState { } } +func (r *RTPMunger) GetTSOffset() uint64 { + return r.tsOffset +} + func (r *RTPMunger) SeedLast(state RTPMungerState) { r.extLastSN = state.ExtLastSN r.extSecondLastSN = state.ExtSecondLastSN diff --git a/pkg/sfu/streamtrackermanager.go b/pkg/sfu/streamtrackermanager.go index 2b80e7c25..b5f36dc97 100644 --- a/pkg/sfu/streamtrackermanager.go +++ b/pkg/sfu/streamtrackermanager.go @@ -50,12 +50,6 @@ type StreamTrackerManagerListener interface { // --------------------------------------------------- -type endsSenderReport struct { - first *buffer.RTCPSenderReportData - newest *buffer.RTCPSenderReportData - lastUpdated time.Time -} - type StreamTrackerManager struct { logger logger.Logger trackInfo atomic.Pointer[livekit.TrackInfo] @@ -76,7 +70,7 @@ type StreamTrackerManager struct { paused bool senderReportMu sync.RWMutex - senderReports [buffer.DefaultMaxLayerSpatial + 1]endsSenderReport + senderReports [buffer.DefaultMaxLayerSpatial + 1]*buffer.RTCPSenderReportData layerOffsets [buffer.DefaultMaxLayerSpatial + 1][buffer.DefaultMaxLayerSpatial + 1]uint32 closed core.Fuse @@ -557,8 +551,8 @@ func (s *StreamTrackerManager) maxExpectedLayerFromTrackInfo() { } func (s *StreamTrackerManager) updateLayerOffsetLocked(ref, other int32) { - srRef := s.senderReports[ref].newest - srOther := s.senderReports[other].newest + srRef := s.senderReports[ref] + srOther := s.senderReports[other] if srRef == nil || srRef.NTPTimestamp == 0 || srOther == nil || srOther.NTPTimestamp == 0 { return } @@ -598,7 +592,7 @@ func (s *StreamTrackerManager) updateLayerOffsetLocked(ref, other int32) { s.layerOffsets[ref][other] = offset } -func (s *StreamTrackerManager) SetRTCPSenderReportData(layer int32, srFirst *buffer.RTCPSenderReportData, srNewest *buffer.RTCPSenderReportData) { +func (s *StreamTrackerManager) SetRTCPSenderReportData(layer int32, srData *buffer.RTCPSenderReportData) { s.senderReportMu.Lock() defer s.senderReportMu.Unlock() @@ -606,9 +600,7 @@ func (s *StreamTrackerManager) SetRTCPSenderReportData(layer int32, srFirst *buf return } - s.senderReports[layer].first = srFirst - s.senderReports[layer].newest = srNewest - s.senderReports[layer].lastUpdated = time.Now() + s.senderReports[layer] = srData // (re)fill offsets as necessary for received layer. for i := int32(0); i < buffer.DefaultMaxLayerSpatial+1; i++ { @@ -624,12 +616,12 @@ func (s *StreamTrackerManager) SetRTCPSenderReportData(layer int32, srFirst *buf } } -func (s *StreamTrackerManager) GetRTCPSenderReportData(layer int32) (*buffer.RTCPSenderReportData, *buffer.RTCPSenderReportData) { +func (s *StreamTrackerManager) GetRTCPSenderReportData(layer int32) *buffer.RTCPSenderReportData { s.senderReportMu.Lock() defer s.senderReportMu.Unlock() if layer < 0 || int(layer) >= len(s.senderReports) { - return nil, nil + return nil } // SVC-TODO: better SVC detection @@ -638,7 +630,7 @@ func (s *StreamTrackerManager) GetRTCPSenderReportData(layer int32) (*buffer.RTC layer = 0 } - return s.senderReports[layer].first, s.senderReports[layer].newest + return s.senderReports[layer] } func (s *StreamTrackerManager) GetReferenceLayerRTPTimestamp(ts uint32, layer int32, referenceLayer int32) (uint32, error) {