Implement a dedicated WebTransport-only dispatcher. This should be a functional no-op. The next step would be adding raw QUIC support to it. PiperOrigin-RevId: 852939580
diff --git a/build/source_list.bzl b/build/source_list.bzl index bde91f4..b48e6f6 100644 --- a/build/source_list.bzl +++ b/build/source_list.bzl
@@ -15,6 +15,7 @@ "common/capsule.h", "common/http/http_header_block.h", "common/http/http_header_storage.h", + "common/http/status_code_mapping.h", "common/internet_checksum.h", "common/lifetime_tracking.h", "common/masque/connect_ip_datagram_payload.h", @@ -272,6 +273,8 @@ "quic/core/http/quic_spdy_stream_body_manager.h", "quic/core/http/spdy_utils.h", "quic/core/http/web_transport_http3.h", + "quic/core/http/web_transport_only_dispatcher.h", + "quic/core/http/web_transport_only_server_session.h", "quic/core/http/web_transport_stream_adapter.h", "quic/core/legacy_quic_stream_id_manager.h", "quic/core/packet_number_indexed_queue.h", @@ -421,6 +424,7 @@ "common/capsule.cc", "common/http/http_header_block.cc", "common/http/http_header_storage.cc", + "common/http/status_code_mapping.cc", "common/internet_checksum.cc", "common/masque/connect_ip_datagram_payload.cc", "common/masque/connect_udp_datagram_payload.cc", @@ -621,6 +625,8 @@ "quic/core/http/quic_spdy_stream_body_manager.cc", "quic/core/http/spdy_utils.cc", "quic/core/http/web_transport_http3.cc", + "quic/core/http/web_transport_only_dispatcher.cc", + "quic/core/http/web_transport_only_server_session.cc", "quic/core/http/web_transport_stream_adapter.cc", "quic/core/legacy_quic_stream_id_manager.cc", "quic/core/qpack/new_qpack_blocking_manager.cc",
diff --git a/build/source_list.gni b/build/source_list.gni index 42c0a86..cd0fefc 100644 --- a/build/source_list.gni +++ b/build/source_list.gni
@@ -15,6 +15,7 @@ "src/quiche/common/capsule.h", "src/quiche/common/http/http_header_block.h", "src/quiche/common/http/http_header_storage.h", + "src/quiche/common/http/status_code_mapping.h", "src/quiche/common/internet_checksum.h", "src/quiche/common/lifetime_tracking.h", "src/quiche/common/masque/connect_ip_datagram_payload.h", @@ -272,6 +273,8 @@ "src/quiche/quic/core/http/quic_spdy_stream_body_manager.h", "src/quiche/quic/core/http/spdy_utils.h", "src/quiche/quic/core/http/web_transport_http3.h", + "src/quiche/quic/core/http/web_transport_only_dispatcher.h", + "src/quiche/quic/core/http/web_transport_only_server_session.h", "src/quiche/quic/core/http/web_transport_stream_adapter.h", "src/quiche/quic/core/legacy_quic_stream_id_manager.h", "src/quiche/quic/core/packet_number_indexed_queue.h", @@ -421,6 +424,7 @@ "src/quiche/common/capsule.cc", "src/quiche/common/http/http_header_block.cc", "src/quiche/common/http/http_header_storage.cc", + "src/quiche/common/http/status_code_mapping.cc", "src/quiche/common/internet_checksum.cc", "src/quiche/common/masque/connect_ip_datagram_payload.cc", "src/quiche/common/masque/connect_udp_datagram_payload.cc", @@ -621,6 +625,8 @@ "src/quiche/quic/core/http/quic_spdy_stream_body_manager.cc", "src/quiche/quic/core/http/spdy_utils.cc", "src/quiche/quic/core/http/web_transport_http3.cc", + "src/quiche/quic/core/http/web_transport_only_dispatcher.cc", + "src/quiche/quic/core/http/web_transport_only_server_session.cc", "src/quiche/quic/core/http/web_transport_stream_adapter.cc", "src/quiche/quic/core/legacy_quic_stream_id_manager.cc", "src/quiche/quic/core/qpack/new_qpack_blocking_manager.cc",
diff --git a/build/source_list.json b/build/source_list.json index da12550..8de1946 100644 --- a/build/source_list.json +++ b/build/source_list.json
@@ -14,6 +14,7 @@ "quiche/common/capsule.h", "quiche/common/http/http_header_block.h", "quiche/common/http/http_header_storage.h", + "quiche/common/http/status_code_mapping.h", "quiche/common/internet_checksum.h", "quiche/common/lifetime_tracking.h", "quiche/common/masque/connect_ip_datagram_payload.h", @@ -271,6 +272,8 @@ "quiche/quic/core/http/quic_spdy_stream_body_manager.h", "quiche/quic/core/http/spdy_utils.h", "quiche/quic/core/http/web_transport_http3.h", + "quiche/quic/core/http/web_transport_only_dispatcher.h", + "quiche/quic/core/http/web_transport_only_server_session.h", "quiche/quic/core/http/web_transport_stream_adapter.h", "quiche/quic/core/legacy_quic_stream_id_manager.h", "quiche/quic/core/packet_number_indexed_queue.h", @@ -420,6 +423,7 @@ "quiche/common/capsule.cc", "quiche/common/http/http_header_block.cc", "quiche/common/http/http_header_storage.cc", + "quiche/common/http/status_code_mapping.cc", "quiche/common/internet_checksum.cc", "quiche/common/masque/connect_ip_datagram_payload.cc", "quiche/common/masque/connect_udp_datagram_payload.cc", @@ -620,6 +624,8 @@ "quiche/quic/core/http/quic_spdy_stream_body_manager.cc", "quiche/quic/core/http/spdy_utils.cc", "quiche/quic/core/http/web_transport_http3.cc", + "quiche/quic/core/http/web_transport_only_dispatcher.cc", + "quiche/quic/core/http/web_transport_only_server_session.cc", "quiche/quic/core/http/web_transport_stream_adapter.cc", "quiche/quic/core/legacy_quic_stream_id_manager.cc", "quiche/quic/core/qpack/new_qpack_blocking_manager.cc",
diff --git a/quiche/common/http/status_code_mapping.cc b/quiche/common/http/status_code_mapping.cc new file mode 100644 index 0000000..c4ba636 --- /dev/null +++ b/quiche/common/http/status_code_mapping.cc
@@ -0,0 +1,41 @@ +// Copyright 2025 The Chromium Authors. All rights reserved. +// Use of this source code is governed by a BSD-style license that can be +// found in the LICENSE file. + +#include "quiche/common/http/status_code_mapping.h" + +#include "absl/status/status.h" + +namespace quiche { + +int StatusCodeAbslToHttp(absl::StatusCode code) { + switch (code) { + case absl::StatusCode::kOk: + return 200; // OK + case absl::StatusCode::kUnauthenticated: + return 401; // Unauthorized + case absl::StatusCode::kPermissionDenied: + return 403; // Forbidden + case absl::StatusCode::kNotFound: + return 404; // Not Found + case absl::StatusCode::kResourceExhausted: + return 429; // Too Many Requests + case absl::StatusCode::kUnavailable: + return 503; // Service Unavailable + + case absl::StatusCode::kOutOfRange: + case absl::StatusCode::kInvalidArgument: + case absl::StatusCode::kFailedPrecondition: + case absl::StatusCode::kAlreadyExists: + return 400; // Bad Request + + default: + return 500; // Internal Server Error + } +} + +int StatusToHttpStatusCode(const absl::Status& status) { + return StatusCodeAbslToHttp(status.code()); +} + +} // namespace quiche
diff --git a/quiche/common/http/status_code_mapping.h b/quiche/common/http/status_code_mapping.h new file mode 100644 index 0000000..470bebc --- /dev/null +++ b/quiche/common/http/status_code_mapping.h
@@ -0,0 +1,20 @@ +// Copyright 2025 The Chromium Authors. All rights reserved. +// Use of this source code is governed by a BSD-style license that can be +// found in the LICENSE file. + +#ifndef QUICHE_COMMON_HTTP_STATUS_CODE_MAPPING_H_ +#define QUICHE_COMMON_HTTP_STATUS_CODE_MAPPING_H_ + +#include "absl/status/status.h" + +namespace quiche { + +// Converts an absl::StatusCode to an HTTP status code. +int StatusCodeAbslToHttp(absl::StatusCode code); + +// Converts an absl::Status to an HTTP status code. +int StatusToHttpStatusCode(const absl::Status& status); + +} // namespace quiche + +#endif // QUICHE_COMMON_HTTP_STATUS_CODE_MAPPING_H_
diff --git a/quiche/quic/core/http/web_transport_only_dispatcher.cc b/quiche/quic/core/http/web_transport_only_dispatcher.cc new file mode 100644 index 0000000..9704db2 --- /dev/null +++ b/quiche/quic/core/http/web_transport_only_dispatcher.cc
@@ -0,0 +1,58 @@ +// Copyright 2025 The Chromium Authors. All rights reserved. +// Use of this source code is governed by a BSD-style license that can be +// found in the LICENSE file. + +#include "quiche/quic/core/http/web_transport_only_dispatcher.h" + +#include <memory> + +#include "absl/strings/string_view.h" +#include "absl/types/span.h" +#include "quiche/quic/core/connection_id_generator.h" +#include "quiche/quic/core/http/web_transport_only_server_session.h" +#include "quiche/quic/core/quic_connection.h" +#include "quiche/quic/core/quic_connection_id.h" +#include "quiche/quic/core/quic_session.h" +#include "quiche/quic/core/quic_types.h" +#include "quiche/quic/core/quic_versions.h" +#include "quiche/quic/platform/api/quic_socket_address.h" +#include "quiche/web_transport/web_transport.h" + +namespace quic { + +std::unique_ptr<QuicSession> WebTransportOnlyDispatcher::CreateQuicSession( + QuicConnectionId server_connection_id, + const QuicSocketAddress& self_address, + const QuicSocketAddress& peer_address, absl::string_view alpn, + const ParsedQuicVersion& version, const ParsedClientHello& parsed_chlo, + ConnectionIdGeneratorInterface& connection_id_generator) { + if (alpn.empty()) { + return nullptr; + } + + // TODO(vasilvv): handle ALPN. + // TODO(vasilvv): support raw QUIC. + + auto connection = std::make_unique<QuicConnection>( + server_connection_id, self_address, peer_address, helper(), + alarm_factory(), writer(), + /*owns_writer=*/false, Perspective::IS_SERVER, + ParsedQuicVersionVector{version}, connection_id_generator); + auto session = std::make_unique<WebTransportOnlyServerSession>( + config(), GetSupportedVersions(), connection.release(), this, + session_helper(), crypto_config(), compressed_certs_cache(), + QuicPriorityType::kWebTransport); + session->SetHandlerFactory( + [this](webtransport::Session* session, + const WebTransportIncomingRequestDetails& details) { + return parameters_.handler_factory(session, details); + }); + session->SetSubprotocolCallback( + [this](absl::Span<const absl::string_view> subprotocols) { + return parameters_.subprotocol_callback(subprotocols); + }); + session->Initialize(); + return session; +} + +} // namespace quic
diff --git a/quiche/quic/core/http/web_transport_only_dispatcher.h b/quiche/quic/core/http/web_transport_only_dispatcher.h new file mode 100644 index 0000000..cbc5358 --- /dev/null +++ b/quiche/quic/core/http/web_transport_only_dispatcher.h
@@ -0,0 +1,68 @@ +// Copyright 2025 The Chromium Authors. All rights reserved. +// Use of this source code is governed by a BSD-style license that can be +// found in the LICENSE file. + +#ifndef QUICHE_QUIC_CORE_HTTP_WEB_TRANSPORT_ONLY_DISPATCHER_H_ +#define QUICHE_QUIC_CORE_HTTP_WEB_TRANSPORT_ONLY_DISPATCHER_H_ + +#include <memory> + +#include "absl/base/nullability.h" +#include "absl/status/status.h" +#include "absl/strings/string_view.h" +#include "absl/types/span.h" +#include "quiche/quic/core/connection_id_generator.h" +#include "quiche/quic/core/http/web_transport_only_server_session.h" +#include "quiche/quic/core/quic_connection_id.h" +#include "quiche/quic/core/quic_dispatcher.h" +#include "quiche/quic/core/quic_session.h" +#include "quiche/quic/core/quic_types.h" +#include "quiche/quic/core/quic_versions.h" +#include "quiche/quic/platform/api/quic_socket_address.h" +#include "quiche/common/platform/api/quiche_export.h" +#include "quiche/web_transport/web_transport.h" + +namespace quic { + +// Parameters defining behavior of WebTransportOnlyDispatcher. +struct QUICHE_EXPORT WebTransportOnlyDispatcherParameters { + // Returns the application-specific handler for an incoming session. If a + // nullptr is returned, the session will not be created. + WebTransportHandlerFactoryCallback handler_factory = + +[](webtransport::Session* absl_nonnull, + const WebTransportIncomingRequestDetails&) { + return absl::UnimplementedError("Backend not configured"); + }; + + // Selects one of the provided subprotocols to be used for the incoming + // session. If the value returned is outside of the [0, subprotocols.size()) + // range, the negotiation is assumed to be unsuccessful. For raw QUIC, this is + // a fatal error. For WebTransport over HTTP/3, an empty subprotocol is + // provided to the resulting session. + WebTransportSelectSubprotocolCallback subprotocol_callback = + +[](absl::Span<const absl::string_view> /*subprotocols*/) { return -1; }; +}; + +// WebTransportOnlyDispatcher is a dedicated dispatcher for applications that +// are written against the webtransport::Session API. +class QUICHE_EXPORT WebTransportOnlyDispatcher : public QuicDispatcher { + public: + using QuicDispatcher::QuicDispatcher; + + WebTransportOnlyDispatcherParameters& parameters() { return parameters_; } + + // QuicDispatcher implementation. + std::unique_ptr<QuicSession> CreateQuicSession( + QuicConnectionId server_connection_id, + const QuicSocketAddress& self_address, + const QuicSocketAddress& peer_address, absl::string_view alpn, + const ParsedQuicVersion& version, const ParsedClientHello& parsed_chlo, + ConnectionIdGeneratorInterface& connection_id_generator) override; + + private: + WebTransportOnlyDispatcherParameters parameters_; +}; + +} // namespace quic + +#endif // QUICHE_QUIC_CORE_HTTP_WEB_TRANSPORT_ONLY_DISPATCHER_H_
diff --git a/quiche/quic/core/http/web_transport_only_server_session.cc b/quiche/quic/core/http/web_transport_only_server_session.cc new file mode 100644 index 0000000..25cdf91 --- /dev/null +++ b/quiche/quic/core/http/web_transport_only_server_session.cc
@@ -0,0 +1,226 @@ +// Copyright 2025 The Chromium Authors. All rights reserved. +// Use of this source code is governed by a BSD-style license that can be +// found in the LICENSE file. + +#include "quiche/quic/core/http/web_transport_only_server_session.h" + +#include <cstddef> +#include <cstdint> +#include <memory> +#include <string> +#include <utility> +#include <vector> + +#include "absl/container/fixed_array.h" +#include "absl/memory/memory.h" +#include "absl/status/statusor.h" +#include "absl/strings/str_cat.h" +#include "absl/strings/string_view.h" +#include "quiche/quic/core/crypto/quic_compressed_certs_cache.h" +#include "quiche/quic/core/http/http_frames.h" +#include "quiche/quic/core/http/quic_header_list.h" +#include "quiche/quic/core/http/quic_server_initiated_spdy_stream.h" +#include "quiche/quic/core/http/quic_server_session_base.h" +#include "quiche/quic/core/http/quic_spdy_server_stream_base.h" +#include "quiche/quic/core/http/quic_spdy_stream.h" +#include "quiche/quic/core/http/spdy_utils.h" +#include "quiche/quic/core/http/web_transport_http3.h" +#include "quiche/quic/core/quic_crypto_server_stream_base.h" +#include "quiche/quic/core/quic_error_codes.h" +#include "quiche/quic/core/quic_stream.h" +#include "quiche/quic/core/quic_types.h" +#include "quiche/quic/platform/api/quic_bug_tracker.h" +#include "quiche/quic/platform/api/quic_logging.h" +#include "quiche/common/http/http_header_block.h" +#include "quiche/common/http/status_code_mapping.h" +#include "quiche/common/platform/api/quiche_bug_tracker.h" +#include "quiche/common/platform/api/quiche_logging.h" +#include "quiche/web_transport/web_transport_headers.h" + +namespace quic { + +namespace { +constexpr QuicStreamCount kDefaultMaxStreamsAcceptedPerLoop = 5; +} + +WebTransportOnlyServerSession::~WebTransportOnlyServerSession() { + DeleteConnection(); +} + +void WebTransportOnlyServerSession::Initialize() { + QUICHE_DCHECK(handler_factory_ != nullptr); + set_max_streams_accepted_per_loop(kDefaultMaxStreamsAcceptedPerLoop); + QuicServerSessionBase::Initialize(); +} + +std::unique_ptr<QuicCryptoServerStreamBase> +WebTransportOnlyServerSession::CreateQuicCryptoServerStream( + const QuicCryptoServerConfig* crypto_config, + QuicCompressedCertsCache* compressed_certs_cache) { + return CreateCryptoServerStream(crypto_config, compressed_certs_cache, this, + stream_helper()); +} + +QuicSpdyStream* WebTransportOnlyServerSession::CreateIncomingStream( + QuicStreamId id) { + if (!ShouldCreateIncomingStream(id)) { + return nullptr; + } + + QuicSpdyStream* stream = new Stream(id, this, BIDIRECTIONAL); + ActivateStream(absl::WrapUnique(stream)); + return stream; +} + +QuicSpdyStream* WebTransportOnlyServerSession::CreateIncomingStream( + PendingStream* pending) { + QuicSpdyStream* stream = new Stream(pending, this); + ActivateStream(absl::WrapUnique(stream)); + return stream; +} + +QuicSpdyStream* +WebTransportOnlyServerSession::CreateOutgoingBidirectionalStream() { + if (!ShouldCreateOutgoingBidirectionalStream()) { + return nullptr; + } + + QuicServerInitiatedSpdyStream* stream = new QuicServerInitiatedSpdyStream( + GetNextOutgoingBidirectionalStreamId(), this, BIDIRECTIONAL); + ActivateStream(absl::WrapUnique(stream)); + return stream; +} + +QuicStream* WebTransportOnlyServerSession::ProcessBidirectionalPendingStream( + PendingStream* pending) { + return CreateIncomingStream(pending); +} + +bool WebTransportOnlyServerSession::OnSettingsFrame( + const SettingsFrame& frame) { + if (!QuicServerSessionBase::OnSettingsFrame(frame)) { + return false; + } + if (!SupportsWebTransport()) { + QUIC_DLOG(ERROR) + << "Refusing connection that does not support WebTransport"; + return false; + } + return true; +} + +void WebTransportOnlyServerSession::Stream::OnBodyAvailable() { + QUIC_BUG(WebTransportOnlyServerSession_OnBodyAvailable) + << "Received body on a WebTransportOnlyServerSession stream"; + OnUnrecoverableError( + QUIC_INTERNAL_ERROR, + "Received HTTP/3 body data within a WebTransport-only session"); +} + +void WebTransportOnlyServerSession::Stream::SendErrorResponse(int code) { + if (!reading_stopped()) { + StopReading(); + } + quiche::HttpHeaderBlock headers; + headers[":status"] = absl::StrCat(code); + WriteHeaders(std::move(headers), /*fin=*/true, /*ack_listener=*/nullptr); +} + +std::string WebTransportOnlyServerSession::Stream::SelectSubprotocolResponse( + absl::string_view client_header_value) { + // The specification requires that a malformed subprotocol header is + // ignored, as is the failure to negotiate the value. + absl::StatusOr<std::vector<std::string>> subprotocols_offered = + webtransport::ParseSubprotocolRequestHeader(client_header_value); + if (!subprotocols_offered.ok()) { + return ""; + } + + // The callback requires a span of string_views instead of strings. + absl::FixedArray<absl::string_view> subprotocol_views( + subprotocols_offered->size()); + for (size_t i = 0; i < subprotocol_views.size(); ++i) { + subprotocol_views[i] = (*subprotocols_offered)[i]; + } + + int selected = static_cast<WebTransportOnlyServerSession*>(session()) + ->subprotocol_callback_(subprotocol_views); + if (selected < 0 || selected >= subprotocol_views.size()) { + return ""; + } + + return std::string(subprotocol_views[selected]); +} + +void WebTransportOnlyServerSession::Stream::OnInitialHeadersComplete( + bool fin, size_t frame_len, const QuicHeaderList& header_list) { + QuicSpdyServerStreamBase::OnInitialHeadersComplete(fin, frame_len, + header_list); + if (write_side_closed()) { + return; + } + + WebTransportIncomingRequestDetails details; + int64_t content_length; + if (!SpdyUtils::CopyAndValidateHeaders(header_list, &content_length, + &details.headers)) { + SendErrorResponse(400); + return; + } + ConsumeHeaderList(); + + if (web_transport() == nullptr) { + // Return 405 Method Not Allowed if any of the QuicSpdySession-level checks + // have prevented creation of the WebTransport session. + SendErrorResponse(405); + return; + } + + auto subprotocol_request_it = + details.headers.find(webtransport::kSubprotocolRequestHeader); + if (subprotocol_request_it != details.headers.end()) { + details.subprotocol = + SelectSubprotocolResponse(subprotocol_request_it->second); + // The specification requires that a malformed subprotocol header is + // ignored. + } + + WebTransportHandlerFactoryCallback& factory = + static_cast<WebTransportOnlyServerSession*>(session())->handler_factory_; + absl::StatusOr<WebTransportConnectResponse> response = + factory(web_transport(), details); + if (!response.ok()) { + SendErrorResponse(quiche::StatusToHttpStatusCode(response.status())); + return; + } + if (response->visitor == nullptr) { + QUICHE_BUG(WebTransportOnlyServerSession_null_visitor) + << "WebTransport request callback has returned a non-error response, " + "but a null visitor."; + SendErrorResponse(500); + return; + } + + response->headers[":status"] = "200"; + if (!details.subprotocol.empty()) { + absl::StatusOr<std::string> serialized = + webtransport::SerializeSubprotocolResponseHeader(details.subprotocol); + if (serialized.ok()) { + response->headers[webtransport::kSubprotocolResponseHeader] = + *std::move(serialized); + } else { + QUIC_DLOG(WARNING) << "Response has invalid subprotocol listed: " + << details.subprotocol; + } + } + + WriteHeaders(std::move(response->headers), /*fin=*/false, + /*ack_listener=*/nullptr); + web_transport()->SetVisitor(std::move(response->visitor)); + // This should be the last call in the sequence, as it will trigger + // OnSessionReady() and the related application logic. + // TODO(vasilvv): add a test ensuring this is called. + web_transport()->HeadersReceived(details.headers); +} + +} // namespace quic
diff --git a/quiche/quic/core/http/web_transport_only_server_session.h b/quiche/quic/core/http/web_transport_only_server_session.h new file mode 100644 index 0000000..b39779b --- /dev/null +++ b/quiche/quic/core/http/web_transport_only_server_session.h
@@ -0,0 +1,118 @@ +// Copyright 2025 The Chromium Authors. All rights reserved. +// Use of this source code is governed by a BSD-style license that can be +// found in the LICENSE file. + +#ifndef QUICHE_QUIC_CORE_HTTP_WEB_TRANSPORT_ONLY_SERVER_SESSION_H_ +#define QUICHE_QUIC_CORE_HTTP_WEB_TRANSPORT_ONLY_SERVER_SESSION_H_ + +#include <cstddef> +#include <memory> +#include <string> +#include <utility> + +#include "absl/base/nullability.h" +#include "absl/status/statusor.h" +#include "absl/strings/string_view.h" +#include "absl/types/span.h" +#include "quiche/quic/core/crypto/quic_compressed_certs_cache.h" +#include "quiche/quic/core/http/http_frames.h" +#include "quiche/quic/core/http/quic_header_list.h" +#include "quiche/quic/core/http/quic_server_session_base.h" +#include "quiche/quic/core/http/quic_spdy_server_stream_base.h" +#include "quiche/quic/core/http/quic_spdy_session.h" +#include "quiche/quic/core/http/quic_spdy_stream.h" +#include "quiche/quic/core/quic_crypto_server_stream_base.h" +#include "quiche/quic/core/quic_stream.h" +#include "quiche/quic/core/quic_types.h" +#include "quiche/common/http/http_header_block.h" +#include "quiche/common/platform/api/quiche_export.h" +#include "quiche/common/quiche_callbacks.h" +#include "quiche/web_transport/web_transport.h" + +namespace quic { + +// Details about an incoming WebTransport request that are provided to the +// application. +struct QUICHE_EXPORT WebTransportIncomingRequestDetails { + // HTTP request headers received for the WebTransport CONNECT request. + quiche::HttpHeaderBlock headers; + // Subprotocol selected during the negotiation phase, if any. + std::string subprotocol; +}; + +// Application-provided response for the incoming WebTransport session. +struct QUICHE_EXPORT WebTransportConnectResponse { + // The visitor handling the WebTransport request. Must be non-null. + std::unique_ptr<webtransport::SessionVisitor> visitor; + // Additional response headers added to the WebTransport CONNECT response. + quiche::HttpHeaderBlock headers; +}; + +using WebTransportSelectSubprotocolCallback = + quiche::MultiUseCallback<int(absl::Span<const absl::string_view>)>; + +// Important note: the callback MUST NOT access `session` in any way other than +// storing it, as it may be not fully initialized until later. +using WebTransportHandlerFactoryCallback = + quiche::MultiUseCallback<absl::StatusOr<WebTransportConnectResponse>( + webtransport::Session* absl_nonnull session, + const WebTransportIncomingRequestDetails& details)>; + +// WebTransportOnlyServerSession is an HTTP/3 session that only accepts incoming +// WebTransport requests. All other requests are rejected with the HTTP error +// 405 Method Not Allowed. +// +// WebTransportOnlyServerSession takes ownership of all underlying connections. +class QUICHE_EXPORT WebTransportOnlyServerSession + : public QuicServerSessionBase { + public: + using QuicServerSessionBase::QuicServerSessionBase; + ~WebTransportOnlyServerSession() override; + + void SetHandlerFactory(WebTransportHandlerFactoryCallback factory) { + handler_factory_ = std::move(factory); + } + void SetSubprotocolCallback(WebTransportSelectSubprotocolCallback callback) { + subprotocol_callback_ = std::move(callback); + } + + void Initialize() override; + std::unique_ptr<QuicCryptoServerStreamBase> CreateQuicCryptoServerStream( + const QuicCryptoServerConfig* crypto_config, + QuicCompressedCertsCache* compressed_certs_cache) override; + bool OnSettingsFrame(const SettingsFrame& frame) override; + + QuicSpdyStream* CreateIncomingStream(QuicStreamId id) override; + QuicSpdyStream* CreateIncomingStream(PendingStream* pending) override; + QuicSpdyStream* CreateOutgoingBidirectionalStream() override; + QuicStream* ProcessBidirectionalPendingStream( + PendingStream* pending) override; + + WebTransportHttp3VersionSet LocallySupportedWebTransportVersions() + const override { + return kDefaultSupportedWebTransportVersions; + } + HttpDatagramSupport LocalHttpDatagramSupport() override { + return HttpDatagramSupport::kRfcAndDraft04; + } + + private: + class Stream : public QuicSpdyServerStreamBase { + public: + using QuicSpdyServerStreamBase::QuicSpdyServerStreamBase; + + void OnBodyAvailable() override; + void OnInitialHeadersComplete(bool fin, size_t frame_len, + const QuicHeaderList& header_list) override; + void SendErrorResponse(int code); + std::string SelectSubprotocolResponse( + absl::string_view client_header_value); + }; + + WebTransportHandlerFactoryCallback handler_factory_; + WebTransportSelectSubprotocolCallback subprotocol_callback_; +}; + +} // namespace quic + +#endif // QUICHE_QUIC_CORE_HTTP_WEB_TRANSPORT_ONLY_SERVER_SESSION_H_
diff --git a/quiche/quic/moqt/tools/moqt_end_to_end_test.cc b/quiche/quic/moqt/tools/moqt_end_to_end_test.cc index 0cb2910..49c60e0 100644 --- a/quiche/quic/moqt/tools/moqt_end_to_end_test.cc +++ b/quiche/quic/moqt/tools/moqt_end_to_end_test.cc
@@ -46,9 +46,9 @@ : server_(quic::test::crypto_test_utils::ProofSourceForTesting(), absl::bind_front(&MoqtEndToEndTest::ServerBackend, this)) { quic::QuicIpAddress host = quic::TestLoopback(); - bool success = server_.CreateUDPSocketAndListen( + absl::Status socket_status = server_.CreateUDPSocketAndListen( quic::QuicSocketAddress(host, /*port=*/0)); - QUICHE_CHECK(success); + QUICHE_CHECK_OK(socket_status); server_address_ = quic::QuicSocketAddress(host, server_.port()); event_loop_ = server_.event_loop(); }
diff --git a/quiche/quic/moqt/tools/moqt_ingestion_server_bin.cc b/quiche/quic/moqt/tools/moqt_ingestion_server_bin.cc index 2b22e92..dd7757e 100644 --- a/quiche/quic/moqt/tools/moqt_ingestion_server_bin.cc +++ b/quiche/quic/moqt/tools/moqt_ingestion_server_bin.cc
@@ -264,8 +264,10 @@ quiche::QuicheIpAddress bind_address; QUICHE_CHECK(bind_address.FromString( quiche::GetQuicheCommandLineFlag(FLAGS_bind_address))); - server.CreateUDPSocketAndListen(quic::QuicSocketAddress( - bind_address, quiche::GetQuicheCommandLineFlag(FLAGS_port))); + absl::Status socket_status = + server.CreateUDPSocketAndListen(quic::QuicSocketAddress( + bind_address, quiche::GetQuicheCommandLineFlag(FLAGS_port))); + QUICHE_CHECK_OK(socket_status); server.HandleEventsForever(); return 0;
diff --git a/quiche/quic/moqt/tools/moqt_relay.cc b/quiche/quic/moqt/tools/moqt_relay.cc index 0d0ea7e..3cda74c 100644 --- a/quiche/quic/moqt/tools/moqt_relay.cc +++ b/quiche/quic/moqt/tools/moqt_relay.cc
@@ -48,17 +48,17 @@ : client_event_loop_(client_event_loop), // TODO(martinduke): Extend MoqtServer so that partial objects can be // received. - server_(std::make_unique<MoqtServer>(std::move(proof_source), - [this](absl::string_view path) { - return IncomingSessionHandler( - path); - })) { + server_(std::make_unique<MoqtServer>( + std::move(proof_source), [this](absl::string_view path) { + return IncomingSessionHandler(path); + })) { quiche::QuicheIpAddress bind_ip_address; QUICHE_CHECK(bind_ip_address.FromString(bind_address)); // CreateUDPSocketAndListen() creates the event loop that we will pass to // MoqtClient. - server_->CreateUDPSocketAndListen( + absl::Status socket_status = server_->CreateUDPSocketAndListen( quic::QuicSocketAddress(bind_ip_address, bind_port)); + QUICHE_CHECK_OK(socket_status); if (!default_upstream.empty()) { quic::QuicUrl url(default_upstream, "https"); if (client_event_loop == nullptr) {
diff --git a/quiche/quic/moqt/tools/moqt_server.cc b/quiche/quic/moqt/tools/moqt_server.cc index daacabe..ffe3b85 100644 --- a/quiche/quic/moqt/tools/moqt_server.cc +++ b/quiche/quic/moqt/tools/moqt_server.cc
@@ -4,51 +4,108 @@ #include "quiche/quic/moqt/tools/moqt_server.h" +#include <cstddef> #include <memory> +#include <string> #include <utility> +#include "absl/status/status.h" #include "absl/status/statusor.h" #include "absl/strings/string_view.h" #include "quiche/quic/core/crypto/proof_source.h" #include "quiche/quic/core/crypto/quic_crypto_server_config.h" +#include "quiche/quic/core/crypto/quic_random.h" +#include "quiche/quic/core/http/web_transport_only_server_session.h" +#include "quiche/quic/core/io/quic_default_event_loop.h" +#include "quiche/quic/core/io/quic_event_loop.h" +#include "quiche/quic/core/io/quic_server_io_harness.h" #include "quiche/quic/core/quic_connection_id.h" +#include "quiche/quic/core/quic_default_clock.h" +#include "quiche/quic/core/quic_default_connection_helper.h" +#include "quiche/quic/core/quic_time.h" #include "quiche/quic/core/quic_types.h" #include "quiche/quic/core/quic_versions.h" #include "quiche/quic/moqt/moqt_messages.h" #include "quiche/quic/moqt/moqt_quic_config.h" #include "quiche/quic/moqt/moqt_session.h" -#include "quiche/quic/tools/quic_server.h" -#include "quiche/quic/tools/web_transport_only_backend.h" +#include "quiche/quic/platform/api/quic_socket_address.h" +#include "quiche/quic/tools/quic_simple_crypto_server_stream_helper.h" +#include "quiche/common/quiche_status_utils.h" #include "quiche/web_transport/web_transport.h" namespace moqt { namespace { -quic::WebTransportRequestCallback CreateWebTransportCallback( - MoqtIncomingSessionCallback callback, quic::QuicServer* server) { - return [server = server, callback = std::move(callback)]( - absl::string_view path, webtransport::Session* session) - -> absl::StatusOr<std::unique_ptr<webtransport::SessionVisitor>> { - absl::StatusOr<MoqtConfigureSessionCallback> configurator = callback(path); + +std::string GenerateRandomTokenSecret() { + constexpr size_t kSize = 256 / 8; + char secret[kSize]; + quic::QuicRandom::GetInstance()->RandBytes(secret, sizeof(secret)); + return std::string(secret, sizeof(secret)); +} + +quic::WebTransportHandlerFactoryCallback CreateWebTransportCallback( + MoqtIncomingSessionCallback callback, quic::QuicEventLoop* event_loop) { + return [event_loop = event_loop, callback = std::move(callback)]( + webtransport::Session* session, + const quic::WebTransportIncomingRequestDetails& details) + -> absl::StatusOr<quic::WebTransportConnectResponse> { + auto path_it = details.headers.find(":path"); + absl::StatusOr<MoqtConfigureSessionCallback> configurator = + callback(path_it != details.headers.end() ? path_it->second : ""); if (!configurator.ok()) { return configurator.status(); } + MoqtSessionParameters parameters(quic::Perspective::IS_SERVER); auto moqt_session = std::make_unique<MoqtSession>( - session, parameters, server->event_loop()->CreateAlarmFactory()); + session, parameters, event_loop->CreateAlarmFactory()); std::move (*configurator)(moqt_session.get()); - return moqt_session; + + quic::WebTransportConnectResponse response; + response.visitor = std::move(moqt_session); + return response; }; } } // namespace MoqtServer::MoqtServer(std::unique_ptr<quic::ProofSource> proof_source, MoqtIncomingSessionCallback callback) - : backend_(CreateWebTransportCallback(std::move(callback), &server_)), - server_(std::move(proof_source), /*proof_verifier=*/nullptr, - GenerateQuicConfig(), - quic::QuicCryptoServerConfig::ConfigOptions(), - quic::CurrentSupportedVersionsWithTls(), &backend_, - quic::kQuicDefaultConnectionIdLength) {} + : config_(GenerateQuicConfig()), + crypto_config_(GenerateRandomTokenSecret(), + quic::QuicRandom::GetInstance(), std::move(proof_source), + quic::KeyExchangeSource::Default()), + version_manager_(quic::CurrentSupportedVersionsWithTls()), + connection_id_generator_(quic::kQuicDefaultConnectionIdLength), + event_loop_( + quic::GetDefaultEventLoop()->Create(quic::QuicDefaultClock::Get())), + dispatcher_(&config_, &crypto_config_, &version_manager_, + std::make_unique<quic::QuicDefaultConnectionHelper>(), + std::make_unique<quic::QuicSimpleCryptoServerStreamHelper>(), + event_loop_->CreateAlarmFactory(), + quic::kQuicDefaultConnectionIdLength, + connection_id_generator_) { + dispatcher_.parameters().handler_factory = + CreateWebTransportCallback(std::move(callback), event_loop_.get()); +} + +absl::Status MoqtServer::CreateUDPSocketAndListen( + const quic::QuicSocketAddress& address) { + QUICHE_ASSIGN_OR_RETURN(fd_, quic::CreateAndBindServerSocket(address)); + QUICHE_ASSIGN_OR_RETURN(io_, quic::QuicServerIoHarness::Create( + event_loop_.get(), &dispatcher_, *fd_)); + io_->InitializeWriter(); + return absl::OkStatus(); +} + +void MoqtServer::WaitForEvents() { + event_loop_->RunEventLoopOnce(quic::QuicTime::Delta::FromMilliseconds(50)); +} + +void MoqtServer::HandleEventsForever() { + while (true) { + WaitForEvents(); + } +} } // namespace moqt
diff --git a/quiche/quic/moqt/tools/moqt_server.h b/quiche/quic/moqt/tools/moqt_server.h index a4996f6..a59e489 100644 --- a/quiche/quic/moqt/tools/moqt_server.h +++ b/quiche/quic/moqt/tools/moqt_server.h
@@ -8,14 +8,20 @@ #include <memory> +#include "absl/status/status.h" #include "absl/status/statusor.h" #include "absl/strings/string_view.h" #include "quiche/quic/core/crypto/proof_source.h" +#include "quiche/quic/core/crypto/quic_crypto_server_config.h" +#include "quiche/quic/core/deterministic_connection_id_generator.h" +#include "quiche/quic/core/http/web_transport_only_dispatcher.h" #include "quiche/quic/core/io/quic_event_loop.h" +#include "quiche/quic/core/io/quic_server_io_harness.h" +#include "quiche/quic/core/io/socket.h" +#include "quiche/quic/core/quic_config.h" +#include "quiche/quic/core/quic_version_manager.h" #include "quiche/quic/moqt/moqt_session.h" #include "quiche/quic/platform/api/quic_socket_address.h" -#include "quiche/quic/tools/quic_server.h" -#include "quiche/quic/tools/web_transport_only_backend.h" #include "quiche/common/platform/api/quiche_export.h" #include "quiche/common/quiche_callbacks.h" @@ -40,18 +46,28 @@ explicit MoqtServer(std::unique_ptr<quic::ProofSource> proof_source, MoqtIncomingSessionCallback callback); - bool CreateUDPSocketAndListen(const quic::QuicSocketAddress& address) { - return server_.CreateUDPSocketAndListen(address); - } - void WaitForEvents() { server_.WaitForEvents(); } - void HandleEventsForever() { server_.HandleEventsForever(); } - quic::QuicEventLoop* event_loop() { return server_.event_loop(); } - int port() { return server_.port(); } + MoqtServer(const MoqtServer&) = delete; + MoqtServer(MoqtServer&&) = delete; + MoqtServer& operator=(const MoqtServer&) = delete; + MoqtServer& operator=(MoqtServer&&) = delete; + + absl::Status CreateUDPSocketAndListen(const quic::QuicSocketAddress& address); + void WaitForEvents(); + void HandleEventsForever(); + quic::QuicEventLoop* event_loop() { return event_loop_.get(); } + int port() { return io_->local_address().port(); } private: friend class test::MoqtServerPeer; - quic::WebTransportOnlyBackend backend_; - quic::QuicServer server_; + quic::QuicConfig config_; + quic::QuicCryptoServerConfig crypto_config_; + quic::QuicVersionManager version_manager_; + quic::DeterministicConnectionIdGenerator connection_id_generator_; + std::unique_ptr<quic::QuicEventLoop> event_loop_; + quic::WebTransportOnlyDispatcher dispatcher_; + + quic::OwnedSocketFd fd_; + std::unique_ptr<quic::QuicServerIoHarness> io_; }; } // namespace moqt
diff --git a/quiche/quic/moqt/tools/moqt_server_test.cc b/quiche/quic/moqt/tools/moqt_server_test.cc index b654819..06fd560 100644 --- a/quiche/quic/moqt/tools/moqt_server_test.cc +++ b/quiche/quic/moqt/tools/moqt_server_test.cc
@@ -4,8 +4,13 @@ #include "quiche/quic/moqt/tools/moqt_server.h" +#include <utility> + +#include "absl/base/nullability.h" #include "absl/memory/memory.h" +#include "absl/status/statusor.h" #include "absl/strings/string_view.h" +#include "quiche/quic/core/http/web_transport_only_server_session.h" #include "quiche/quic/core/quic_alarm.h" #include "quiche/quic/core/quic_time.h" #include "quiche/quic/moqt/moqt_session.h" @@ -16,14 +21,18 @@ #include "quiche/quic/tools/web_transport_only_backend.h" #include "quiche/common/http/http_header_block.h" #include "quiche/common/quiche_ip_address.h" +#include "quiche/common/test_tools/quiche_test_utils.h" #include "quiche/web_transport/test_tools/mock_web_transport.h" +#include "quiche/web_transport/web_transport.h" namespace moqt::test { class MoqtServerPeer { public: - static quic::WebTransportOnlyBackend* backend(MoqtServer& server) { - return &server.backend_; + static absl::StatusOr<quic::WebTransportConnectResponse> CallHandlerFactory( + MoqtServer& server, webtransport::Session* session, + const quic::WebTransportIncomingRequestDetails& details) { + return server.dispatcher_.parameters().handler_factory(session, details); } }; @@ -42,7 +51,7 @@ quiche::QuicheIpAddress bind_address; bind_address.FromString("127.0.0.1"); // This will create an event loop that makes alarm factories. - EXPECT_TRUE(server_.CreateUDPSocketAndListen( + QUICHE_EXPECT_OK(server_.CreateUDPSocketAndListen( quic::QuicSocketAddress(bind_address, 0))); } @@ -55,9 +64,12 @@ TEST_F(MoqtServerTest, NewSessionHasAlarmFactory) { quiche::HttpHeaderBlock headers; headers.AppendValueOrAddHeader(":path", "/foo"); - quic::WebTransportOnlyBackend::WebTransportResponse response = - MoqtServerPeer::backend(server_)->ProcessWebTransportRequest( - headers, &mock_session_); + absl::StatusOr<quic::WebTransportConnectResponse> response = + MoqtServerPeer::CallHandlerFactory( + server_, &mock_session_, + quic::WebTransportIncomingRequestDetails{.headers = + std::move(headers)}); + QUICHE_EXPECT_OK(response.status()); ASSERT_NE(session_, nullptr); ASSERT_NE(MoqtSessionPeer::GetAlarmFactory(session_), nullptr); auto delegate = new MockAlarmDelegate();
diff --git a/quiche/quic/tools/web_transport_only_backend.cc b/quiche/quic/tools/web_transport_only_backend.cc index d908de3..0e5fd0e 100644 --- a/quiche/quic/tools/web_transport_only_backend.cc +++ b/quiche/quic/tools/web_transport_only_backend.cc
@@ -8,10 +8,11 @@ #include <string> #include <utility> -#include "absl/status/status.h" #include "absl/status/statusor.h" +#include "absl/strings/str_cat.h" #include "quiche/quic/tools/quic_backend_response.h" #include "quiche/common/http/http_header_block.h" +#include "quiche/common/http/status_code_mapping.h" #include "quiche/web_transport/web_transport.h" namespace quic { @@ -45,24 +46,14 @@ absl::StatusOr<std::unique_ptr<webtransport::SessionVisitor>> processed = callback_(path->second, session); - switch (processed.status().code()) { - case absl::StatusCode::kOk: - response.response_headers[":status"] = "200"; - response.visitor = *std::move(processed); - return response; - case absl::StatusCode::kNotFound: - response.response_headers[":status"] = "404"; - return response; - case absl::StatusCode::kInvalidArgument: - response.response_headers[":status"] = "400"; - return response; - case absl::StatusCode::kResourceExhausted: - response.response_headers[":status"] = "429"; - return response; - default: - response.response_headers[":status"] = "500"; - return response; + if (!processed.ok()) { + response.response_headers[":status"] = + absl::StrCat(quiche::StatusToHttpStatusCode(processed.status())); + return response; } + response.response_headers[":status"] = "200"; + response.visitor = *std::move(processed); + return response; } } // namespace quic