Make RtpVideoSender::OnBitrateUpdated lock-free SetFecAllowed() is called by the encoder, and OnBitrateUpdated() reads the value on the transport queue. This was the last reason for OnBitrateUpdated() to take mutex_, which blocked the worker thread whenever a bitrate update coincided with packetization of an encoded frame. Post the value to the transport queue instead and guard it with transport_checker_. FecControllerOverride is public API, so this works for calls from any thread. Bug: webrtc:42223727 Change-Id: Iad4757cfd3092f39fda1086fcb942c3d56c5b8bf Reviewed-on: https://webrtc-review.googlesource.com/c/src/+/506163 Commit-Queue: Tomas Gunnarsson <tommi@webrtc.org> Reviewed-by: Erik Språng <sprang@webrtc.org> Cr-Commit-Position: refs/heads/main@{#48789}
diff --git a/call/BUILD.gn b/call/BUILD.gn index 9546e3f..5dda1e1 100644 --- a/call/BUILD.gn +++ b/call/BUILD.gn
@@ -815,6 +815,7 @@ ":rtp_sender", ":video_send_stream_api", "../api:bitrate_allocation", + "../api:fec_controller_api", "../api:field_trials", "../api:frame_transformer_interface", "../api:make_ref_counted",
diff --git a/call/rtp_video_sender.cc b/call/rtp_video_sender.cc index 6797ed2..c283d3f 100644 --- a/call/rtp_video_sender.cc +++ b/call/rtp_video_sender.cc
@@ -842,7 +842,6 @@ int framerate) { RTC_DCHECK_RUN_ON(&transport_checker_); // Substract overhead from bitrate. - MutexLock lock(&mutex_); if (transport_overhead_bytes_per_packet_ != update.packet_overhead.bytes<size_t>()) { transport_overhead_bytes_per_packet_ = @@ -977,8 +976,12 @@ } void RtpVideoSender::SetFecAllowed(bool fec_allowed) { - MutexLock lock(&mutex_); - fec_allowed_ = fec_allowed; + // Called by the encoder, which may run on any thread. `fec_allowed_` is only + // used on the transport queue, so apply the value there. + transport_queue_.PostTask(SafeTask(safety_.flag(), [this, fec_allowed] { + RTC_DCHECK_RUN_ON(&transport_checker_); + fec_allowed_ = fec_allowed; + })); } void RtpVideoSender::OnPacketFeedbackVector(
diff --git a/call/rtp_video_sender.h b/call/rtp_video_sender.h index 11ae567..7a12287 100644 --- a/call/rtp_video_sender.h +++ b/call/rtp_video_sender.h
@@ -194,7 +194,7 @@ bool active_ RTC_GUARDED_BY(mutex_) = false; const std::unique_ptr<FecController> fec_controller_; - bool fec_allowed_ RTC_GUARDED_BY(mutex_) = true; + bool fec_allowed_ RTC_GUARDED_BY(transport_checker_) = true; // Rtp modules are assumed to be sorted in simulcast index order. const std::vector<webrtc_internal_rtp_video_sender::RtpStreamSender>
diff --git a/call/rtp_video_sender_unittest.cc b/call/rtp_video_sender_unittest.cc index acb135a..8406c2c 100644 --- a/call/rtp_video_sender_unittest.cc +++ b/call/rtp_video_sender_unittest.cc
@@ -16,6 +16,7 @@ #include <memory> #include <optional> #include <span> +#include <utility> #include <vector> #include "absl/strings/string_view.h" @@ -23,6 +24,7 @@ #include "api/call/transport.h" #include "api/crypto/crypto_options.h" #include "api/environment/environment.h" +#include "api/fec_controller.h" #include "api/frame_transformer_interface.h" #include "api/make_ref_counted.h" #include "api/rtp_header_extension_id.h" @@ -77,6 +79,7 @@ namespace { using ::testing::_; +using ::testing::DoAll; using ::testing::ElementsAre; using ::testing::ElementsAreArray; using ::testing::Ge; @@ -84,6 +87,7 @@ using ::testing::IsNull; using ::testing::NiceMock; using ::testing::NotNull; +using ::testing::Return; using ::testing::SaveArg; using ::testing::SizeIs; @@ -105,6 +109,39 @@ MOCK_METHOD(void, OnReceivedIntraFrameRequest, (uint32_t), (override)); }; +class MockFecController : public FecController { + public: + MOCK_METHOD(void, + SetProtectionCallback, + (VCMProtectionCallback * protection_callback), + (override)); + MOCK_METHOD(void, + SetProtectionMethod, + (bool enable_fec, bool enable_nack), + (override)); + MOCK_METHOD(void, + SetEncodingData, + (size_t width, + size_t height, + size_t num_temporal_layers, + size_t max_payload_size), + (override)); + MOCK_METHOD(uint32_t, + UpdateFecRates, + (uint32_t estimated_bitrate_bps, + int actual_framerate, + uint8_t fraction_lost, + std::vector<bool> loss_mask_vector, + int64_t round_trip_time_ms), + (override)); + MOCK_METHOD(void, + UpdateWithEncodedData, + (size_t encoded_image_length, + VideoFrameType encoded_image_frametype), + (override)); + MOCK_METHOD(bool, UseLossVectorMask, (), (override)); +}; + RtpSenderObservers CreateObservers( RtcpIntraFrameObserver* intra_frame_callback, ReportBlockDataObserver* report_block_data_observer, @@ -180,7 +217,8 @@ scoped_refptr<FrameTransformerInterface> frame_transformer, const std::vector<int>& payload_types, absl::string_view field_trials = "", - const std::vector<uint32_t>& csrcs = {}) + const std::vector<uint32_t>& csrcs = {}, + std::unique_ptr<FecController> fec_controller = nullptr) : time_controller_(Timestamp::Millis(1000000)), env_(CreateTestEnvironment( {.field_trials = field_trials, .time = &time_controller_})), @@ -200,14 +238,17 @@ VideoEncoderConfig::ContentType::kRealtimeVideo) { transport_controller_.EnsureStarted(); std::map<uint32_t, RtpState> suspended_ssrcs; + if (!fec_controller) { + fec_controller = std::make_unique<FecControllerDefault>(env_); + } router_ = std::make_unique<RtpVideoSender>( env_, time_controller_.GetMainThread(), suspended_ssrcs, suspended_payload_states, config_.rtp, config_.rtcp_report_interval_ms, &transport_, CreateObservers(&encoder_feedback_, &stats_proxy_, &stats_proxy_, &stats_proxy_, frame_count_observer, &stats_proxy_), - &transport_controller_, std::make_unique<FecControllerDefault>(env_), - nullptr, CryptoOptions{}, frame_transformer); + &transport_controller_, std::move(fec_controller), nullptr, + CryptoOptions{}, frame_transformer); } RtpVideoSenderTestFixture( const std::vector<uint32_t>& ssrcs, @@ -1568,6 +1609,40 @@ } } +TEST(RtpVideoSenderTest, ReservesBitrateForFecOnlyWhenAllowed) { + constexpr uint32_t kTargetBitrateBps = 300'000; + // The rate that the FEC controller leaves for the encoder. + constexpr uint32_t kRateWithFecBps = 100'000; + auto fec_controller = std::make_unique<NiceMock<MockFecController>>(); + uint32_t payload_bitrate_bps = 0; + ON_CALL(*fec_controller, UpdateFecRates) + .WillByDefault( + DoAll(SaveArg<0>(&payload_bitrate_bps), Return(kRateWithFecBps))); + RtpVideoSenderTestFixture test( + {kSsrc1}, {}, kPayloadType, {}, /*frame_count_observer=*/nullptr, + /*frame_transformer=*/nullptr, /*payload_types=*/{}, + /*field_trials=*/"", /*csrcs=*/{}, std::move(fec_controller)); + test.SetSending(true); + const BitrateAllocationUpdate update = + CreateBitrateAllocationUpdate(kTargetBitrateBps); + + test.router()->OnBitrateUpdated(update, /*framerate=*/30); + EXPECT_EQ(test.router()->GetPayloadBitrateBps(), kRateWithFecBps); + + // SetFecAllowed() may be called on any thread. The value applies to bitrate + // updates once it has been handed to the transport queue. + test.router()->SetFecAllowed(false); + test.AdvanceTime(TimeDelta::Zero()); + test.router()->OnBitrateUpdated(update, /*framerate=*/30); + EXPECT_GT(payload_bitrate_bps, kRateWithFecBps); + EXPECT_EQ(test.router()->GetPayloadBitrateBps(), payload_bitrate_bps); + + test.router()->SetFecAllowed(true); + test.AdvanceTime(TimeDelta::Zero()); + test.router()->OnBitrateUpdated(update, /*framerate=*/30); + EXPECT_EQ(test.router()->GetPayloadBitrateBps(), kRateWithFecBps); +} + TEST(RtpVideoSenderTest, ClearsPendingPacketsOnInactivation) { RtpVideoSenderTestFixture test({kSsrc1}, {kRtxSsrc1}, kPayloadType, {}); test.SetSending(true);