Fold all of the control message handlers from the control stream object to the session itself.

The control stream class is likely going to go away once we split the control streams into two.

PiperOrigin-RevId: 936917501
diff --git a/quiche/quic/moqt/moqt_session.cc b/quiche/quic/moqt/moqt_session.cc
index 9a7c7e4..1b4d14b 100644
--- a/quiche/quic/moqt/moqt_session.cc
+++ b/quiche/quic/moqt/moqt_session.cc
@@ -970,94 +970,92 @@
 absl::Status MoqtSession::ControlStream::OnRawControlMessage(
     const MoqtRawControlMessage& message) {
   return ControlMessageDispatcher::DispatchControlMessage(
-      *this, message_parser(), message, "control");
+      *session_, message_parser(), message, "control");
 }
 
-absl::Status MoqtSession::ControlStream::OnControlMessage(
-    const MoqtSetup& message) {
-  if (session_->parameters_.perspective == Perspective::IS_SERVER) {
-    session_->peer_supports_object_ack_ =
-        message.parameters.support_object_acks.value_or(
-            kDefaultSupportObjectAcks);
-    session_->peer_max_request_id_ =
+absl::Status MoqtSession::OnControlMessage(const MoqtSetup& message) {
+  if (parameters_.perspective == Perspective::IS_SERVER) {
+    peer_supports_object_ack_ = message.parameters.support_object_acks.value_or(
+        kDefaultSupportObjectAcks);
+    peer_max_request_id_ =
         message.parameters.max_request_id.value_or(kDefaultMaxRequestId);
     QUICHE_DLOG(INFO) << "Received CLIENT_SETUP";
     MoqtSetup response;
-    session_->parameters_.ToSetupParameters(response.parameters);
-    QUICHE_RETURN_IF_ERROR(
-        SendOrBufferMessage(session_->framer_.SerializeSetup(response)));
+    parameters_.ToSetupParameters(response.parameters);
+    SendControlMessage(framer_.SerializeSetup(response));
     QUICHE_DLOG(INFO) << "Sent SERVER_SETUP";
     // TODO: handle path.
-    std::move(session_->callbacks_.session_established_callback)();
+    std::move(callbacks_.session_established_callback)();
     return absl::OkStatus();
   } else {
-    session_->peer_supports_object_ack_ =
-        message.parameters.support_object_acks.value_or(
-            kDefaultSupportObjectAcks);
+    peer_supports_object_ack_ = message.parameters.support_object_acks.value_or(
+        kDefaultSupportObjectAcks);
     QUIC_DLOG(INFO) << ENDPOINT << "Received the SETUP message";
     // TODO: handle path.
-    session_->peer_max_request_id_ =
+    peer_max_request_id_ =
         message.parameters.max_request_id.value_or(kDefaultMaxRequestId);
-    std::move(session_->callbacks_.session_established_callback)();
+    std::move(callbacks_.session_established_callback)();
     return absl::OkStatus();
   }
 }
 
-absl::Status MoqtSession::ControlStream::OnControlMessage(
-    const MoqtSubscribe& message) {
-  if (!session_->ValidateRequestId(message.request_id)) {
+absl::Status MoqtSession::OnControlMessage(const MoqtSubscribe& message) {
+  if (!ValidateRequestId(message.request_id)) {
     return absl::OkStatus();
   }
   QUIC_DLOG(INFO) << ENDPOINT << "Received a SUBSCRIBE for "
                   << message.full_track_name;
-  if (session_->sent_goaway_) {
+  if (sent_goaway_) {
     QUIC_DLOG(INFO) << ENDPOINT << "Received a SUBSCRIBE after GOAWAY";
-    return SendRequestError(message.request_id, RequestErrorCode::kUnauthorized,
-                            std::nullopt, "SUBSCRIBE after GOAWAY");
+    SendRequestErrorOnControlStream(message.request_id,
+                                    RequestErrorCode::kUnauthorized,
+                                    std::nullopt, "SUBSCRIBE after GOAWAY");
+    return absl::OkStatus();
   }
-  if (session_->subscribed_track_names_.contains(message.full_track_name)) {
-    return SendRequestError(message.request_id,
-                            RequestErrorCode::kDuplicateSubscription,
-                            std::nullopt, "");
+  if (subscribed_track_names_.contains(message.full_track_name)) {
+    SendRequestErrorOnControlStream(message.request_id,
+                                    RequestErrorCode::kDuplicateSubscription,
+                                    std::nullopt, "");
+    return absl::OkStatus();
   }
   const FullTrackName& track_name = message.full_track_name;
   std::shared_ptr<MoqtTrackPublisher> track_publisher =
-      session_->publisher_->GetTrack(track_name);
+      publisher_->GetTrack(track_name);
   if (track_publisher == nullptr) {
     QUIC_DLOG(INFO) << ENDPOINT << "SUBSCRIBE for " << track_name
                     << " rejected by the application: does not exist";
-    return SendRequestError(message.request_id, RequestErrorCode::kDoesNotExist,
-                            std::nullopt, "not found");
+    SendRequestErrorOnControlStream(message.request_id,
+                                    RequestErrorCode::kDoesNotExist,
+                                    std::nullopt, "not found");
+    return absl::OkStatus();
   }
 
   MoqtPublishingMonitorInterface* monitoring = nullptr;
   auto monitoring_it =
-      session_->monitoring_interfaces_for_published_tracks_.find(track_name);
-  if (monitoring_it !=
-      session_->monitoring_interfaces_for_published_tracks_.end()) {
+      monitoring_interfaces_for_published_tracks_.find(track_name);
+  if (monitoring_it != monitoring_interfaces_for_published_tracks_.end()) {
     monitoring = monitoring_it->second;
-    session_->monitoring_interfaces_for_published_tracks_.erase(monitoring_it);
+    monitoring_interfaces_for_published_tracks_.erase(monitoring_it);
   }
 
   MoqtTrackPublisher* track_publisher_ptr = track_publisher.get();
   auto subscription = std::make_unique<SubscriptionPublisher>(
-      session_->framer_, track_publisher, this, message.request_id,
-      session_->next_local_track_alias_++, message.parameters, session_,
-      monitoring, session_->callbacks_.clock, session_->trace_recorder_, false);
+      framer_, track_publisher, GetControlStream(), message.request_id,
+      next_local_track_alias_++, message.parameters, this, monitoring,
+      callbacks_.clock, trace_recorder_, false);
   SubscriptionPublisher* subscription_ptr = subscription.get();
-  auto [it, success] = session_->published_subscriptions_.emplace(
+  auto [it, success] = published_subscriptions_.emplace(
       message.request_id, std::move(subscription));
   if (!success) {
     QUICHE_NOTREACHED();  // ValidateRequestId() should have caught this.
   }
-  session_->subscribed_track_names_.insert(message.full_track_name);
+  subscribed_track_names_.insert(message.full_track_name);
   track_publisher_ptr->AddObjectListener(subscription_ptr);
   return absl::OkStatus();
 }
 
-absl::Status MoqtSession::ControlStream::OnControlMessage(
-    const MoqtSubscribeOk& message) {
-  RemoteTrack* track = session_->RemoteTrackById(message.request_id);
+absl::Status MoqtSession::OnControlMessage(const MoqtSubscribeOk& message) {
+  RemoteTrack* track = RemoteTrackById(message.request_id);
   if (track == nullptr) {
     QUIC_DLOG(INFO) << ENDPOINT << "Received the SUBSCRIBE_OK for "
                     << "request_id = " << message.request_id
@@ -1081,39 +1079,37 @@
   SubscribeRemoteTrack* subscribe =
       absl::down_cast<SubscribeRemoteTrack*>(track);
   if (!subscribe->set_track_alias(message.track_alias)) {
-    OnFatalError(absl::AlreadyExistsError(""));
-    return absl::OkStatus();
+    return absl::AlreadyExistsError("Duplicate track alias");
   }
   subscribe->OnObjectOrOk(
       SubscribeOkData(message.parameters, message.extensions));
   return absl::OkStatus();
 }
 
-absl::Status MoqtSession::ControlStream::OnControlMessage(
-    const MoqtRequestOk& message) {
-  if (session_->upstream_by_id_.contains(message.request_id)) {
+absl::Status MoqtSession::OnControlMessage(const MoqtRequestOk& message) {
+  if (upstream_by_id_.contains(message.request_id)) {
     return absl::InvalidArgumentError(
         "Received REQUEST_OK for SUBSCRIBE, FETCH, or PUBLISH");
   }
   // Response to REQUEST_UPDATE for a subscribe.
-  auto ru_it = session_->pending_subscribe_updates_.find(message.request_id);
-  if (ru_it != session_->pending_subscribe_updates_.end()) {
-    auto sub_it = session_->subscribe_by_name_.find(ru_it->second.name);
-    if (sub_it == session_->subscribe_by_name_.end()) {
+  auto ru_it = pending_subscribe_updates_.find(message.request_id);
+  if (ru_it != pending_subscribe_updates_.end()) {
+    auto sub_it = subscribe_by_name_.find(ru_it->second.name);
+    if (sub_it == subscribe_by_name_.end()) {
       std::move(ru_it->second.response_callback)(
           MoqtRequestErrorInfo{RequestErrorCode::kDoesNotExist, std::nullopt,
                                "subscription does not exist anymore"});
-      session_->pending_subscribe_updates_.erase(ru_it);
+      pending_subscribe_updates_.erase(ru_it);
       return absl::OkStatus();
     }
     sub_it->second->Update(ru_it->second.parameters);
     std::move(ru_it->second.response_callback)(MessageParameters());
-    session_->pending_subscribe_updates_.erase(ru_it);
+    pending_subscribe_updates_.erase(ru_it);
     return absl::OkStatus();
   }
   // Response to PUBLISH_NAMESPACE.
-  auto pn_it = session_->publish_namespace_by_id_.find(message.request_id);
-  if (pn_it != session_->publish_namespace_by_id_.end()) {
+  auto pn_it = publish_namespace_by_id_.find(message.request_id);
+  if (pn_it != publish_namespace_by_id_.end()) {
     if (pn_it->second.response_callback == nullptr) {
       return absl::InvalidArgumentError(
           "Multiple responses for PUBLISH_NAMESPACE");
@@ -1130,12 +1126,11 @@
   return absl::OkStatus();
 }
 
-absl::Status MoqtSession::ControlStream::OnControlMessage(
-    const MoqtRequestError& message) {
+absl::Status MoqtSession::OnControlMessage(const MoqtRequestError& message) {
   MoqtRequestErrorInfo error_info{message.error_code, message.retry_interval,
                                   message.reason_phrase};
   // TODO(martinduke): Do something with retry_interval.
-  RemoteTrack* track = session_->RemoteTrackById(message.request_id);
+  RemoteTrack* track = RemoteTrackById(message.request_id);
   if (track != nullptr) {
     // It's in response to SUBSCRIBE or FETCH.
     if (!track->ErrorIsAllowed()) {
@@ -1159,30 +1154,29 @@
         subscribe->visitor()->OnReply(subscribe->full_track_name(), error_info);
       }
     }
-    if (!session_->is_closing_) {
+    if (!is_closing_) {
       // The visitor might have closed the session.
       track->Destroy();
     }
     return absl::OkStatus();
   }
   // Response to REQUEST_UPDATE for a subscribe.
-  auto ru_it = session_->pending_subscribe_updates_.find(message.request_id);
-  if (ru_it != session_->pending_subscribe_updates_.end()) {
+  auto ru_it = pending_subscribe_updates_.find(message.request_id);
+  if (ru_it != pending_subscribe_updates_.end()) {
     std::move(ru_it->second.response_callback)(error_info);
-    session_->pending_subscribe_updates_.erase(ru_it);
+    pending_subscribe_updates_.erase(ru_it);
     return absl::OkStatus();
   }
   // Response to PUBLISH_NAMESPACE.
-  auto pn_it = session_->publish_namespace_by_id_.find(message.request_id);
-  if (pn_it != session_->publish_namespace_by_id_.end()) {
+  auto pn_it = publish_namespace_by_id_.find(message.request_id);
+  if (pn_it != publish_namespace_by_id_.end()) {
     if (pn_it->second.response_callback == nullptr) {
       return absl::InvalidArgumentError(
           "Multiple responses for PUBLISH_NAMESPACE");
     }
     std::move(pn_it->second.response_callback)(error_info);
-    session_->publish_namespace_by_namespace_.erase(
-        pn_it->second.track_namespace);
-    session_->publish_namespace_by_id_.erase(pn_it);
+    publish_namespace_by_namespace_.erase(pn_it->second.track_namespace);
+    publish_namespace_by_id_.erase(pn_it);
     return absl::OkStatus();
   }
   // Response to SUBSCRIBE_NAMESPACE is handled in the NamespaceStream.
@@ -1193,51 +1187,47 @@
   return absl::OkStatus();
 }
 
-absl::Status MoqtSession::ControlStream::OnControlMessage(
-    const MoqtUnsubscribe& message) {
-  auto it = session_->published_subscriptions_.find(message.request_id);
-  if (it == session_->published_subscriptions_.end()) {
+absl::Status MoqtSession::OnControlMessage(const MoqtUnsubscribe& message) {
+  auto it = published_subscriptions_.find(message.request_id);
+  if (it == published_subscriptions_.end()) {
     return absl::OkStatus();
   }
   QUIC_DLOG(INFO) << ENDPOINT << "Received an UNSUBSCRIBE for "
                   << it->second->publisher().GetTrackName();
-  session_->PublishIsDone(message.request_id);
+  PublishIsDone(message.request_id);
   return absl::OkStatus();
 }
 
-absl::Status MoqtSession::ControlStream::OnControlMessage(
-    const MoqtPublishDone& message) {
-  auto it = session_->upstream_by_id_.find(message.request_id);
-  if (it == session_->upstream_by_id_.end()) {
+absl::Status MoqtSession::OnControlMessage(const MoqtPublishDone& message) {
+  auto it = upstream_by_id_.find(message.request_id);
+  if (it == upstream_by_id_.end()) {
     return absl::OkStatus();
   }
   auto* subscribe = absl::down_cast<SubscribeRemoteTrack*>(it->second.get());
   QUIC_DLOG(INFO) << ENDPOINT << "Received a PUBLISH_DONE for "
                   << it->second->full_track_name();
-  subscribe->OnPublishDone(message.stream_count, session_->callbacks_.clock,
-                           session_->alarm_factory_.get());
+  subscribe->OnPublishDone(message.stream_count, callbacks_.clock,
+                           alarm_factory_.get());
   return absl::OkStatus();
 }
 
-absl::Status MoqtSession::ControlStream::OnControlMessage(
-    const MoqtRequestUpdate& message) {
-  auto it =
-      session_->published_subscriptions_.find(message.existing_request_id);
-  if (it != session_->published_subscriptions_.end()) {
+absl::Status MoqtSession::OnControlMessage(const MoqtRequestUpdate& message) {
+  auto it = published_subscriptions_.find(message.existing_request_id);
+  if (it != published_subscriptions_.end()) {
     // It's updating SUBSCRIBE.
     it->second->Update(message.parameters);
     // TODO(martinduke): There should be an MoqtResponseCallback sent to the
     // application, rather than automatic OK.
-    return SendRequestOk(message.request_id, MessageParameters());
+    SendControlMessage(framer_.SerializeRequestOk(
+        MoqtRequestOk{.request_id = message.request_id}));
+    return absl::OkStatus();
   }
-  auto pn_it =
-      session_->publish_namespace_by_id_.find(message.existing_request_id);
-  if (pn_it != session_->publish_namespace_by_id_.end()) {
+  auto pn_it = publish_namespace_by_id_.find(message.existing_request_id);
+  if (pn_it != publish_namespace_by_id_.end()) {
     // It's updating PUBLISH_NAMESPACE.
-    quiche::QuicheWeakPtr<MoqtSessionInterface> session_weakptr =
-        session_->GetWeakPtr();
+    quiche::QuicheWeakPtr<MoqtSessionInterface> session_weakptr = GetWeakPtr();
     TrackNamespace track_namespace = pn_it->second.track_namespace;
-    session_->callbacks().incoming_publish_namespace_callback(
+    callbacks().incoming_publish_namespace_callback(
         track_namespace, message.parameters,
         [&](std::variant<MessageParameters, MoqtRequestErrorInfo> response) {
           MoqtSession* session =
@@ -1251,14 +1241,16 @@
                       const MessageParameters& parameters) {
                     // In draft-18, there are no useful parameters in
                     // PUBLISH_NAMESPACE_OK, but Issue #1639 would change that.
-                    CheckStatus(SendRequestOk(request_id, parameters));
+                    SendControlMessage(framer_.SerializeRequestOk(MoqtRequestOk{
+                        .request_id = request_id, .parameters = parameters}));
                   },
                   [this, id = message.request_id, track_ns = track_namespace](
                       const MoqtRequestErrorInfo& error_info) {
-                    CheckStatus(SendRequestError(id, error_info));
-                    session_->incoming_publish_namespaces_by_id_.erase(id);
-                    session_->incoming_publish_namespaces_by_namespace_.erase(
-                        track_ns);
+                    SendRequestErrorOnControlStream(id, error_info.error_code,
+                                                    error_info.retry_interval,
+                                                    error_info.reason_phrase);
+                    incoming_publish_namespaces_by_id_.erase(id);
+                    incoming_publish_namespaces_by_namespace_.erase(track_ns);
                   }},
               response);
         });
@@ -1266,35 +1258,38 @@
   }
   // TODO(martinduke): Check all the request types.
   // Does not match any known request.
-  return SendRequestError(message.request_id, RequestErrorCode::kNotSupported,
-                          std::nullopt, "No support for update of this type");
+  SendRequestErrorOnControlStream(message.request_id,
+                                  RequestErrorCode::kNotSupported, std::nullopt,
+                                  "No support for update of this type");
+  return absl::OkStatus();
 }
 
-absl::Status MoqtSession::ControlStream::OnControlMessage(
+absl::Status MoqtSession::OnControlMessage(
     const MoqtPublishNamespace& message) {
-  if (!session_->ValidateRequestId(message.request_id)) {
+  if (!ValidateRequestId(message.request_id)) {
     return absl::OkStatus();
   }
-  if (session_->sent_goaway_) {
+  if (sent_goaway_) {
     QUIC_DLOG(INFO) << ENDPOINT << "Received a PUBLISH_NAMESPACE after GOAWAY";
-    return SendRequestError(message.request_id, RequestErrorCode::kUnauthorized,
-                            std::nullopt, "PUBLISH_NAMESPACE after GOAWAY");
+    SendRequestErrorOnControlStream(
+        message.request_id, RequestErrorCode::kUnauthorized, std::nullopt,
+        "PUBLISH_NAMESPACE after GOAWAY");
+    return absl::OkStatus();
   }
   QUIC_DLOG(INFO) << ENDPOINT << "Received a PUBLISH_NAMESPACE for "
                   << message.track_namespace;
-  auto [it, inserted] =
-      session_->incoming_publish_namespaces_by_namespace_.emplace(
-          message.track_namespace, message.request_id);
+  auto [it, inserted] = incoming_publish_namespaces_by_namespace_.emplace(
+      message.track_namespace, message.request_id);
   if (!inserted) {
-    return SendRequestError(message.request_id,
-                            RequestErrorCode::kDuplicateSubscription,
-                            std::nullopt, "Duplicate PUBLISH_NAMESPACE");
+    SendRequestErrorOnControlStream(
+        message.request_id, RequestErrorCode::kDuplicateSubscription,
+        std::nullopt, "Duplicate PUBLISH_NAMESPACE");
+    return absl::OkStatus();
   }
-  quiche::QuicheWeakPtr<MoqtSessionInterface> session_weakptr =
-      session_->GetWeakPtr();
-  session_->incoming_publish_namespaces_by_id_[message.request_id] =
+  quiche::QuicheWeakPtr<MoqtSessionInterface> session_weakptr = GetWeakPtr();
+  incoming_publish_namespaces_by_id_[message.request_id] =
       message.track_namespace;
-  session_->callbacks_.incoming_publish_namespace_callback(
+  callbacks_.incoming_publish_namespace_callback(
       message.track_namespace, message.parameters,
       [&](std::variant<MessageParameters, MoqtRequestErrorInfo> response) {
         MoqtSession* session =
@@ -1308,114 +1303,116 @@
                     const MessageParameters& parameters) {
                   // In draft-18, there are no useful parameters in
                   // PUBLISH_NAMESPACE_OK, but Issue #1639 would change that.
-                  CheckStatus(SendRequestOk(request_id, parameters));
+                  SendControlMessage(framer_.SerializeRequestOk(MoqtRequestOk{
+                      .request_id = request_id, .parameters = parameters}));
                 },
                 [this, id = message.request_id,
                  track_ns = message.track_namespace](
                     const MoqtRequestErrorInfo& error_info) {
-                  CheckStatus(SendRequestError(id, error_info));
-                  session_->incoming_publish_namespaces_by_id_.erase(id);
-                  session_->incoming_publish_namespaces_by_namespace_.erase(
-                      track_ns);
+                  SendRequestErrorOnControlStream(id, error_info.error_code,
+                                                  error_info.retry_interval,
+                                                  error_info.reason_phrase);
+                  incoming_publish_namespaces_by_id_.erase(id);
+                  incoming_publish_namespaces_by_namespace_.erase(track_ns);
                 }},
             response);
       });
   return absl::OkStatus();
 }
 
-absl::Status MoqtSession::ControlStream::OnControlMessage(
+absl::Status MoqtSession::OnControlMessage(
     const MoqtPublishNamespaceDone& message) {
-  auto it =
-      session_->incoming_publish_namespaces_by_id_.find(message.request_id);
-  if (it == session_->incoming_publish_namespaces_by_id_.end()) {
+  auto it = incoming_publish_namespaces_by_id_.find(message.request_id);
+  if (it == incoming_publish_namespaces_by_id_.end()) {
     return absl::OkStatus();
   }
-  session_->callbacks_.incoming_publish_namespace_callback(
-      it->second, std::nullopt, nullptr);
-  session_->incoming_publish_namespaces_by_namespace_.erase(it->second);
-  session_->incoming_publish_namespaces_by_id_.erase(it);
+  callbacks_.incoming_publish_namespace_callback(it->second, std::nullopt,
+                                                 nullptr);
+  incoming_publish_namespaces_by_namespace_.erase(it->second);
+  incoming_publish_namespaces_by_id_.erase(it);
   return absl::OkStatus();
 }
 
-absl::Status MoqtSession::ControlStream::OnControlMessage(
+absl::Status MoqtSession::OnControlMessage(
     const MoqtPublishNamespaceCancel& message) {
-  auto it = session_->publish_namespace_by_id_.find(message.request_id);
-  if (it == session_->publish_namespace_by_id_.end()) {
+  auto it = publish_namespace_by_id_.find(message.request_id);
+  if (it == publish_namespace_by_id_.end()) {
     return absl::OkStatus();  // State might have been destroyed due to
                               // PUBLISH_NAMESPACE_DONE.
   }
   std::move(it->second.cancel_callback)(MoqtRequestErrorInfo{
       message.error_code, std::nullopt, std::string(message.error_reason)});
-  session_->publish_namespace_by_namespace_.erase(it->second.track_namespace);
-  session_->publish_namespace_by_id_.erase(it);
+  publish_namespace_by_namespace_.erase(it->second.track_namespace);
+  publish_namespace_by_id_.erase(it);
   return absl::OkStatus();
 }
 
-absl::Status MoqtSession::ControlStream::OnControlMessage(
-    const MoqtTrackStatus& message) {
-  if (!session_->ValidateRequestId(message.request_id)) {
+absl::Status MoqtSession::OnControlMessage(const MoqtTrackStatus& message) {
+  if (!ValidateRequestId(message.request_id)) {
     return absl::OkStatus();
   }
-  if (session_->sent_goaway_) {
+  if (sent_goaway_) {
     QUIC_DLOG(INFO) << ENDPOINT
                     << "Received a TRACK_STATUS_REQUEST after GOAWAY";
-    return SendRequestError(message.request_id, RequestErrorCode::kUnauthorized,
-                            std::nullopt, "TRACK_STATUS_REQUEST after GOAWAY");
+    SendRequestErrorOnControlStream(
+        message.request_id, RequestErrorCode::kUnauthorized, std::nullopt,
+        "TRACK_STATUS_REQUEST after GOAWAY");
+    return absl::OkStatus();
   }
   // TODO(martinduke): Handle authentication.
   std::shared_ptr<MoqtTrackPublisher> track =
-      session_->publisher_->GetTrack(message.full_track_name);
+      publisher_->GetTrack(message.full_track_name);
   if (track == nullptr) {
-    return SendRequestError(message.request_id, RequestErrorCode::kDoesNotExist,
-                            std::nullopt, "Track does not exist");
+    SendRequestErrorOnControlStream(message.request_id,
+                                    RequestErrorCode::kDoesNotExist,
+                                    std::nullopt, "Track does not exist");
+    return absl::OkStatus();
   }
-  auto [it, inserted] = session_->incoming_track_status_.emplace(
+  auto [it, inserted] = incoming_track_status_.emplace(
       message.request_id, std::make_unique<DownstreamTrackStatus>(
-                              message.request_id, session_, track.get()));
+                              message.request_id, this, track.get()));
   track->AddObjectListener(it->second.get());
   return absl::OkStatus();
 }
 
-absl::Status MoqtSession::ControlStream::OnControlMessage(
-    const MoqtGoAway& message) {
+absl::Status MoqtSession::OnControlMessage(const MoqtGoAway& message) {
   if (!message.new_session_uri.empty() &&
       perspective() == quic::Perspective::IS_SERVER) {
     return absl::InvalidArgumentError(
         "Received GOAWAY with new_session_uri on the server");
   }
-  if (session_->received_goaway_) {
+  if (received_goaway_) {
     return absl::InvalidArgumentError("Received multiple GOAWAY messages");
   }
-  session_->received_goaway_ = true;
-  if (session_->callbacks_.goaway_received_callback != nullptr) {
-    std::move(session_->callbacks_.goaway_received_callback)(
-        message.new_session_uri);
+  received_goaway_ = true;
+  if (callbacks_.goaway_received_callback != nullptr) {
+    std::move(callbacks_.goaway_received_callback)(message.new_session_uri);
   }
   return absl::OkStatus();
 }
 
-absl::Status MoqtSession::ControlStream::OnControlMessage(
-    const MoqtMaxRequestId& message) {
-  if (message.max_request_id < session_->peer_max_request_id_) {
+absl::Status MoqtSession::OnControlMessage(const MoqtMaxRequestId& message) {
+  if (message.max_request_id < peer_max_request_id_) {
     QUIC_DLOG(INFO) << ENDPOINT
                     << "Peer sent MAX_REQUEST_ID message with "
                        "lower value than previous";
     return absl::InvalidArgumentError(
         "MAX_REQUEST_ID has lower value than previous");
   }
-  session_->peer_max_request_id_ = message.max_request_id;
+  peer_max_request_id_ = message.max_request_id;
   return absl::OkStatus();
 }
 
-absl::Status MoqtSession::ControlStream::OnControlMessage(
-    const MoqtFetch& message) {
-  if (!session_->ValidateRequestId(message.request_id)) {
+absl::Status MoqtSession::OnControlMessage(const MoqtFetch& message) {
+  if (!ValidateRequestId(message.request_id)) {
     return absl::OkStatus();
   }
-  if (session_->sent_goaway_) {
+  if (sent_goaway_) {
     QUIC_DLOG(INFO) << ENDPOINT << "Received a FETCH after GOAWAY";
-    return SendRequestError(message.request_id, RequestErrorCode::kUnauthorized,
-                            std::nullopt, "FETCH after GOAWAY");
+    SendRequestErrorOnControlStream(message.request_id,
+                                    RequestErrorCode::kUnauthorized,
+                                    std::nullopt, "FETCH after GOAWAY");
+    return absl::OkStatus();
   }
   std::unique_ptr<MoqtFetchTask> fetch;
   FullTrackName track_name;
@@ -1424,13 +1421,14 @@
         std::get<StandaloneFetch>(message.fetch);
     track_name = standalone_fetch.full_track_name;
     std::shared_ptr<MoqtTrackPublisher> track_publisher =
-        session_->publisher_->GetTrack(track_name);
+        publisher_->GetTrack(track_name);
     if (track_publisher == nullptr) {
       QUIC_DLOG(INFO) << ENDPOINT << "FETCH for " << track_name
                       << " rejected by the application: not found";
-      return SendRequestError(message.request_id,
-                              RequestErrorCode::kDoesNotExist, std::nullopt,
-                              "not found");
+      SendRequestErrorOnControlStream(message.request_id,
+                                      RequestErrorCode::kDoesNotExist,
+                                      std::nullopt, "not found");
+      return absl::OkStatus();
     }
     QUIC_DLOG(INFO) << ENDPOINT << "Received a StandaloneFETCH for "
                     << track_name;
@@ -1446,14 +1444,15 @@
             ? std::get<struct JoiningFetchRelative>(message.fetch)
                   .joining_request_id
             : std::get<JoiningFetchAbsolute>(message.fetch).joining_request_id;
-    auto it = session_->published_subscriptions_.find(joining_request_id);
-    if (it == session_->published_subscriptions_.end()) {
+    auto it = published_subscriptions_.find(joining_request_id);
+    if (it == published_subscriptions_.end()) {
       QUIC_DLOG(INFO) << ENDPOINT << "Received a JOINING_FETCH for "
                       << "request_id " << joining_request_id
                       << " that does not exist";
-      return SendRequestError(
+      SendRequestErrorOnControlStream(
           message.request_id, RequestErrorCode::kInvalidJoiningRequestId,
           std::nullopt, "Joining Fetch for non-existent request");
+      return absl::OkStatus();
     }
     if (!it->second->can_have_joining_fetch()) {
       QUIC_DLOG(INFO) << ENDPOINT << "Received a JOINING_FETCH for "
@@ -1466,9 +1465,10 @@
     if (it->second->established()) {
       if (!it->second->parameters().largest_object.has_value()) {
         // Nothing to Fetch.
-        return SendRequestError(message.request_id,
-                                RequestErrorCode::kDoesNotExist, std::nullopt,
-                                "not found");
+        SendRequestErrorOnControlStream(message.request_id,
+                                        RequestErrorCode::kDoesNotExist,
+                                        std::nullopt, "not found");
+        return absl::OkStatus();
       }
       const Location largest_location =
           *it->second->parameters().largest_object;
@@ -1485,9 +1485,10 @@
             std::get<JoiningFetchAbsolute>(message.fetch);
         start_group = absolute_fetch.joining_start;
         if (start_group > largest_location.group) {
-          return SendRequestError(message.request_id,
-                                  RequestErrorCode::kInvalidRange, std::nullopt,
-                                  "invalid range");
+          SendRequestErrorOnControlStream(message.request_id,
+                                          RequestErrorCode::kInvalidRange,
+                                          std::nullopt, "invalid range");
+          return absl::OkStatus();
         }
       }
       fetch = it->second->publisher().StandaloneFetch(
@@ -1512,38 +1513,40 @@
   if (!fetch->GetStatus().ok()) {
     QUIC_DLOG(INFO) << ENDPOINT << "FETCH for " << track_name
                     << " could not initialize the task";
-    return SendRequestError(message.request_id, RequestErrorCode::kInvalidRange,
-                            std::nullopt, fetch->GetStatus().message());
+    SendRequestErrorOnControlStream(message.request_id,
+                                    RequestErrorCode::kInvalidRange,
+                                    std::nullopt, fetch->GetStatus().message());
+    return absl::OkStatus();
   }
   auto published_fetch =
       std::make_unique<PublishedFetch>(message.request_id, std::move(fetch));
-  auto result = session_->incoming_fetches_.emplace(message.request_id,
-                                                    std::move(published_fetch));
+  auto result =
+      incoming_fetches_.emplace(message.request_id, std::move(published_fetch));
   if (!result.second) {  // Emplace failed.
     QUIC_DLOG(INFO) << ENDPOINT << "FETCH for " << track_name
                     << " could not be added to the session";
-    return SendRequestError(message.request_id,
-                            RequestErrorCode::kInternalError, std::nullopt,
-                            "Could not initialize FETCH state");
+    SendRequestErrorOnControlStream(
+        message.request_id, RequestErrorCode::kInternalError, std::nullopt,
+        "Could not initialize FETCH state");
+    return absl::OkStatus();
   }
   MoqtFetchTask* fetch_task = result.first->second->fetch_task_ptr();
   fetch_task->SetFetchResponseCallback(
       [this, request_id = message.request_id](
           std::variant<MoqtFetchOk, MoqtRequestError> message) {
-        if (!session_->incoming_fetches_.contains(request_id)) {
+        if (!incoming_fetches_.contains(request_id)) {
           return;  // FETCH was cancelled.
         }
         if (std::holds_alternative<MoqtFetchOk>(message)) {
           MoqtFetchOk& fetch_ok = std::get<MoqtFetchOk>(message);
           fetch_ok.request_id = request_id;
-          CheckStatus(SendOrBufferMessage(
-              session_->framer_.SerializeFetchOk(fetch_ok)));
+          SendControlMessage(framer_.SerializeFetchOk(fetch_ok));
           return;
         }
-        CheckStatus(SendRequestError(
+        SendRequestErrorOnControlStream(
             request_id, std::get<MoqtRequestError>(message).error_code,
             std::get<MoqtRequestError>(message).retry_interval,
-            std::get<MoqtRequestError>(message).reason_phrase));
+            std::get<MoqtRequestError>(message).reason_phrase);
       });
   // Set a temporary new-object callback that creates a data stream. When
   // created, the stream visitor will replace this callback.
@@ -1552,25 +1555,23 @@
        subscriber_priority = message.parameters.subscriber_priority.value_or(
            kDefaultSubscriberPriority),
        request_id = message.request_id]() {
-        auto it = session_->incoming_fetches_.find(request_id);
-        if (it == session_->incoming_fetches_.end()) {
+        auto it = incoming_fetches_.find(request_id);
+        if (it == incoming_fetches_.end()) {
           return;
         }
-        if (!session_->session()->CanOpenNextOutgoingUnidirectionalStream() ||
-            !session_->OpenDataStream(it->second.get(),
-                                      SendOrderForFetch(subscriber_priority))) {
-          session_->UpdateTrackPriority(
-              request_id, std::nullopt,
-              MoqtTrackPriority(subscriber_priority,
-                                kDefaultPublisherPriority));
+        if (!session()->CanOpenNextOutgoingUnidirectionalStream() ||
+            !OpenDataStream(it->second.get(),
+                            SendOrderForFetch(subscriber_priority))) {
+          UpdateTrackPriority(request_id, std::nullopt,
+                              MoqtTrackPriority(subscriber_priority,
+                                                kDefaultPublisherPriority));
         }
       });
   return absl::OkStatus();
 }
 
-absl::Status MoqtSession::ControlStream::OnControlMessage(
-    const MoqtFetchOk& message) {
-  RemoteTrack* track = session_->RemoteTrackById(message.request_id);
+absl::Status MoqtSession::OnControlMessage(const MoqtFetchOk& message) {
+  RemoteTrack* track = RemoteTrackById(message.request_id);
   if (track == nullptr) {
     QUIC_DLOG(INFO) << ENDPOINT << "Received the FETCH_OK for "
                     << "request_id = " << message.request_id
@@ -1584,31 +1585,28 @@
   QUIC_DLOG(INFO) << ENDPOINT << "Received the FETCH_OK for request_id = "
                   << message.request_id << " " << track->full_track_name();
   UpstreamFetch* fetch = absl::down_cast<UpstreamFetch*>(track);
-  fetch->OnFetchResult(
-      message.end_location, absl::OkStatus(),
-      [=, session = session_]() { session->CancelFetch(message.request_id); });
+  fetch->OnFetchResult(message.end_location, absl::OkStatus(),
+                       [=, this]() { CancelFetch(message.request_id); });
   return absl::OkStatus();
 }
 
-absl::Status MoqtSession::ControlStream::OnControlMessage(
-    const MoqtRequestsBlocked& message) {
+absl::Status MoqtSession::OnControlMessage(const MoqtRequestsBlocked& message) {
   // TODO(martinduke): Derive logic for granting more subscribes.
   return absl::OkStatus();
 }
 
-absl::Status MoqtSession::ControlStream::OnControlMessage(
-    const MoqtPublish& message) {
-  if (!session_->ValidateRequestId(message.request_id)) {
+absl::Status MoqtSession::OnControlMessage(const MoqtPublish& message) {
+  if (!ValidateRequestId(message.request_id)) {
     return absl::OkStatus();
   }
-  RequestErrorCode error_code = session_->sent_goaway_
-                                    ? RequestErrorCode::kUnauthorized
-                                    : RequestErrorCode::kNotSupported;
-  absl::string_view error_reason = session_->sent_goaway_
+  RequestErrorCode error_code = sent_goaway_ ? RequestErrorCode::kUnauthorized
+                                             : RequestErrorCode::kNotSupported;
+  absl::string_view error_reason = sent_goaway_
                                        ? "Received a PUBLISH after GOAWAY"
                                        : "PUBLISH is not supported";
-  return SendRequestError(message.request_id, error_code, std::nullopt,
-                          error_reason);
+  SendRequestErrorOnControlStream(message.request_id, error_code, std::nullopt,
+                                  error_reason);
+  return absl::OkStatus();
 }
 
 void MoqtSession::OnMalformedTrack(RemoteTrack* track) {
diff --git a/quiche/quic/moqt/moqt_session.h b/quiche/quic/moqt/moqt_session.h
index 81eb53c..b47935f 100644
--- a/quiche/quic/moqt/moqt_session.h
+++ b/quiche/quic/moqt/moqt_session.h
@@ -210,6 +210,7 @@
   MoqtTraceRecorder& trace_recorder() { return trace_recorder_; }
 
  private:
+  friend class ControlMessageDispatcher;
   friend class test::MoqtSessionPeer;
 
   struct Empty {};
@@ -266,36 +267,6 @@
         const MoqtRawControlMessage& message) override;
 
     // MoqtControlParserVisitor implementation.
-    absl::Status OnControlMessage(const MoqtSetup& message);
-    absl::Status OnControlMessage(const MoqtRequestOk& message);
-    absl::Status OnControlMessage(const MoqtRequestError& message);
-    absl::Status OnControlMessage(const MoqtSubscribe& message);
-    absl::Status OnControlMessage(const MoqtSubscribeOk& message);
-    absl::Status OnControlMessage(const MoqtUnsubscribe& message);
-    absl::Status OnControlMessage(const MoqtPublishDone& /*message*/);
-    absl::Status OnControlMessage(const MoqtRequestUpdate& message);
-    absl::Status OnControlMessage(const MoqtPublishNamespace& message);
-    absl::Status OnControlMessage(const MoqtPublishNamespaceDone& /*message*/);
-    absl::Status OnControlMessage(const MoqtPublishNamespaceCancel& message);
-    absl::Status OnControlMessage(const MoqtTrackStatus& message);
-    absl::Status OnControlMessage(const MoqtGoAway& /*message*/);
-    absl::Status OnControlMessage(const MoqtMaxRequestId& message);
-    absl::Status OnControlMessage(const MoqtFetch& message);
-    absl::Status OnControlMessage(const MoqtFetchCancel& /*message*/) {
-      return absl::OkStatus();
-    }
-    absl::Status OnControlMessage(const MoqtFetchOk& message);
-    absl::Status OnControlMessage(const MoqtRequestsBlocked& message);
-    absl::Status OnControlMessage(const MoqtPublish& message);
-    absl::Status OnControlMessage(const MoqtObjectAck& message) {
-      auto subscription_it =
-          session_->published_subscriptions_.find(message.subscribe_id);
-      if (subscription_it == session_->published_subscriptions_.end()) {
-        return absl::OkStatus();
-      }
-      subscription_it->second->ProcessObjectAck(message);
-      return absl::OkStatus();
-    }
 
     // webtransport::StreamVisitor overrides
     void OnResetStreamReceived(webtransport::StreamErrorCode error) override {
@@ -468,6 +439,51 @@
                                     parameters_.perspective);
   }
 
+  // Handlers for the control messages on the main control stream.
+  absl::Status OnControlMessage(const MoqtSetup& message);
+  absl::Status OnControlMessage(const MoqtRequestOk& message);
+  absl::Status OnControlMessage(const MoqtRequestError& message);
+  absl::Status OnControlMessage(const MoqtSubscribe& message);
+  absl::Status OnControlMessage(const MoqtSubscribeOk& message);
+  absl::Status OnControlMessage(const MoqtUnsubscribe& message);
+  absl::Status OnControlMessage(const MoqtPublishDone& /*message*/);
+  absl::Status OnControlMessage(const MoqtRequestUpdate& message);
+  absl::Status OnControlMessage(const MoqtPublishNamespace& message);
+  absl::Status OnControlMessage(const MoqtPublishNamespaceDone& /*message*/);
+  absl::Status OnControlMessage(const MoqtPublishNamespaceCancel& message);
+  absl::Status OnControlMessage(const MoqtTrackStatus& message);
+  absl::Status OnControlMessage(const MoqtGoAway& /*message*/);
+  absl::Status OnControlMessage(const MoqtMaxRequestId& message);
+  absl::Status OnControlMessage(const MoqtFetch& message);
+  absl::Status OnControlMessage(const MoqtFetchCancel& /*message*/) {
+    return absl::OkStatus();
+  }
+  absl::Status OnControlMessage(const MoqtFetchOk& message);
+  absl::Status OnControlMessage(const MoqtRequestsBlocked& message);
+  absl::Status OnControlMessage(const MoqtPublish& message);
+  absl::Status OnControlMessage(const MoqtObjectAck& message) {
+    auto subscription_it = published_subscriptions_.find(message.subscribe_id);
+    if (subscription_it == published_subscriptions_.end()) {
+      return absl::OkStatus();
+    }
+    subscription_it->second->ProcessObjectAck(message);
+    return absl::OkStatus();
+  }
+
+  // TODO(vasilvv): remove this once all requests are moved into individual
+  // streams.
+  void SendRequestErrorOnControlStream(
+      uint64_t request_id, RequestErrorCode error_code,
+      std::optional<quic::QuicTimeDelta> retry_interval,
+      absl::string_view reason_phrase) {
+    MoqtRequestError request_error;
+    request_error.request_id = request_id;
+    request_error.error_code = error_code;
+    request_error.retry_interval = retry_interval;
+    request_error.reason_phrase = reason_phrase;
+    SendControlMessage(framer_.SerializeRequestError(request_error));
+  }
+
   bool is_closing_ = false;
   webtransport::Session* session_;
   MoqtSessionParameters parameters_;
diff --git a/quiche/quic/moqt/moqt_session_test.cc b/quiche/quic/moqt/moqt_session_test.cc
index 2b6bbb9..f704f8f 100644
--- a/quiche/quic/moqt/moqt_session_test.cc
+++ b/quiche/quic/moqt/moqt_session_test.cc
@@ -1176,7 +1176,8 @@
   subscribe_ok.request_id += 2;
   EXPECT_CALL(
       mock_session_,
-      CloseSession(static_cast<uint64_t>(MoqtError::kDuplicateTrackAlias), ""));
+      CloseSession(static_cast<uint64_t>(MoqtError::kDuplicateTrackAlias),
+                   "Duplicate track alias"));
   control_stream->ReceiveMessage(subscribe_ok);
 }