NEW_GROUP_REQUEST support, part 1. MoqtTrackPublisher::AddObjectListener accepts MessageParameters. This enables NEW_GROUP_REQUEST support in MoqtOutgoingQueue for SUBSCRIBE. To follow: PUBLISH_OK, REQUEST_UPDATE, and MoqtRelayTrackPublisher. PiperOrigin-RevId: 983192591
diff --git a/quiche/quic/moqt/moqt_integration_test.cc b/quiche/quic/moqt/moqt_integration_test.cc index c0c5ab4..ee02ac3 100644 --- a/quiche/quic/moqt/moqt_integration_test.cc +++ b/quiche/quic/moqt/moqt_integration_test.cc
@@ -559,7 +559,7 @@ bool received_ok = false; ON_CALL(*track_publisher, expiration).WillByDefault(Return(std::nullopt)); EXPECT_CALL(*track_publisher, AddObjectListener) - .WillOnce([&](MoqtObjectListener* listener) { + .WillOnce([&](MoqtObjectListener* listener, const MessageParameters&) { listener->OnSubscribeAccepted(); }); EXPECT_CALL(subscribe_visitor_, OnReply) @@ -590,7 +590,7 @@ ON_CALL(*track_publisher, expiration) .WillByDefault(Return(quic::QuicTimeDelta::Zero())); EXPECT_CALL(*track_publisher, AddObjectListener) - .WillOnce([&](MoqtObjectListener* listener) { + .WillOnce([&](MoqtObjectListener* listener, const MessageParameters&) { listener->OnSubscribeAccepted(); }); EXPECT_CALL(subscribe_visitor_, OnReply) @@ -621,7 +621,7 @@ ON_CALL(*track_publisher, expiration) .WillByDefault(Return(quic::QuicTimeDelta::Zero())); EXPECT_CALL(*track_publisher, AddObjectListener) - .WillOnce([&](MoqtObjectListener* listener) { + .WillOnce([&](MoqtObjectListener* listener, const MessageParameters&) { listener->OnSubscribeAccepted(); }); EXPECT_CALL(subscribe_visitor_, OnReply)
diff --git a/quiche/quic/moqt/moqt_outgoing_queue.cc b/quiche/quic/moqt/moqt_outgoing_queue.cc index 77bf351..ffa22fe 100644 --- a/quiche/quic/moqt/moqt_outgoing_queue.cc +++ b/quiche/quic/moqt/moqt_outgoing_queue.cc
@@ -46,6 +46,9 @@ "flag."; return; } + if (key) { + expect_new_group_ = false; + } if (closed_) { QUICHE_BUG(MoqtOutgoingQueue_AddObject_closed) << "Trying to send objects on a closed queue.";
diff --git a/quiche/quic/moqt/moqt_outgoing_queue.h b/quiche/quic/moqt/moqt_outgoing_queue.h index f46f37e..15d6c24 100644 --- a/quiche/quic/moqt/moqt_outgoing_queue.h +++ b/quiche/quic/moqt/moqt_outgoing_queue.h
@@ -9,7 +9,6 @@ #include <cstdint> #include <memory> #include <optional> -#include <string> #include <utility> #include <vector> @@ -19,16 +18,15 @@ #include "quiche/quic/core/quic_clock.h" #include "quiche/quic/core/quic_default_clock.h" #include "quiche/quic/core/quic_time.h" -#include "quiche/quic/moqt/moqt_error.h" #include "quiche/quic/moqt/moqt_fetch_task.h" #include "quiche/quic/moqt/moqt_key_value_pair.h" -#include "quiche/quic/moqt/moqt_messages.h" #include "quiche/quic/moqt/moqt_names.h" #include "quiche/quic/moqt/moqt_object.h" #include "quiche/quic/moqt/moqt_priority.h" #include "quiche/quic/moqt/moqt_publisher.h" #include "quiche/quic/moqt/moqt_session_callbacks.h" #include "quiche/quic/moqt/moqt_types.h" +#include "quiche/common/quiche_callbacks.h" #include "quiche/common/quiche_circular_deque.h" #include "quiche/common/quiche_mem_slice.h" @@ -43,9 +41,20 @@ // frames that they produce. class MoqtOutgoingQueue : public MoqtTrackPublisher { public: - MoqtOutgoingQueue(FullTrackName track, const quic::QuicClock* clock = - quic::QuicDefaultClock::Get()) - : clock_(clock), track_(std::move(track)) {} + // If the caller does not provide a new_group_callback, then the track + // property DYNAMIC_GROUPS will be set to false. If a callback is provided, + // the caller commits to creating a new group. + MoqtOutgoingQueue( + FullTrackName track, + const quic::QuicClock* clock = quic::QuicDefaultClock::Get(), + quiche::MultiUseCallback<void()> new_group_callback = nullptr) + : clock_(clock), + track_(std::move(track)), + extensions_(std::nullopt, std::nullopt, std::nullopt, std::nullopt, + new_group_callback != nullptr ? std::optional<bool>(true) + : std::nullopt, + std::nullopt), + new_group_callback_(std::move(new_group_callback)) {} MoqtOutgoingQueue(const MoqtOutgoingQueue&) = delete; MoqtOutgoingQueue(MoqtOutgoingQueue&&) = default; @@ -61,9 +70,18 @@ std::optional<PublishedObject> GetCachedObject( uint64_t group, std::optional<uint64_t> subgroup, uint64_t min_object, uint64_t offset = 0) const override; - void AddObjectListener(MoqtObjectListener* listener) override { + void AddObjectListener(MoqtObjectListener* listener, + const MessageParameters& parameters) override { listeners_.insert(listener); listener->OnSubscribeAccepted(); + if (extensions_.dynamic_groups() && !expect_new_group_ && + parameters.new_group_request.has_value() && + (*parameters.new_group_request == 0 || queue_.empty() || + *parameters.new_group_request > current_group_id_) && + new_group_callback_ != nullptr) { + expect_new_group_ = true; + new_group_callback_(); + } } void RemoveObjectListener(MoqtObjectListener* listener) override { listeners_.erase(listener); @@ -158,6 +176,8 @@ absl::InlinedVector<Group, kMaxQueuedGroups> queue_; uint64_t current_group_id_ = -1; absl::flat_hash_set<MoqtObjectListener*> listeners_; + bool expect_new_group_ = false; + quiche::MultiUseCallback<void()> new_group_callback_; }; } // namespace moqt
diff --git a/quiche/quic/moqt/moqt_outgoing_queue_test.cc b/quiche/quic/moqt/moqt_outgoing_queue_test.cc index 4b7b76e..93267b7 100644 --- a/quiche/quic/moqt/moqt_outgoing_queue_test.cc +++ b/quiche/quic/moqt/moqt_outgoing_queue_test.cc
@@ -19,6 +19,7 @@ #include "quiche/quic/core/quic_time.h" #include "quiche/quic/moqt/moqt_error.h" #include "quiche/quic/moqt/moqt_fetch_task.h" +#include "quiche/quic/moqt/moqt_key_value_pair.h" #include "quiche/quic/moqt/moqt_names.h" #include "quiche/quic/moqt/moqt_object.h" #include "quiche/quic/moqt/moqt_priority.h" @@ -28,6 +29,7 @@ #include "quiche/quic/moqt/test_tools/moqt_mock_visitor.h" #include "quiche/common/platform/api/quiche_expect_bug.h" #include "quiche/common/platform/api/quiche_test.h" +#include "quiche/common/quiche_callbacks.h" #include "quiche/common/quiche_mem_slice.h" #include "quiche/common/test_tools/quiche_test_utils.h" #include "quiche/web_transport/web_transport.h" @@ -37,6 +39,7 @@ using ::quiche::test::IsOkAndHolds; using ::quiche::test::StatusIs; +using ::testing::_; using ::testing::AnyOf; using ::testing::ElementsAre; using ::testing::Field; @@ -45,9 +48,13 @@ class TestMoqtOutgoingQueue : public MoqtOutgoingQueue, public MoqtObjectListener { public: - TestMoqtOutgoingQueue() : MoqtOutgoingQueue(FullTrackName{"test", "track"}) { + TestMoqtOutgoingQueue( + quiche::MultiUseCallback<void()> new_group_callback = nullptr) + : MoqtOutgoingQueue(FullTrackName{"test", "track"}, + quic::QuicDefaultClock::Get(), + std::move(new_group_callback)) { EXPECT_CALL(*this, OnSubscribeAccepted).WillOnce(Return()); - AddObjectListener(this); + AddObjectListener(this, MessageParameters()); } void OnNewObjectAvailable(Location sequence, std::optional<uint64_t> subgroup, @@ -467,12 +474,96 @@ EXPECT_CALL(listener, OnTrackPublisherGone).WillOnce([&] { queue.RemoveObjectListener(&listener); }); - queue.AddObjectListener(&listener); + queue.AddObjectListener(&listener, MessageParameters()); } queue.RemoveAllSubscriptions(); EXPECT_FALSE(queue.HasSubscribers()); } +TEST(MoqtOutgoingQueue, NewGroupRequest) { + testing::MockFunction<void()> callback; + TestMoqtOutgoingQueue queue(callback.AsStdFunction()); + { + testing::InSequence seq; + // new_group_request = 1 when queue is empty. + EXPECT_CALL(callback, Call()); + EXPECT_CALL(queue, PublishObject(0, 0, "a")); + // new_group_request = 0 when current_group_id_ is 0. + EXPECT_CALL(callback, Call()); + EXPECT_CALL(queue, CloseStreamForGroup(0)); + EXPECT_CALL(queue, PublishObject(1, 0, "b")); + // new_group_request = 1 when current_group_id_ is 1 (stale; no callback). + EXPECT_CALL(queue, PublishObject(1, 1, "c")); + // new_group_request = 2 when current_group_id_ is 1. + EXPECT_CALL(callback, Call()); + EXPECT_CALL(queue, CloseStreamForGroup(1)); + EXPECT_CALL(queue, PublishObject(2, 0, "d")); + } + + MockMoqtObjectListener listener; + EXPECT_CALL(listener, OnSubscribeAccepted).Times(4); + EXPECT_CALL(listener, + OnNewObjectAvailable(Location(0, 0), testing::Optional(0), _)); + EXPECT_CALL(listener, + OnNewObjectAvailable(Location(0, 1), testing::Optional(0), _)); + EXPECT_CALL(listener, + OnNewObjectAvailable(Location(1, 0), testing::Optional(0), _)); + EXPECT_CALL(listener, + OnNewObjectAvailable(Location(1, 1), testing::Optional(0), _)); + EXPECT_CALL(listener, + OnNewObjectAvailable(Location(1, 2), testing::Optional(0), _)); + EXPECT_CALL(listener, + OnNewObjectAvailable(Location(2, 0), testing::Optional(0), _)); + + MessageParameters parameters; + parameters.new_group_request = 1; + queue.AddObjectListener(&listener, parameters); + queue.AddObject(quiche::QuicheMemSlice::Copy("a"), true); + + parameters.new_group_request = 0; + queue.AddObjectListener(&listener, parameters); + queue.AddObject(quiche::QuicheMemSlice::Copy("b"), true); + + parameters.new_group_request = 1; + queue.AddObjectListener(&listener, parameters); + queue.AddObject(quiche::QuicheMemSlice::Copy("c"), false); + + parameters.new_group_request = 2; + queue.AddObjectListener(&listener, parameters); + queue.AddObject(quiche::QuicheMemSlice::Copy("d"), true); +} + +TEST(MoqtOutgoingQueue, NewGroupRequestIgnoredWithoutCallback) { + TestMoqtOutgoingQueue queue; + { + testing::InSequence seq; + EXPECT_CALL(queue, PublishObject(0, 0, "a")); + EXPECT_CALL(queue, PublishObject(0, 1, "b")); + } + queue.AddObject(quiche::QuicheMemSlice::Copy("a"), true); + + MockMoqtObjectListener listener; + EXPECT_CALL(listener, OnSubscribeAccepted); + EXPECT_CALL(listener, + OnNewObjectAvailable(Location(0, 1), testing::Optional(0), _)); + MessageParameters parameters; + parameters.new_group_request = 0; + queue.AddObjectListener(&listener, parameters); + queue.AddObject(quiche::QuicheMemSlice::Copy("b"), false); +} + +TEST(MoqtOutgoingQueue, DynamicGroupsExtension) { + TestMoqtOutgoingQueue queue_without_callback; + EXPECT_FALSE(queue_without_callback.extensions().dynamic_groups()); + EXPECT_FALSE(queue_without_callback.extensions().contains( + static_cast<uint64_t>(ExtensionHeader::kDynamicGroups))); + + TestMoqtOutgoingQueue queue_with_callback([]() {}); + EXPECT_TRUE(queue_with_callback.extensions().dynamic_groups()); + EXPECT_TRUE(queue_with_callback.extensions().contains( + static_cast<uint64_t>(ExtensionHeader::kDynamicGroups))); +} + } // namespace } // namespace moqt::test
diff --git a/quiche/quic/moqt/moqt_publisher.h b/quiche/quic/moqt/moqt_publisher.h index 8d7ccab..88a28c7 100644 --- a/quiche/quic/moqt/moqt_publisher.h +++ b/quiche/quic/moqt/moqt_publisher.h
@@ -96,7 +96,8 @@ // Registers a listener with the track. The listener will be notified of all // newly arriving objects. The pointer to the listener must be valid until // removed. - virtual void AddObjectListener(MoqtObjectListener* listener) = 0; + virtual void AddObjectListener(MoqtObjectListener* listener, + const MessageParameters& parameters) = 0; virtual void RemoveObjectListener(MoqtObjectListener* listener) = 0; // Methods to return various track properties. Returns nullopt if the value is
diff --git a/quiche/quic/moqt/moqt_relay_publisher_test.cc b/quiche/quic/moqt/moqt_relay_publisher_test.cc index 024b3ef..7fb1f9a 100644 --- a/quiche/quic/moqt/moqt_relay_publisher_test.cc +++ b/quiche/quic/moqt/moqt_relay_publisher_test.cc
@@ -86,7 +86,7 @@ publisher_.GetTrack(FullTrackName("foo", "bar")); EXPECT_NE(track, nullptr); EXPECT_CALL(session_, Subscribe); - track->AddObjectListener(&object_listener_); + track->AddObjectListener(&object_listener_, MessageParameters()); track->RemoveObjectListener(&object_listener_); publisher_.OnPublishNamespaceDone(TrackNamespace({"foo"}), &session_); EXPECT_EQ(publisher_.GetTrack(FullTrackName("foo", "bar")), nullptr);
diff --git a/quiche/quic/moqt/moqt_relay_track_publisher.cc b/quiche/quic/moqt/moqt_relay_track_publisher.cc index 9ed5270..54bd709 100644 --- a/quiche/quic/moqt/moqt_relay_track_publisher.cc +++ b/quiche/quic/moqt/moqt_relay_track_publisher.cc
@@ -352,7 +352,8 @@ return object_it->second.ToPublishedObject(offset); } -void MoqtRelayTrackPublisher::AddObjectListener(MoqtObjectListener* listener) { +void MoqtRelayTrackPublisher::AddObjectListener(MoqtObjectListener* listener, + const MessageParameters&) { if (is_closing_) { return; }
diff --git a/quiche/quic/moqt/moqt_relay_track_publisher.h b/quiche/quic/moqt/moqt_relay_track_publisher.h index 44cea94..a869a23 100644 --- a/quiche/quic/moqt/moqt_relay_track_publisher.h +++ b/quiche/quic/moqt/moqt_relay_track_publisher.h
@@ -95,7 +95,8 @@ std::optional<PublishedObject> GetCachedObject( uint64_t group_id, std::optional<uint64_t> subgroup_id, uint64_t min_object, uint64_t offset = 0) const override; - void AddObjectListener(MoqtObjectListener* listener) override; + void AddObjectListener(MoqtObjectListener* listener, + const MessageParameters& parameters) override; void RemoveObjectListener(MoqtObjectListener* listener) override; std::optional<Location> largest_location() const override; const TrackExtensions& extensions() const override { return extensions_; }
diff --git a/quiche/quic/moqt/moqt_relay_track_publisher_test.cc b/quiche/quic/moqt/moqt_relay_track_publisher_test.cc index aae3a93..5ccd216 100644 --- a/quiche/quic/moqt/moqt_relay_track_publisher_test.cc +++ b/quiche/quic/moqt/moqt_relay_track_publisher_test.cc
@@ -62,7 +62,7 @@ void SubscribeAndOk() { EXPECT_CALL(*session_, Subscribe).WillOnce(testing::Return(true)); - publisher_.AddObjectListener(&listener_); + publisher_.AddObjectListener(&listener_, MessageParameters()); EXPECT_CALL(listener_, OnSubscribeAccepted); MessageParameters parameters; parameters.largest_object = kLargestLocation; @@ -119,7 +119,7 @@ TEST_F(MoqtRelayTrackPublisherTest, FiniteExpiration) { EXPECT_CALL(*session_, Subscribe).WillOnce(testing::Return(true)); - publisher_.AddObjectListener(&listener_); + publisher_.AddObjectListener(&listener_, MessageParameters()); EXPECT_CALL(listener_, OnSubscribeAccepted); MessageParameters parameters; parameters.largest_object = kLargestLocation; @@ -310,7 +310,7 @@ TEST_F(MoqtRelayTrackPublisherTest, SubscribeRejected) { EXPECT_CALL(*session_, Subscribe).WillOnce(testing::Return(true)); - publisher_.AddObjectListener(&listener_); + publisher_.AddObjectListener(&listener_, MessageParameters()); EXPECT_CALL(listener_, OnSubscribeRejected).WillOnce([this] { publisher_.RemoveObjectListener(&listener_); }); @@ -322,7 +322,7 @@ TEST_F(MoqtRelayTrackPublisherTest, LastListenerGone) { EXPECT_CALL(*session_, Subscribe).WillOnce(testing::Return(true)); - publisher_.AddObjectListener(&listener_); + publisher_.AddObjectListener(&listener_, MessageParameters()); EXPECT_CALL(*session_, Unsubscribe(kTrackName)); publisher_.RemoveObjectListener(&listener_); EXPECT_TRUE(track_deleted_); @@ -331,17 +331,17 @@ TEST_F(MoqtRelayTrackPublisherTest, SessionDies) { session_.reset(); EXPECT_CALL(listener_, OnSubscribeRejected); - publisher_.AddObjectListener(&listener_); + publisher_.AddObjectListener(&listener_, MessageParameters()); EXPECT_TRUE(track_deleted_); } TEST_F(MoqtRelayTrackPublisherTest, SecondListenerNoSubscribe) { EXPECT_CALL(*session_, Subscribe).WillOnce(testing::Return(true)); - publisher_.AddObjectListener(&listener_); + publisher_.AddObjectListener(&listener_, MessageParameters()); EXPECT_CALL(*session_, Subscribe).Times(0); EXPECT_CALL(listener_, OnSubscribeAccepted).Times(0); MockMoqtObjectListener listener2; - publisher_.AddObjectListener(&listener2); + publisher_.AddObjectListener(&listener2, MessageParameters()); EXPECT_CALL(listener_, OnSubscribeAccepted); EXPECT_CALL(listener2, OnSubscribeAccepted); MessageParameters parameters; @@ -352,7 +352,7 @@ TEST_F(MoqtRelayTrackPublisherTest, OnMalformedObject) { EXPECT_CALL(*session_, Subscribe).WillOnce(testing::Return(true)); - publisher_.AddObjectListener(&listener_); + publisher_.AddObjectListener(&listener_, MessageParameters()); EXPECT_CALL(listener_, OnTrackPublisherGone); publisher_.OnMalformedTrack(kTrackName); EXPECT_TRUE(track_deleted_); @@ -360,7 +360,7 @@ TEST_F(MoqtRelayTrackPublisherTest, DuplicateObject) { EXPECT_CALL(*session_, Subscribe).WillOnce(testing::Return(true)); - publisher_.AddObjectListener(&listener_); + publisher_.AddObjectListener(&listener_, MessageParameters()); Location location = kLargestLocation.Next(); EXPECT_CALL(listener_, OnNewObjectAvailable(location, Optional(0), /*publisher_priority=*/128)); @@ -384,7 +384,7 @@ TEST_F(MoqtRelayTrackPublisherTest, DuplicateObjectChangedMetadata) { EXPECT_CALL(*session_, Subscribe).WillOnce(testing::Return(true)); - publisher_.AddObjectListener(&listener_); + publisher_.AddObjectListener(&listener_, MessageParameters()); Location location = kLargestLocation.Next(); EXPECT_CALL(listener_, OnNewObjectAvailable(location, Optional(0), /*publisher_priority=*/128)); @@ -406,7 +406,7 @@ TEST_F(MoqtRelayTrackPublisherTest, DuplicateObjectChangedPayload) { EXPECT_CALL(*session_, Subscribe).WillOnce(testing::Return(true)); - publisher_.AddObjectListener(&listener_); + publisher_.AddObjectListener(&listener_, MessageParameters()); Location location = kLargestLocation.Next(); EXPECT_CALL(listener_, OnNewObjectAvailable(location, Optional(0), /*publisher_priority=*/128)); @@ -456,7 +456,7 @@ EXPECT_CALL(*session_, Subscribe).Times(0); MockMoqtObjectListener listener2; EXPECT_CALL(listener2, OnSubscribeAccepted); - publisher_.AddObjectListener(&listener2); + publisher_.AddObjectListener(&listener2, MessageParameters()); } TEST_F(MoqtRelayTrackPublisherTest, DatagramPreference) { @@ -643,7 +643,7 @@ testing::Field(&MessageParameters::oack_window_size, quic::QuicTimeDelta::FromMilliseconds(50)))) .WillOnce(testing::Return(true)); - publisher_.AddObjectListener(&listener_); + publisher_.AddObjectListener(&listener_, MessageParameters()); } } // namespace
diff --git a/quiche/quic/moqt/moqt_session.cc b/quiche/quic/moqt/moqt_session.cc index 31578a9..2570852 100644 --- a/quiche/quic/moqt/moqt_session.cc +++ b/quiche/quic/moqt/moqt_session.cc
@@ -606,7 +606,7 @@ stream_visitor_ptr->BindStream(stream); next_request_id_ += 2; ++next_local_track_alias_; - publisher->AddObjectListener(publisher_ptr); + publisher->AddObjectListener(publisher_ptr, parameters); return true; }
diff --git a/quiche/quic/moqt/moqt_session_test.cc b/quiche/quic/moqt/moqt_session_test.cc index a4fc9c8..f59b00b 100644 --- a/quiche/quic/moqt/moqt_session_test.cc +++ b/quiche/quic/moqt/moqt_session_test.cc
@@ -237,7 +237,7 @@ TrackExtensions extensions = TrackExtensions()) { MoqtObjectListener* listener_ptr = nullptr; EXPECT_CALL(*publisher, AddObjectListener) - .WillOnce([&](MoqtObjectListener* listener) { + .WillOnce([&](MoqtObjectListener* listener, const MessageParameters&) { listener_ptr = listener; listener->OnSubscribeAccepted(); }); @@ -673,8 +673,8 @@ MockTrackPublisher* track = CreateTrackPublisher(); MoqtObjectListener* listener; EXPECT_CALL(*track, AddObjectListener) - .WillOnce( - [&](MoqtObjectListener* listener_ptr) { listener = listener_ptr; }); + .WillOnce([&](MoqtObjectListener* listener_ptr, + const MessageParameters&) { listener = listener_ptr; }); bidi_wrapper_->ReceiveMessage(request); EXPECT_CALL(mock_bidi_stream_, @@ -691,8 +691,8 @@ MockTrackPublisher* track = CreateTrackPublisher(); MoqtObjectListener* listener; EXPECT_CALL(*track, AddObjectListener) - .WillOnce( - [&](MoqtObjectListener* listener_ptr) { listener = listener_ptr; }); + .WillOnce([&](MoqtObjectListener* listener_ptr, + const MessageParameters&) { listener = listener_ptr; }); bidi_wrapper_->ReceiveMessage(request); EXPECT_CALL(mock_bidi_stream_, Writev(ControlMessageOfType(MoqtMessageType::kRequestError), _)) @@ -711,7 +711,7 @@ MoqtSubscribe request = DefaultSubscribe(); MockTrackPublisher* track = CreateTrackPublisher(); EXPECT_CALL(*track, AddObjectListener) - .WillOnce([&](MoqtObjectListener* listener) { + .WillOnce([&](MoqtObjectListener* listener, const MessageParameters&) { EXPECT_CALL(*track, RemoveObjectListener); listener->OnSubscribeRejected(MoqtRequestErrorInfo( RequestErrorCode::kInternalError, std::nullopt, "Test error")); @@ -2458,7 +2458,7 @@ MoqtTrackStatus track_status = DefaultSubscribe(); EXPECT_CALL(*track, AddObjectListener) - .WillOnce([&](MoqtObjectListener* listener) { + .WillOnce([&](MoqtObjectListener* listener, const MessageParameters&) { EXPECT_CALL(*track, expiration) .WillRepeatedly( Return(quic::QuicTimeDelta::FromMilliseconds(10000))); @@ -2509,7 +2509,7 @@ MoqtTrackStatus track_status = DefaultSubscribe(); bool executed_AddObjectListener = false; EXPECT_CALL(*track, AddObjectListener) - .WillOnce([&](MoqtObjectListener* listener) { + .WillOnce([&](MoqtObjectListener* listener, const MessageParameters&) { EXPECT_CALL( mock_bidi_stream_, Writev(ControlMessageOfType(MoqtMessageType::kRequestError), _));
diff --git a/quiche/quic/moqt/moqt_subscribe_stream.cc b/quiche/quic/moqt/moqt_subscribe_stream.cc index 7b9f517..35578d8 100644 --- a/quiche/quic/moqt/moqt_subscribe_stream.cc +++ b/quiche/quic/moqt/moqt_subscribe_stream.cc
@@ -199,7 +199,7 @@ } } // Don't add the publisher until we know it's successful. - track_publisher->AddObjectListener(subscription_.get()); + track_publisher->AddObjectListener(subscription_.get(), message.parameters); return absl::OkStatus(); }
diff --git a/quiche/quic/moqt/moqt_track_status_stream.cc b/quiche/quic/moqt/moqt_track_status_stream.cc index 62c3ea1..8e4863b 100644 --- a/quiche/quic/moqt/moqt_track_status_stream.cc +++ b/quiche/quic/moqt/moqt_track_status_stream.cc
@@ -122,7 +122,7 @@ } // If the upstream subscription is already established, the code below will // invoke `OnSubscribeAccepted` immediately. - publisher_->AddObjectListener(this); + publisher_->AddObjectListener(this, message.parameters); return absl::OkStatus(); }
diff --git a/quiche/quic/moqt/test_tools/moqt_mock_visitor.h b/quiche/quic/moqt/test_tools/moqt_mock_visitor.h index 238235c..a237fde 100644 --- a/quiche/quic/moqt/test_tools/moqt_mock_visitor.h +++ b/quiche/quic/moqt/test_tools/moqt_mock_visitor.h
@@ -83,7 +83,9 @@ MOCK_METHOD(std::optional<PublishedObject>, GetCachedObject, (uint64_t, std::optional<uint64_t>, uint64_t, uint64_t), (const, override)); - MOCK_METHOD(void, AddObjectListener, (MoqtObjectListener * listener), + MOCK_METHOD(void, AddObjectListener, + (MoqtObjectListener * listener, + const MessageParameters& parameters), (override)); MOCK_METHOD(void, RemoveObjectListener, (MoqtObjectListener * listener), (override)); @@ -120,7 +122,8 @@ } return it->second.ToPublishedObject(); } - void AddObjectListener(MoqtObjectListener* listener) override { + void AddObjectListener(MoqtObjectListener* listener, + const MessageParameters&) override { listeners_.insert(listener); listener->OnSubscribeAccepted(); }
diff --git a/quiche/quic/moqt/tools/moqt_relay_test.cc b/quiche/quic/moqt/tools/moqt_relay_test.cc index 0bc2584..3915f5c 100644 --- a/quiche/quic/moqt/tools/moqt_relay_test.cc +++ b/quiche/quic/moqt/tools/moqt_relay_test.cc
@@ -139,7 +139,7 @@ std::shared_ptr<MoqtTrackPublisher> track = upstream_.publisher()->GetTrack(FullTrackName("foo", "bar")); EXPECT_NE(track, nullptr); - track->AddObjectListener(&object_listener); + track->AddObjectListener(&object_listener, MessageParameters()); track->RemoveObjectListener(&object_listener); // Track should have been destroyed.