< Summary

Line coverage
0%
Covered lines: 0
Uncovered lines: 529
Coverable lines: 529
Total lines: 997
Line coverage: 0%
Branch coverage
0%
Covered branches: 0
Total branches: 208
Branch coverage: 0%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

File(s)

https://raw.githubusercontent.com/dotnet/runtime/811a7eabb75c42db53440e8ba3f60c07511cfd1f/src/libraries/System.Net.Http/src/System/Net/Http/SocketsHttpHandler/Http3Connection.cs

#LineLine coverage
 1// Licensed to the .NET Foundation under one or more agreements.
 2// The .NET Foundation licenses this file to you under the MIT license.
 3
 4using System.Collections.Generic;
 5using System.Diagnostics;
 6using System.Globalization;
 7using System.IO;
 8using System.Net.Http.Headers;
 9using System.Net.Http.Metrics;
 10using System.Net.Quic;
 11using System.Runtime.CompilerServices;
 12using System.Runtime.Versioning;
 13using System.Threading;
 14using System.Threading.Tasks;
 15
 16namespace System.Net.Http
 17{
 18    [SupportedOSPlatform("linux")]
 19    [SupportedOSPlatform("macos")]
 20    [SupportedOSPlatform("windows")]
 21    internal sealed class Http3Connection : HttpConnectionBase
 22    {
 23        private readonly HttpAuthority _authority;
 24        private readonly byte[]? _altUsedEncodedHeader;
 25        private QuicConnection? _connection;
 26        private Task? _connectionClosedTask;
 27
 28        // Keep a collection of requests around so we can process GOAWAY.
 029        private readonly Dictionary<QuicStream, Http3RequestStream> _activeRequests = new Dictionary<QuicStream, Http3Re
 30
 31        // Set when GOAWAY is being processed, when aborting, or when disposing.
 032        private long _firstRejectedStreamId = -1;
 33
 34        // Our control stream.
 35        private QuicStream? _clientControl;
 36        private Task? _sendSettingsTask;
 37
 38        // Server-advertised SETTINGS_MAX_FIELD_SECTION_SIZE
 39        // https://www.rfc-editor.org/rfc/rfc9114.html#section-7.2.4.1-2.2.1
 040        private uint _maxHeaderListSize = uint.MaxValue; // Defaults to infinite
 41
 42        // Once the server's streams are received, these are set to true. Further receipt of these streams results in a 
 43        private bool _haveServerControlStream;
 44        private bool _haveServerQpackDecodeStream;
 45        private bool _haveServerQpackEncodeStream;
 46
 47        // A connection-level error will abort any future operations.
 48        private Exception? _abortException;
 49
 050        public HttpAuthority Authority => _authority;
 051        public HttpConnectionPool Pool => _pool;
 052        public uint MaxHeaderListSize => _maxHeaderListSize;
 053        public byte[]? AltUsedEncodedHeaderBytes => _altUsedEncodedHeader;
 054        public Exception? AbortException => Volatile.Read(ref _abortException);
 055        private object SyncObj => _activeRequests;
 56
 57        private int _availableRequestStreamsCount;
 58        private TaskCompletionSource<bool>? _availableStreamsWaiter;
 59
 60        /// <summary>
 61        /// If true, we've received GOAWAY, are aborting due to a connection-level error, or are disposing due to pool l
 62        /// </summary>
 63        private bool ShuttingDown
 64        {
 65            get
 066            {
 067                Debug.Assert(Monitor.IsEntered(SyncObj));
 068                return _firstRejectedStreamId != -1;
 069            }
 70        }
 71
 72        public Http3Connection(HttpConnectionPool pool, HttpAuthority authority, bool includeAltUsedHeader)
 073            : base(pool, GetNextConnectionId())
 074        {
 075            _authority = authority;
 76
 077            if (includeAltUsedHeader)
 078            {
 079                bool altUsedDefaultPort = pool.Kind == HttpConnectionKind.Http && authority.Port == HttpConnectionPool.D
 080                string altUsedValue = altUsedDefaultPort ? authority.IdnHost : string.Create(CultureInfo.InvariantCultur
 081                _altUsedEncodedHeader = QPack.QPackEncoder.EncodeLiteralHeaderFieldWithoutNameReferenceToArray(KnownHead
 082            }
 83
 084            uint maxHeaderListSize = _pool._lastSeenHttp3MaxHeaderListSize;
 085            if (maxHeaderListSize > 0)
 086            {
 87                // Previous connections to the same host advertised a limit.
 88                // Use this as an initial value before we receive the SETTINGS frame.
 089                _maxHeaderListSize = maxHeaderListSize;
 090            }
 091        }
 92
 93        public void InitQuicConnection(QuicConnection connection, Activity? connectionSetupActivity, DnsEndPoint connect
 094        {
 95            // Report the exact DnsEndPoint used to establish the QUIC connection (Alt-Svc may point it at an authority
 96            // distinct from the pool's origin), consistent with the connection's RemoteEndPoint.
 097            MarkConnectionAsEstablished(connectionSetupActivity: connectionSetupActivity, remoteEndPoint: connection.Rem
 98
 099            _connection = connection;
 100
 101            // Avoid capturing the initial request's ExecutionContext for the entire lifetime of the new connection.
 0102            using (ExecutionContext.SuppressFlow())
 0103            {
 104                // Errors are observed via Abort().
 0105                _sendSettingsTask = SendSettingsAsync();
 106
 107                // This process is cleaned up when _connection is disposed, and errors are observed via Abort().
 0108                _ = AcceptStreamsAsync();
 0109            }
 0110        }
 111
 112        /// <summary>
 113        /// Starts shutting down the <see cref="Http3Connection"/>. Final cleanup will happen when there are no more act
 114        /// </summary>
 115        public override void Dispose()
 0116        {
 0117            lock (SyncObj)
 0118            {
 0119                if (_firstRejectedStreamId == -1)
 0120                {
 0121                    _firstRejectedStreamId = long.MaxValue;
 0122                    CheckForShutdown();
 0123                }
 0124            }
 0125        }
 126
 127        /// <summary>
 128        /// Called when shutting down, this checks for when shutdown is complete (no more active requests) and does actu
 129        /// </summary>
 130        /// <remarks>Requires <see cref="SyncObj"/> to be locked.</remarks>
 131        private void CheckForShutdown()
 0132        {
 0133            Debug.Assert(Monitor.IsEntered(SyncObj));
 0134            Debug.Assert(ShuttingDown);
 135
 0136            if (_activeRequests.Count != 0)
 0137            {
 0138                return;
 139            }
 140
 0141            if (_connection != null)
 0142            {
 143                // Close the QuicConnection in the background.
 144
 0145                _availableStreamsWaiter?.SetResult(false);
 0146                _availableStreamsWaiter = null;
 147
 0148                _connectionClosedTask ??= _connection.CloseAsync((long)Http3ErrorCode.NoError).AsTask();
 149
 0150                QuicConnection connection = _connection;
 0151                _connection = null;
 152
 0153                _ = _connectionClosedTask.ContinueWith(async closeTask =>
 0154                {
 0155                    if (closeTask.IsFaulted && NetEventSource.Log.IsEnabled())
 0156                    {
 0157                        Trace($"{nameof(QuicConnection)} failed to close: {closeTask.Exception!.InnerException}");
 0158                    }
 0159
 0160                    try
 0161                    {
 0162                        await connection.DisposeAsync().ConfigureAwait(false);
 0163                    }
 0164                    catch (Exception ex)
 0165                    {
 0166                        Trace($"{nameof(QuicConnection)} failed to dispose: {ex}");
 0167                    }
 0168
 0169                    if (_clientControl != null)
 0170                    {
 0171                        await _sendSettingsTask!.ConfigureAwait(false);
 0172                        await _clientControl.DisposeAsync().ConfigureAwait(false);
 0173                        _clientControl = null;
 0174                    }
 0175
 0176                }, CancellationToken.None, TaskContinuationOptions.ExecuteSynchronously, TaskScheduler.Default);
 177
 0178                MarkConnectionAsClosed();
 0179            }
 0180        }
 181
 182        /// <summary>
 183        /// When EnableMultipleHttp3Connections is false: always reserve a stream, return a bool indicating if the strea
 184        /// When EnableMultipleHttp3Connections is true: reserve a stream only if it's available meaning that the return
 185        /// </summary>
 186        public bool TryReserveStream()
 0187        {
 0188            bool singleConnection = !_pool.Settings.EnableMultipleHttp3Connections;
 189
 0190            lock (SyncObj)
 0191            {
 192                // For the single connection case, we allow the counter to go below zero.
 0193                Debug.Assert(singleConnection || _availableRequestStreamsCount >= 0);
 194
 0195                if (NetEventSource.Log.IsEnabled()) Trace($"_availableRequestStreamsCount = {_availableRequestStreamsCou
 196
 0197                bool streamAvailable = _availableRequestStreamsCount > 0;
 198
 199                // Do not let the counter to go below zero when EnableMultipleHttp3Connections is true.
 200                // This equivalent to an immediate ReleaseStream() for the case no stream is immediately available.
 0201                if (singleConnection || _availableRequestStreamsCount > 0)
 0202                {
 0203                    --_availableRequestStreamsCount;
 0204                }
 205
 0206                return streamAvailable;
 207            }
 0208        }
 209
 210        public void ReleaseStream()
 0211        {
 0212            lock (SyncObj)
 0213            {
 0214                Debug.Assert(!_pool.Settings.EnableMultipleHttp3Connections || _availableRequestStreamsCount >= 0);
 215
 0216                if (NetEventSource.Log.IsEnabled()) Trace($"_availableRequestStreamsCount = {_availableRequestStreamsCou
 0217                ++_availableRequestStreamsCount;
 218
 0219                _availableStreamsWaiter?.SetResult(!ShuttingDown);
 0220                _availableStreamsWaiter = null;
 0221            }
 0222        }
 223
 224        public void StreamCapacityCallback(QuicConnection connection, QuicStreamCapacityChangedArgs args)
 0225        {
 0226            Debug.Assert(_connection is null || connection == _connection);
 227
 0228            lock (SyncObj)
 0229            {
 0230                Debug.Assert(_availableStreamsWaiter is null || _availableRequestStreamsCount >= 0);
 231
 0232                if (NetEventSource.Log.IsEnabled()) Trace($"_availableRequestStreamsCount = {_availableRequestStreamsCou
 233
 234                // Since _availableStreamsWaiter is only used in the multi-connection case, when _availableRequestStream
 235                // we don't need to check the value of _availableRequestStreamsCount here.
 0236                _availableRequestStreamsCount += args.BidirectionalIncrement;
 0237                _availableStreamsWaiter?.SetResult(!ShuttingDown);
 0238                _availableStreamsWaiter = null;
 0239            }
 0240        }
 241
 242        public Task<bool> WaitForAvailableStreamsAsync()
 0243        {
 244            // In the single connection case, _availableStreamsWaiter notifications do not guarantee that _availableRequ
 0245            Debug.Assert(_pool.Settings.EnableMultipleHttp3Connections, "Calling WaitForAvailableStreamsAsync() is inval
 246
 0247            lock (SyncObj)
 0248            {
 0249                Debug.Assert(_availableRequestStreamsCount >= 0);
 250
 0251                if (ShuttingDown)
 0252                {
 0253                    return Task.FromResult(false);
 254                }
 0255                if (_availableRequestStreamsCount > 0)
 0256                {
 0257                    return Task.FromResult(true);
 258                }
 259
 0260                Debug.Assert(_availableStreamsWaiter is null);
 0261                _availableStreamsWaiter = new TaskCompletionSource<bool>(TaskCreationOptions.RunContinuationsAsynchronou
 0262                return _availableStreamsWaiter.Task;
 263            }
 0264        }
 265
 266        public async Task<HttpResponseMessage> SendAsync(HttpRequestMessage request, WaitForHttp3ConnectionActivity wait
 0267        {
 0268            request.ConnectionId = Id;
 269
 270            // Allocate an active request
 0271            QuicStream? quicStream = null;
 0272            Http3RequestStream? requestStream = null;
 273
 274            try
 0275            {
 0276                Exception? exception = null;
 277                try
 0278                {
 0279                    QuicConnection? conn = _connection;
 0280                    if (conn != null)
 0281                    {
 282                        // We found a connection in the pool, but it did not have available streams, OpenOutboundStreamA
 0283                        if (!waitForConnectionActivity.Started && !streamAvailable)
 0284                        {
 0285                            waitForConnectionActivity.Start();
 0286                        }
 287
 0288                        quicStream = await conn.OpenOutboundStreamAsync(QuicStreamType.Bidirectional, cancellationToken)
 289
 0290                        requestStream = new Http3RequestStream(request, this, quicStream);
 0291                        lock (SyncObj)
 0292                        {
 0293                            if (_activeRequests.Count == 0)
 0294                            {
 0295                                MarkConnectionAsNotIdle();
 0296                            }
 0297                            _activeRequests.Add(quicStream, requestStream);
 0298                        }
 0299                    }
 0300                }
 301                // Swallow any exceptions caused by the connection being closed locally or even disposed due to a race.
 302                // Since quicStream will stay `null`, the code below will throw appropriate exception to retry the reque
 0303                catch (ObjectDisposedException e)
 0304                {
 0305                    exception = e;
 0306                }
 0307                catch (QuicException e) when (e.QuicError != QuicError.OperationAborted)
 0308                {
 0309                    exception = e;
 0310                }
 311                finally
 0312                {
 0313                    waitForConnectionActivity.Stop(request, Pool, exception);
 0314                }
 315
 0316                if (quicStream == null)
 0317                {
 0318                    throw new HttpRequestException(HttpRequestError.Unknown, SR.net_http_request_aborted, null, RequestR
 319                }
 320
 0321                requestStream!.StreamId = quicStream.Id;
 322
 323                bool goAway;
 0324                lock (SyncObj)
 0325                {
 0326                    goAway = _firstRejectedStreamId != -1 && requestStream.StreamId >= _firstRejectedStreamId;
 0327                }
 328
 0329                if (goAway)
 0330                {
 0331                    throw new HttpRequestException(HttpRequestError.Unknown, SR.net_http_request_aborted, null, RequestR
 332                }
 333
 0334                waitForConnectionActivity.AssertActivityNotRunning();
 0335                if (ConnectionSetupActivity is not null) ConnectionSetupDistributedTracing.AddConnectionLinkToRequestAct
 0336                if (NetEventSource.Log.IsEnabled()) Trace($"Sending request: {request}");
 337
 0338                Task<HttpResponseMessage> responseTask = requestStream.SendAsync(cancellationToken);
 339
 340                // null out requestStream to avoid disposing in finally block. It is now in charge of disposing itself.
 0341                requestStream = null;
 342
 0343                return await responseTask.ConfigureAwait(false);
 344            }
 0345            catch (QuicException ex) when (ex.QuicError == QuicError.OperationAborted)
 346            {
 347                // This will happen if we aborted _connection somewhere and we have pending OpenOutboundStreamAsync call
 348                // note that _abortException may be null if we closed the connection in response to a GOAWAY frame
 349                throw new HttpRequestException(HttpRequestError.Unknown, SR.net_http_client_execution_error, _abortExcep
 350            }
 351            finally
 0352            {
 0353                if (requestStream is not null)
 0354                {
 0355                    await requestStream.DisposeAsync().ConfigureAwait(false);
 0356                }
 0357            }
 0358        }
 359
 360        /// <summary>
 361        /// Aborts the connection with an error.
 362        /// </summary>
 363        /// <remarks>
 364        /// Used for e.g. I/O or connection-level frame parsing errors.
 365        /// </remarks>
 366        internal Exception Abort(Exception abortException)
 0367        {
 368            // Only observe the first exception we get.
 0369            Exception? firstException = Interlocked.CompareExchange(ref _abortException, abortException, null);
 370
 0371            if (firstException != null)
 0372            {
 0373                if (NetEventSource.Log.IsEnabled() && !ReferenceEquals(firstException, abortException))
 0374                {
 375                    // Lost the race to set the field to another exception, so just trace this one.
 0376                    Trace($"{nameof(abortException)}=={abortException}");
 0377                }
 378
 0379                return firstException;
 380            }
 381
 382            // Stop sending requests to this connection.
 383            // Do not dispose the connection when invalidating as the rest of this method does exactly that:
 384            //   set up _firstRejectedStreamId, close the connection with proper error code and CheckForShutdown.
 0385            _pool.InvalidateHttp3Connection(this, dispose: false);
 386
 0387            long connectionResetErrorCode = (abortException as HttpProtocolException)?.ErrorCode ?? (long)Http3ErrorCode
 388
 0389            lock (SyncObj)
 0390            {
 391                // Set _firstRejectedStreamId != -1 to make ShuttingDown = true.
 392                // It's possible GOAWAY is already being processed, in which case this would already be != -1.
 0393                if (_firstRejectedStreamId == -1)
 0394                {
 0395                    _firstRejectedStreamId = long.MaxValue;
 0396                }
 397
 398                // Abort the connection. This will cause all of our streams to abort on their next I/O.
 0399                if (_connection != null && _connectionClosedTask == null)
 0400                {
 0401                    _connectionClosedTask = _connection.CloseAsync((long)connectionResetErrorCode).AsTask();
 0402                }
 403
 0404                CheckForShutdown();
 0405            }
 406
 0407            return abortException;
 0408        }
 409
 410        private void OnServerGoAway(long firstRejectedStreamId)
 0411        {
 0412            if (NetEventSource.Log.IsEnabled())
 0413            {
 0414                Trace($"GOAWAY received. First rejected stream ID = {firstRejectedStreamId}");
 0415            }
 416
 417            // Stop sending requests to this connection.
 418            // Do not dispose the connection when invalidating as the rest of this method does exactly that:
 419            //   set up _firstRejectedStreamId to the stream id from GO_AWAY frame and CheckForShutdown.
 0420            _pool.InvalidateHttp3Connection(this, dispose: false);
 421
 0422            var streamsToGoAway = new List<Http3RequestStream>();
 423
 0424            lock (SyncObj)
 0425            {
 0426                if (_firstRejectedStreamId != -1 && firstRejectedStreamId > _firstRejectedStreamId)
 0427                {
 428                    // Server can send multiple GOAWAY frames.
 429                    // Spec says a server MUST NOT increase the stream ID in subsequent GOAWAYs,
 430                    // but doesn't specify what client should do if that is violated. Ignore for now.
 0431                    if (NetEventSource.Log.IsEnabled())
 0432                    {
 0433                        Trace("HTTP/3 server sent GOAWAY with increasing stream ID. Retried requests may have been doubl
 0434                    }
 0435                    return;
 436                }
 437
 0438                _firstRejectedStreamId = firstRejectedStreamId;
 439
 0440                foreach (KeyValuePair<QuicStream, Http3RequestStream> request in _activeRequests)
 0441                {
 0442                    if (request.Value.StreamId >= firstRejectedStreamId)
 0443                    {
 0444                        streamsToGoAway.Add(request.Value);
 0445                    }
 0446                }
 447
 0448                CheckForShutdown();
 0449            }
 450
 451            // GOAWAY each stream outside of the lock, so they can acquire the lock to remove themselves from _activeReq
 0452            foreach (Http3RequestStream stream in streamsToGoAway)
 0453            {
 0454                stream.GoAway();
 0455            }
 0456        }
 457
 458        public void RemoveStream(QuicStream stream)
 0459        {
 0460            lock (SyncObj)
 0461            {
 0462                if (_activeRequests.Remove(stream))
 0463                {
 0464                    if (_activeRequests.Count == 0)
 0465                    {
 0466                        MarkConnectionAsIdle();
 0467                    }
 468
 0469                    if (ShuttingDown)
 0470                    {
 0471                        CheckForShutdown();
 0472                    }
 0473                }
 0474            }
 0475        }
 476
 477        public override void Trace(string message, [CallerMemberName] string? memberName = null) =>
 0478            Trace(0, _connection is not null ? $"{_connection} {message}" : message, memberName);
 479
 480        internal void Trace(long streamId, string message, [CallerMemberName] string? memberName = null) =>
 0481            NetEventSource.Log.HandlerMessage(
 0482                _pool?.GetHashCode() ?? 0,    // pool ID
 0483                GetHashCode(),                // connection ID
 0484                (int)streamId,                // stream ID
 0485                memberName,                   // method name
 0486                message);                     // message
 487
 488        private async Task SendSettingsAsync()
 0489        {
 490            try
 0491            {
 0492                _clientControl = await _connection!.OpenOutboundStreamAsync(QuicStreamType.Unidirectional).ConfigureAwai
 493
 494                // Server MUST NOT abort our control stream, setup a continuation which will react accordingly
 0495                _ = _clientControl.WritesClosed.ContinueWith(t =>
 0496                {
 0497                    if (t.Exception?.InnerException is QuicException ex && ex.QuicError == QuicError.StreamAborted)
 0498                    {
 0499                        Abort(HttpProtocolException.CreateHttp3ConnectionException(Http3ErrorCode.ClosedCriticalStream))
 0500                    }
 0501                }, CancellationToken.None, TaskContinuationOptions.ExecuteSynchronously, TaskScheduler.Current);
 502
 0503                await _clientControl.WriteAsync(_pool.Settings.Http3SettingsFrame, CancellationToken.None).ConfigureAwai
 0504            }
 0505            catch (QuicException ex) when (ex.QuicError == QuicError.ConnectionAborted)
 0506            {
 0507                Debug.Assert(ex.ApplicationErrorCode.HasValue);
 0508                Http3ErrorCode code = (Http3ErrorCode)ex.ApplicationErrorCode.Value;
 509
 0510                Abort(HttpProtocolException.CreateHttp3ConnectionException(code, SR.net_http_http3_connection_close));
 0511            }
 0512            catch (Exception ex)
 0513            {
 0514                Abort(ex);
 0515            }
 0516        }
 517
 518        public static unsafe byte[] BuildSettingsFrame(HttpConnectionSettings settings)
 0519        {
 0520            Span<byte> buffer = stackalloc byte[4 + VariableLengthIntegerHelper.MaximumEncodedLength];
 521
 0522            int integerLength = VariableLengthIntegerHelper.WriteInteger(buffer.Slice(4), settings.MaxResponseHeadersByt
 0523            int payloadLength = 1 + integerLength; // includes the setting ID and the integer value.
 0524            Debug.Assert(payloadLength <= VariableLengthIntegerHelper.OneByteLimit);
 525
 0526            buffer[0] = (byte)Http3StreamType.Control;
 0527            buffer[1] = (byte)Http3FrameType.Settings;
 0528            buffer[2] = (byte)payloadLength;
 0529            buffer[3] = (byte)Http3SettingType.MaxHeaderListSize;
 530
 0531            return buffer.Slice(0, 4 + integerLength).ToArray();
 0532        }
 533
 534        /// <summary>
 535        /// Accepts unidirectional streams (control, QPack, ...) from the server.
 536        /// </summary>
 537        private async Task AcceptStreamsAsync()
 0538        {
 539            try
 0540            {
 0541                while (true)
 0542                {
 543                    ValueTask<QuicStream> streamTask;
 544
 0545                    lock (SyncObj)
 0546                    {
 0547                        if (ShuttingDown)
 0548                        {
 0549                            return;
 550                        }
 551
 552                        // No cancellation token is needed here; we expect the operation to cancel itself when _connecti
 0553                        streamTask = _connection!.AcceptInboundStreamAsync(CancellationToken.None);
 0554                    }
 555
 0556                    QuicStream stream = await streamTask.ConfigureAwait(false);
 557
 558                    // This process is cleaned up when _connection is disposed, and errors are observed via Abort().
 0559                    _ = ProcessServerStreamAsync(stream);
 0560                }
 561            }
 0562            catch (QuicException ex) when (ex.QuicError == QuicError.OperationAborted)
 0563            {
 564                // Shutdown initiated by us, no need to abort.
 0565            }
 0566            catch (QuicException ex) when (ex.QuicError == QuicError.ConnectionAborted)
 0567            {
 0568                Debug.Assert(ex.ApplicationErrorCode.HasValue);
 0569                Http3ErrorCode code = (Http3ErrorCode)ex.ApplicationErrorCode.Value;
 570
 0571                Abort(HttpProtocolException.CreateHttp3ConnectionException(code, SR.net_http_http3_connection_close));
 0572            }
 0573            catch (Exception ex)
 0574            {
 0575                Abort(ex);
 0576            }
 0577        }
 578
 579        /// <summary>
 580        /// Routes a stream to an appropriate stream-type-specific processor
 581        /// </summary>
 582        private async Task ProcessServerStreamAsync(QuicStream stream)
 0583        {
 0584            ArrayBuffer buffer = default;
 585
 586            try
 0587            {
 0588                await using (stream.ConfigureAwait(false))
 0589                {
 590                    // Check if this is a bidirectional stream (which we don't support from the server).
 0591                    if (stream.CanWrite)
 0592                    {
 593                        // Server initiated bidirectional streams are either push streams or extensions, and we support 
 0594                        throw HttpProtocolException.CreateHttp3ConnectionException(Http3ErrorCode.StreamCreationError);
 595                    }
 596
 0597                    buffer = new ArrayBuffer(initialSize: 32, usePool: true);
 598
 599                    // Read the stream type, which is a variable-length integer.
 600                    // This may require multiple reads if the integer is encoded in multiple bytes.
 601                    long streamType;
 0602                    while (true)
 0603                    {
 604                        int bytesRead;
 605                        try
 0606                        {
 0607                            bytesRead = await stream.ReadAsync(buffer.AvailableMemory, CancellationToken.None).Configure
 0608                        }
 0609                        catch (QuicException ex) when (ex.QuicError == QuicError.StreamAborted)
 0610                        {
 611                            // Treat identical to receiving 0. See below comment.
 0612                            bytesRead = 0;
 0613                        }
 614
 0615                        if (bytesRead == 0)
 0616                        {
 617                            // https://www.rfc-editor.org/rfc/rfc9114.html#name-unidirectional-streams
 618                            // A sender can close or reset a unidirectional stream unless otherwise specified. A receive
 619                            // tolerate unidirectional streams being closed or reset prior to the reception of the unidi
 620                            // stream header.
 0621                            return;
 622                        }
 623
 0624                        buffer.Commit(bytesRead);
 625
 0626                        if (VariableLengthIntegerHelper.TryRead(buffer.ActiveSpan, out streamType, out int streamTypeLen
 0627                        {
 628                            // Successfully read the stream type.
 0629                            buffer.Discard(streamTypeLength);
 0630                            break;
 631                        }
 0632                    }
 633
 0634                    if (NetEventSource.Log.IsEnabled())
 0635                    {
 0636                        NetEventSource.Info(this, $"Received server-initiated unidirectional stream of type {streamType}
 0637                    }
 638
 639                    // Process the stream based on its type.
 0640                    switch ((Http3StreamType)streamType)
 641                    {
 642                        case Http3StreamType.Control:
 0643                            if (Interlocked.Exchange(ref _haveServerControlStream, true))
 0644                            {
 645                                // A second control stream has been received.
 0646                                throw HttpProtocolException.CreateHttp3ConnectionException(Http3ErrorCode.StreamCreation
 647                            }
 648
 649                            // Ownership of buffer is transferred to ProcessServerControlStreamAsync.
 0650                            ArrayBuffer bufferCopy = buffer;
 0651                            buffer = default;
 652
 0653                            await ProcessServerControlStreamAsync(stream, bufferCopy).ConfigureAwait(false);
 0654                            return;
 655                        case Http3StreamType.QPackDecoder:
 0656                            if (Interlocked.Exchange(ref _haveServerQpackDecodeStream, true))
 0657                            {
 658                                // A second QPack decode stream has been received.
 0659                                throw HttpProtocolException.CreateHttp3ConnectionException(Http3ErrorCode.StreamCreation
 660                            }
 661
 662                            // The stream must not be closed, but we aren't using QPACK right now -- ignore.
 0663                            buffer.Dispose();
 0664                            await stream.CopyToAsync(Stream.Null).ConfigureAwait(false);
 0665                            return;
 666                        case Http3StreamType.QPackEncoder:
 0667                            if (Interlocked.Exchange(ref _haveServerQpackEncodeStream, true))
 0668                            {
 669                                // A second QPack encode stream has been received.
 0670                                throw HttpProtocolException.CreateHttp3ConnectionException(Http3ErrorCode.StreamCreation
 671                            }
 672
 673                            // We haven't enabled QPack in our SETTINGS frame, so we shouldn't receive any meaningful da
 674                            // However, the standard says the stream must not be closed for the lifetime of the connecti
 0675                            buffer.Dispose();
 0676                            await stream.CopyToAsync(Stream.Null).ConfigureAwait(false);
 0677                            return;
 678                        case Http3StreamType.Push:
 679                            // We don't support push streams.
 680                            // Because no maximum push stream ID was negotiated via a MAX_PUSH_ID frame, server should n
 0681                            throw HttpProtocolException.CreateHttp3ConnectionException(Http3ErrorCode.IdError);
 682                        default:
 683                            // Unknown stream type. Per spec, these must be ignored and aborted but not be considered a 
 0684                            stream.Abort(QuicAbortDirection.Read, (long)Http3ErrorCode.StreamCreationError);
 0685                            return;
 686                    }
 687                }
 0688            }
 0689            catch (QuicException ex) when (ex.QuicError == QuicError.OperationAborted)
 690            {
 691                // ignore the exception, we have already closed the connection
 692            }
 0693            catch (QuicException ex) when (ex.QuicError == QuicError.ConnectionAborted)
 694            {
 695                Debug.Assert(ex.ApplicationErrorCode.HasValue);
 696                Http3ErrorCode code = (Http3ErrorCode)ex.ApplicationErrorCode.Value;
 697
 698                Abort(HttpProtocolException.CreateHttp3ConnectionException(code, SR.net_http_http3_connection_close));
 699            }
 0700            catch (Exception ex)
 0701            {
 0702                Abort(ex);
 0703            }
 704            finally
 0705            {
 0706                buffer.Dispose();
 0707            }
 0708        }
 709
 710        /// <summary>
 711        /// Reads the server's control stream.
 712        /// </summary>
 713        private async Task ProcessServerControlStreamAsync(QuicStream stream, ArrayBuffer buffer)
 0714        {
 715            try
 0716            {
 0717                using (buffer)
 0718                {
 719                    // Read the first frame of the control stream. Per spec:
 720                    // A SETTINGS frame MUST be sent as the first frame of each control stream.
 721
 0722                    (Http3FrameType? frameType, long payloadLength) = await ReadFrameEnvelopeAsync().ConfigureAwait(fals
 723
 0724                    if (frameType == null)
 0725                    {
 726                        // Connection closed prematurely, expected SETTINGS frame.
 0727                        throw HttpProtocolException.CreateHttp3ConnectionException(Http3ErrorCode.ClosedCriticalStream);
 728                    }
 729
 0730                    if (frameType != Http3FrameType.Settings)
 0731                    {
 0732                        throw HttpProtocolException.CreateHttp3ConnectionException(Http3ErrorCode.MissingSettings);
 733                    }
 734
 0735                    await ProcessSettingsFrameAsync(payloadLength).ConfigureAwait(false);
 736
 737                    // Read subsequent frames.
 738
 0739                    while (true)
 0740                    {
 0741                        (frameType, payloadLength) = await ReadFrameEnvelopeAsync().ConfigureAwait(false);
 742
 0743                        switch (frameType)
 744                        {
 745                            case Http3FrameType.GoAway:
 0746                                await ProcessGoAwayFrameAsync(payloadLength).ConfigureAwait(false);
 0747                                break;
 748                            case Http3FrameType.Settings:
 749                                // If an endpoint receives a second SETTINGS frame on the control stream, the endpoint M
 0750                                throw HttpProtocolException.CreateHttp3ConnectionException(Http3ErrorCode.UnexpectedFram
 751                            case Http3FrameType.Headers: // Servers should not send these frames to a control stream.
 752                            case Http3FrameType.Data:
 753                            case Http3FrameType.MaxPushId:
 754                            case Http3FrameType.ReservedHttp2Priority: // These frames are explicitly reserved and must 
 755                            case Http3FrameType.ReservedHttp2Ping:
 756                            case Http3FrameType.ReservedHttp2WindowUpdate:
 757                            case Http3FrameType.ReservedHttp2Continuation:
 0758                                if (NetEventSource.Log.IsEnabled())
 0759                                {
 0760                                    Trace($"Received reserved frame: {frameType}");
 0761                                }
 0762                                throw HttpProtocolException.CreateHttp3ConnectionException(Http3ErrorCode.UnexpectedFram
 763                            case Http3FrameType.PushPromise:
 764                            case Http3FrameType.CancelPush:
 765                                // Because we haven't sent any MAX_PUSH_ID frame, it is invalid to receive any push-rela
 0766                                throw HttpProtocolException.CreateHttp3ConnectionException(Http3ErrorCode.IdError);
 767                            case null:
 768                                // End of stream reached. If we're shutting down, stop looping. Otherwise, this is an er
 769                                bool shuttingDown;
 0770                                lock (SyncObj)
 0771                                {
 0772                                    shuttingDown = ShuttingDown;
 0773                                }
 0774                                if (!shuttingDown)
 0775                                {
 0776                                    if (NetEventSource.Log.IsEnabled())
 0777                                    {
 0778                                        Trace($"Control stream closed by the server.");
 0779                                    }
 0780                                    throw HttpProtocolException.CreateHttp3ConnectionException(Http3ErrorCode.ClosedCrit
 781                                }
 0782                                return;
 783                            default:
 0784                                await SkipUnknownPayloadAsync(payloadLength).ConfigureAwait(false);
 0785                                break;
 786                        }
 0787                    }
 788                }
 789            }
 0790            catch (QuicException ex) when (ex.QuicError == QuicError.StreamAborted)
 0791            {
 792                // Peers MUST NOT close the control stream
 0793                throw HttpProtocolException.CreateHttp3ConnectionException(Http3ErrorCode.ClosedCriticalStream);
 794            }
 795
 796            async ValueTask<(Http3FrameType? frameType, long payloadLength)> ReadFrameEnvelopeAsync()
 0797            {
 798                long frameType, payloadLength;
 799                int bytesRead;
 800
 0801                while (!Http3Frame.TryReadIntegerPair(buffer.ActiveSpan, out frameType, out payloadLength, out bytesRead
 0802                {
 0803                    buffer.EnsureAvailableSpace(VariableLengthIntegerHelper.MaximumEncodedLength * 2);
 0804                    bytesRead = await stream.ReadAsync(buffer.AvailableMemory, CancellationToken.None).ConfigureAwait(fa
 805
 0806                    if (bytesRead != 0)
 0807                    {
 0808                        buffer.Commit(bytesRead);
 0809                    }
 0810                    else if (buffer.ActiveLength == 0)
 0811                    {
 812                        // End of stream.
 0813                        return (null, 0);
 814                    }
 815                    else
 0816                    {
 817                        // Our buffer has partial frame data in it but not enough to complete the read: bail out.
 0818                        throw HttpProtocolException.CreateHttp3ConnectionException(Http3ErrorCode.FrameError);
 819                    }
 0820                }
 821
 0822                buffer.Discard(bytesRead);
 823
 0824                return ((Http3FrameType)frameType, payloadLength);
 0825            }
 826
 827            async ValueTask ProcessSettingsFrameAsync(long settingsPayloadLength)
 0828            {
 0829                while (settingsPayloadLength != 0)
 0830                {
 831                    long settingId, settingValue;
 832                    int bytesRead;
 833
 0834                    while (!Http3Frame.TryReadIntegerPair(buffer.ActiveSpan, out settingId, out settingValue, out bytesR
 0835                    {
 0836                        buffer.EnsureAvailableSpace(VariableLengthIntegerHelper.MaximumEncodedLength * 2);
 0837                        bytesRead = await stream.ReadAsync(buffer.AvailableMemory, CancellationToken.None).ConfigureAwai
 838
 0839                        if (bytesRead != 0)
 0840                        {
 0841                            buffer.Commit(bytesRead);
 0842                        }
 843                        else
 0844                        {
 845                            // Our buffer has partial frame data in it but not enough to complete the read: bail out.
 0846                            throw HttpProtocolException.CreateHttp3ConnectionException(Http3ErrorCode.FrameError);
 847                        }
 0848                    }
 849
 0850                    settingsPayloadLength -= bytesRead;
 851
 0852                    if (settingsPayloadLength < 0)
 0853                    {
 854                        // An integer was encoded past the payload length.
 855                        // A frame payload that contains additional bytes after the identified fields or a frame payload
 0856                        throw HttpProtocolException.CreateHttp3ConnectionException(Http3ErrorCode.FrameError);
 857                    }
 858
 0859                    buffer.Discard(bytesRead);
 860
 0861                    if (NetEventSource.Log.IsEnabled()) Trace($"Applying setting {(Http3SettingType)settingId}={settingV
 862
 0863                    switch ((Http3SettingType)settingId)
 864                    {
 865                        case Http3SettingType.MaxHeaderListSize:
 0866                            _maxHeaderListSize = (uint)Math.Min((ulong)settingValue, uint.MaxValue);
 0867                            _pool._lastSeenHttp3MaxHeaderListSize = _maxHeaderListSize;
 0868                            break;
 869                        case Http3SettingType.ReservedHttp2EnablePush:
 870                        case Http3SettingType.ReservedHttp2MaxConcurrentStreams:
 871                        case Http3SettingType.ReservedHttp2InitialWindowSize:
 872                        case Http3SettingType.ReservedHttp2MaxFrameSize:
 873                            // Per https://tools.ietf.org/html/draft-ietf-quic-http-31#section-7.2.4.1
 874                            // these settings IDs are reserved and must never be sent.
 0875                            throw HttpProtocolException.CreateHttp3ConnectionException(Http3ErrorCode.SettingsError);
 876                    }
 0877                }
 0878            }
 879
 880            async ValueTask ProcessGoAwayFrameAsync(long goawayPayloadLength)
 0881            {
 882                long firstRejectedStreamId;
 883                int bytesRead;
 884
 0885                while (!VariableLengthIntegerHelper.TryRead(buffer.ActiveSpan, out firstRejectedStreamId, out bytesRead)
 0886                {
 0887                    buffer.EnsureAvailableSpace(VariableLengthIntegerHelper.MaximumEncodedLength);
 0888                    bytesRead = await stream.ReadAsync(buffer.AvailableMemory, CancellationToken.None).ConfigureAwait(fa
 889
 0890                    if (bytesRead != 0)
 0891                    {
 0892                        buffer.Commit(bytesRead);
 0893                    }
 894                    else
 0895                    {
 896                        // Our buffer has partial frame data in it but not enough to complete the read: bail out.
 0897                        throw HttpProtocolException.CreateHttp3ConnectionException(Http3ErrorCode.FrameError);
 898                    }
 0899                }
 900
 0901                buffer.Discard(bytesRead);
 0902                if (bytesRead != goawayPayloadLength)
 0903                {
 904                    // Frame contains unknown extra data after the integer.
 0905                    throw HttpProtocolException.CreateHttp3ConnectionException(Http3ErrorCode.FrameError);
 906                }
 907
 0908                OnServerGoAway(firstRejectedStreamId);
 0909            }
 910
 911            async ValueTask SkipUnknownPayloadAsync(long payloadLength)
 0912            {
 0913                while (payloadLength != 0)
 0914                {
 0915                    if (buffer.ActiveLength == 0)
 0916                    {
 0917                        int bytesRead = await stream.ReadAsync(buffer.AvailableMemory, CancellationToken.None).Configure
 918
 0919                        if (bytesRead != 0)
 0920                        {
 0921                            buffer.Commit(bytesRead);
 0922                        }
 923                        else
 0924                        {
 925                            // Our buffer has partial frame data in it but not enough to complete the read: bail out.
 0926                            throw HttpProtocolException.CreateHttp3ConnectionException(Http3ErrorCode.FrameError);
 927                        }
 0928                    }
 929
 0930                    long readLength = Math.Min(payloadLength, buffer.ActiveLength);
 0931                    buffer.Discard((int)readLength);
 0932                    payloadLength -= readLength;
 0933                }
 0934            }
 0935        }
 936    }
 937
 938    /// <summary>
 939    /// Tracks telemetry signals associated with the time period an HTTP/3 request spends waiting for a usable HTTP/3 co
 940    /// the wait_for_connection Activity, the RequestLeftQueue EventSource event and the http.client.request.time_in_que
 941    /// </summary>
 942    internal struct WaitForHttp3ConnectionActivity
 943    {
 944        // The HttpConnectionSettings -> SocketsHttpHandlerMetrics indirection is needed for the trimmer.
 945        private HttpConnectionSettings _settings;
 946        private readonly HttpAuthority _authority;
 947        private Activity? _activity;
 948        private long _startTimestamp;
 949
 950        public WaitForHttp3ConnectionActivity(HttpConnectionSettings settings, HttpAuthority authority)
 951        {
 952            _settings = settings;
 953            _authority = authority;
 954        }
 955
 956        public bool Started { get; private set; }
 957
 958        public void Start()
 959        {
 960            Debug.Assert(!Started);
 961            _startTimestamp = HttpTelemetry.Log.IsEnabled() || (GlobalHttpSettings.MetricsHandler.IsGloballyEnabled && _
 962            _activity = ConnectionSetupDistributedTracing.StartWaitForConnectionActivity(_authority);
 963            Started = true;
 964        }
 965
 966        public void Stop(HttpRequestMessage request, HttpConnectionPool pool, Exception? exception)
 967        {
 968            if (exception is not null)
 969            {
 970                ConnectionSetupDistributedTracing.ReportError(_activity, exception);
 971            }
 972
 973            _activity?.Stop();
 974
 975            if (_startTimestamp != 0)
 976            {
 977                TimeSpan duration = Stopwatch.GetElapsedTime(_startTimestamp);
 978
 979                if (GlobalHttpSettings.MetricsHandler.IsGloballyEnabled)
 980                {
 981                    _settings._metrics!.RequestLeftQueue(request, pool, duration, versionMajor: 3);
 982                }
 983                if (HttpTelemetry.Log.IsEnabled())
 984                {
 985                    HttpTelemetry.Log.RequestLeftQueue(3, duration);
 986                }
 987            }
 988        }
 989
 990        [Conditional("DEBUG")]
 991        public void AssertActivityNotRunning()
 992        {
 993            Debug.Assert(_activity?.IsStopped != false);
 994        }
 995    }
 996}
 997