< Summary

Information
Class: System.Net.Http.HttpConnection
Assembly: System.Net.Http
File(s): File 1: https://raw.githubusercontent.com/dotnet/runtime/811a7eabb75c42db53440e8ba3f60c07511cfd1f/src/libraries/System.Net.Http/src/System/Net/Http/SocketsHttpHandler/ChunkedEncodingReadStream.cs
File 2: https://raw.githubusercontent.com/dotnet/runtime/811a7eabb75c42db53440e8ba3f60c07511cfd1f/src/libraries/System.Net.Http/src/System/Net/Http/SocketsHttpHandler/ChunkedEncodingWriteStream.cs
File 3: https://raw.githubusercontent.com/dotnet/runtime/811a7eabb75c42db53440e8ba3f60c07511cfd1f/src/libraries/System.Net.Http/src/System/Net/Http/SocketsHttpHandler/ConnectionCloseReadStream.cs
File 4: https://raw.githubusercontent.com/dotnet/runtime/811a7eabb75c42db53440e8ba3f60c07511cfd1f/src/libraries/System.Net.Http/src/System/Net/Http/SocketsHttpHandler/ContentLengthReadStream.cs
File 5: https://raw.githubusercontent.com/dotnet/runtime/811a7eabb75c42db53440e8ba3f60c07511cfd1f/src/libraries/System.Net.Http/src/System/Net/Http/SocketsHttpHandler/ContentLengthWriteStream.cs
File 6: https://raw.githubusercontent.com/dotnet/runtime/811a7eabb75c42db53440e8ba3f60c07511cfd1f/src/libraries/System.Net.Http/src/System/Net/Http/SocketsHttpHandler/HttpConnection.cs
File 7: https://raw.githubusercontent.com/dotnet/runtime/811a7eabb75c42db53440e8ba3f60c07511cfd1f/src/libraries/System.Net.Http/src/System/Net/Http/SocketsHttpHandler/HttpContentReadStream.cs
File 8: https://raw.githubusercontent.com/dotnet/runtime/811a7eabb75c42db53440e8ba3f60c07511cfd1f/src/libraries/System.Net.Http/src/System/Net/Http/SocketsHttpHandler/HttpContentWriteStream.cs
File 9: https://raw.githubusercontent.com/dotnet/runtime/811a7eabb75c42db53440e8ba3f60c07511cfd1f/src/libraries/System.Net.Http/src/System/Net/Http/SocketsHttpHandler/RawConnectionStream.cs
Line coverage
0%
Covered lines: 0
Uncovered lines: 2083
Coverable lines: 2083
Total lines: 3661
Line coverage: 0%
Branch coverage
0%
Covered branches: 0
Total branches: 802
Branch coverage: 0%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Cyclomatic complexity NPath complexity Sequence coverage
File 1: .ctor(...)100%110%
File 1: Read(...)0%28280%
File 1: ReadAsync(...)0%12120%
File 1: ReadAsyncCore(...)0%20200%
File 1: CopyToAsync(...)0%440%
File 1: CopyToAsyncCore(...)0%660%
File 1: PeekChunkFromConnectionBuffer()100%110%
File 1: ReadChunksFromConnectionBuffer(...)0%660%
File 1: ReadChunkFromConnectionBuffer(...)0%34340%
File 1: ValidateChunkExtension(...)0%880%
File 1: DrainAsync(...)0%16160%
File 1: Fill()0%220%
File 1: FillAsync()0%220%
File 2: .cctor()100%110%
File 2: .ctor(...)100%110%
File 2: Write(...)0%220%
File 2: WriteAsync(...)0%220%
File 2: WriteChunkAsync(System.Net.Http.HttpConnection,System.ReadOnlyMemory`1<System.Byte>)100%110%
File 2: FinishAsync(...)0%220%
File 3: .ctor(...)100%110%
File 3: Read(...)0%660%
File 3: ReadAsync(...)0%880%
File 3: CopyToAsync(...)0%660%
File 3: CompleteCopyToAsync(...)100%110%
File 3: Finish(...)100%110%
File 4: .ctor(...)100%110%
File 4: Read(...)0%10100%
File 4: ReadAsync(...)0%12120%
File 4: CopyToAsync(...)0%660%
File 4: CompleteCopyToAsync(...)100%110%
File 4: Finish()100%110%
File 4: ReadFromConnectionBuffer(...)0%220%
File 4: DrainAsync(...)0%12120%
File 5: .ctor(...)100%110%
File 5: Write(...)0%220%
File 5: WriteAsync(...)0%220%
File 5: FinishAsync(...)0%220%
File 6: .cctor()100%110%
File 6: .ctor(...)0%220%
File 6: Finalize()100%110%
File 6: Dispose()100%110%
File 6: Dispose(...)0%880%
File 6: PrepareForReuse(...)0%14140%
File 6: TryOwnScavengingTaskCompletion()0%220%
File 6: TryReturnScavengingTaskCompletionOwnership()0%440%
File 6: CheckUsabilityOnScavenge()0%440%
File 6: ReadAheadWithZeroByteReadAsync()0%660%
File 6: TransitionToCompletedAndTryOwnCompletion()0%220%
File 6: CheckKeepAliveTimeoutExceeded()0%220%
File 6: ConsumeFromRemainingBuffer(...)100%110%
File 6: WriteHeaders(...)0%32320%
File 6: WriteHost(System.Uri)0%440%
File 6: WriteHeaderCollection(...)0%18180%
File 6: WriteCRLF()100%110%
File 6: WriteBytes(...)100%110%
File 6: WriteAsciiString(...)100%110%
File 6: WriteString(...)0%440%
File 6: ThrowForInvalidCharEncoding()100%110%
File 6: SendAsync(...)0%1201200%
File 6: MapSendException(...)0%10100%
File 6: CreateRequestContentStream(...)0%440%
File 6: RegisterCancellation(...)0%220%
File 6: SendRequestContentAsync(...)0%880%
File 6: SendRequestContentWithExpect100ContinueAsync(...)0%660%
File 6: ParseStatusLine(...)0%880%
File 6: ParseStatusLineCore(...)0%26260%
File 6: ParseHeaders(...)0%440%
File 6: ParseHeadersCore(...)0%26260%
File 6: ThrowForInvalidHeaderLine(System.ReadOnlySpan`1<System.Byte>,System.Int32)100%110%
File 6: AddResponseHeader(...)0%30300%
File 6: ThrowForEmptyHeaderName()100%110%
File 6: ThrowForInvalidHeaderName(System.ReadOnlySpan`1<System.Byte>)100%110%
File 6: ThrowExceededAllowedReadLineBytes()100%110%
File 6: ProcessKeepAliveHeader(...)0%18180%
File 6: WriteToBuffer(...)100%110%
File 6: Write(...)0%660%
File 6: WriteAsync(...)0%880%
File 6: AwaitFlushAndWriteAsync(System.Threading.Tasks.ValueTask,System.ReadOnlyMemory`1<System.Byte>)0%220%
File 6: WriteWithoutBuffering(...)0%440%
File 6: WriteWithoutBufferingAsync(...)0%440%
File 6: FlushThenWriteWithoutBufferingAsync(...)100%110%
File 6: WriteHexInt32Async(...)0%440%
File 6: Flush()0%220%
File 6: FlushAsync(...)0%220%
File 6: WriteToStream(...)0%220%
File 6: WriteToStreamAsync(...)0%440%
File 6: TryReadNextChunkedLine(...)0%14140%
File 6: InitialFillAsync(...)0%440%
File 6: FillAsync(...)0%660%
File 6: FillForHeadersAsync(...)0%220%
File 6: ReadUntilEndOfHeaderAsync(System.Boolean)0%12120%
File 6: TryFindEndOfLine(System.ReadOnlySpan`1<System.Byte>,System.Int32&)0%880%
File 6: ReadFromBuffer(...)100%110%
File 6: Read(...)0%440%
File 6: ReadAsync(...)0%440%
File 6: ReadAndLogBytesReadAsync(System.Memory`1<System.Byte>)0%220%
File 6: ReadBuffered(...)0%660%
File 6: ReadBufferedAsync(...)0%440%
File 6: ReadBufferedAsyncCore()0%440%
File 6: CopyFromBufferAsync(...)0%440%
File 6: CopyToUntilEofAsync(...)0%440%
File 6: CopyToUntilEofWithExistingBufferedDataAsync(...)0%220%
File 6: CopyToContentLengthAsync(...)0%18180%
File 6: Acquire()100%110%
File 6: Release()0%220%
File 6: DetachFromPool()100%110%
File 6: CompleteResponse()0%880%
File 6: DrainResponseAsync(...)0%12120%
File 6: ReturnConnectionToPool()0%440%
File 6: ToString()100%110%
File 6: Trace(...)0%440%
File 7: .ctor(...)100%110%
File 7: Write(...)100%110%
File 7: WriteAsync(...)100%110%
File 7: DrainAsync(...)100%110%
File 7: Dispose(...)0%660%
File 7: DrainOnDisposeAsync()0%660%
File 8: .ctor(...)100%110%
File 8: Flush()0%220%
File 8: FlushAsync(...)0%220%
File 8: Read(...)100%110%
File 8: ReadAsync(...)100%110%
File 8: CopyToAsync(...)100%110%
File 9: .ctor(...)0%220%
File 9: Read(...)0%660%
File 9: ReadAsync()0%880%
File 9: CopyToAsync(...)0%660%
File 9: CompleteCopyToAsync(...)100%110%
File 9: Finish(...)100%110%
File 9: Write(...)0%440%
File 9: WriteAsync(...)0%880%
File 9: Flush()0%220%
File 9: FlushAsync(...)0%660%
File 9: WaitWithConnectionCancellationAsync(...)100%110%

File(s)

https://raw.githubusercontent.com/dotnet/runtime/811a7eabb75c42db53440e8ba3f60c07511cfd1f/src/libraries/System.Net.Http/src/System/Net/Http/SocketsHttpHandler/ChunkedEncodingReadStream.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.Buffers.Text;
 5using System.Diagnostics;
 6using System.IO;
 7using System.Text;
 8using System.Threading;
 9using System.Threading.Tasks;
 10
 11namespace System.Net.Http
 12{
 13    internal sealed partial class HttpConnection
 14    {
 15        private sealed class ChunkedEncodingReadStream : HttpContentReadStream
 16        {
 17            /// <summary>The number of bytes remaining in the chunk.</summary>
 18            private ulong _chunkBytesRemaining;
 19            /// <summary>The current state of the parsing state machine for the chunked response.</summary>
 020            private ParsingState _state = ParsingState.ExpectChunkHeader;
 21            private readonly HttpResponseMessage _response;
 22
 023            public ChunkedEncodingReadStream(HttpConnection connection, HttpResponseMessage response) : base(connection)
 024            {
 025                Debug.Assert(response != null, "The HttpResponseMessage cannot be null.");
 026                _response = response;
 027            }
 28
 29            public override int Read(Span<byte> buffer)
 030            {
 031                if (_connection == null)
 032                {
 33                    // Response body fully consumed
 034                    return 0;
 35                }
 36
 037                if (buffer.Length == 0)
 038                {
 039                    if (PeekChunkFromConnectionBuffer())
 040                    {
 041                        return 0;
 42                    }
 043                }
 44                else
 045                {
 46                    // Try to consume from data we already have in the buffer.
 047                    int bytesRead = ReadChunksFromConnectionBuffer(buffer, cancellationRegistration: default);
 048                    if (bytesRead > 0)
 049                    {
 050                        return bytesRead;
 51                    }
 052                }
 53
 54                // Nothing available to consume.  Fall back to I/O.
 055                while (true)
 056                {
 057                    if (_connection == null)
 058                    {
 59                        // Fully consumed the response in ReadChunksFromConnectionBuffer.
 060                        return 0;
 61                    }
 62
 063                    if (_state == ParsingState.ExpectChunkData &&
 064                        buffer.Length >= _connection.ReadBufferSize &&
 065                        _chunkBytesRemaining >= (ulong)_connection.ReadBufferSize)
 066                    {
 67                        // As an optimization, we skip going through the connection's read buffer if both
 68                        // the remaining chunk data and the buffer are both at least as large
 69                        // as the connection buffer.  That avoids an unnecessary copy while still reading
 70                        // the maximum amount we'd otherwise read at a time.
 071                        Debug.Assert(_connection.RemainingBuffer.Length == 0);
 072                        Debug.Assert(buffer.Length != 0);
 073                        int bytesRead = _connection.Read(buffer.Slice(0, (int)Math.Min((ulong)buffer.Length, _chunkBytes
 074                        if (bytesRead == 0)
 075                        {
 076                            throw new HttpIOException(HttpRequestError.ResponseEnded, SR.Format(SR.net_http_invalid_resp
 77                        }
 078                        _chunkBytesRemaining -= (ulong)bytesRead;
 079                        if (_chunkBytesRemaining == 0)
 080                        {
 081                            _state = ParsingState.ExpectChunkTerminator;
 082                        }
 083                        return bytesRead;
 84                    }
 85
 086                    if (buffer.Length == 0)
 087                    {
 88                        // User requested a zero-byte read, and we have no data available in the buffer for processing.
 89                        // This zero-byte read indicates their desire to trade off the extra cost of a zero-byte read
 90                        // for reduced memory consumption when data is not immediately available.
 91                        // So, we will issue our own zero-byte read against the underlying stream to allow it to make us
 92                        // optimizations, such as deferring buffer allocation until data is actually available.
 093                        _connection.Read(buffer);
 094                    }
 95
 96                    // We're only here if we need more data to make forward progress.
 097                    Fill();
 98
 99                    // Now that we have more, see if we can get any response data, and if
 100                    // we can we're done.
 0101                    if (buffer.Length == 0)
 0102                    {
 0103                        if (PeekChunkFromConnectionBuffer())
 0104                        {
 0105                            return 0;
 106                        }
 0107                    }
 108                    else
 0109                    {
 0110                        int bytesCopied = ReadChunksFromConnectionBuffer(buffer, cancellationRegistration: default);
 0111                        if (bytesCopied > 0)
 0112                        {
 0113                            return bytesCopied;
 114                        }
 0115                    }
 0116                }
 0117            }
 118
 119            public override ValueTask<int> ReadAsync(Memory<byte> buffer, CancellationToken cancellationToken)
 0120            {
 0121                if (cancellationToken.IsCancellationRequested)
 0122                {
 123                    // Cancellation requested.
 0124                    return ValueTask.FromCanceled<int>(cancellationToken);
 125                }
 126
 0127                if (_connection == null)
 0128                {
 129                    // Response body fully consumed
 0130                    return new ValueTask<int>(0);
 131                }
 132
 0133                if (buffer.Length == 0)
 0134                {
 0135                    if (PeekChunkFromConnectionBuffer())
 0136                    {
 0137                        return new ValueTask<int>(0);
 138                    }
 0139                }
 140                else
 0141                {
 142                    // Try to consume from data we already have in the buffer.
 0143                    int bytesRead = ReadChunksFromConnectionBuffer(buffer.Span, cancellationRegistration: default);
 0144                    if (bytesRead > 0)
 0145                    {
 0146                        return new ValueTask<int>(bytesRead);
 147                    }
 0148                }
 149
 150                // We may have just consumed the remainder of the response (with no actual data
 151                // available), so check again.
 0152                if (_connection == null)
 0153                {
 0154                    Debug.Assert(_state == ParsingState.Done);
 0155                    return new ValueTask<int>(0);
 156                }
 157
 158                // Nothing available to consume.  Fall back to I/O.
 0159                return ReadAsyncCore(buffer, cancellationToken);
 0160            }
 161
 162            private async ValueTask<int> ReadAsyncCore(Memory<byte> buffer, CancellationToken cancellationToken)
 0163            {
 164                // Should only be called if ReadChunksFromConnectionBuffer returned 0.
 165
 0166                Debug.Assert(_connection != null);
 167
 0168                CancellationTokenRegistration ctr = _connection.RegisterCancellation(cancellationToken);
 169                try
 0170                {
 0171                    while (true)
 0172                    {
 0173                        if (_connection == null)
 0174                        {
 175                            // Fully consumed the response in ReadChunksFromConnectionBuffer.
 0176                            return 0;
 177                        }
 178
 0179                        if (_state == ParsingState.ExpectChunkData &&
 0180                            buffer.Length >= _connection.ReadBufferSize &&
 0181                            _chunkBytesRemaining >= (ulong)_connection.ReadBufferSize)
 0182                        {
 183                            // As an optimization, we skip going through the connection's read buffer if both
 184                            // the remaining chunk data and the buffer are both at least as large
 185                            // as the connection buffer.  That avoids an unnecessary copy while still reading
 186                            // the maximum amount we'd otherwise read at a time.
 0187                            Debug.Assert(_connection.RemainingBuffer.Length == 0);
 0188                            Debug.Assert(buffer.Length != 0);
 0189                            int bytesRead = await _connection.ReadAsync(buffer.Slice(0, (int)Math.Min((ulong)buffer.Leng
 0190                            if (bytesRead == 0)
 0191                            {
 0192                                throw new HttpIOException(HttpRequestError.ResponseEnded, SR.Format(SR.net_http_invalid_
 193                            }
 0194                            _chunkBytesRemaining -= (ulong)bytesRead;
 0195                            if (_chunkBytesRemaining == 0)
 0196                            {
 0197                                _state = ParsingState.ExpectChunkTerminator;
 0198                            }
 0199                            return bytesRead;
 200                        }
 201
 0202                        if (buffer.Length == 0)
 0203                        {
 204                            // User requested a zero-byte read, and we have no data available in the buffer for processi
 205                            // This zero-byte read indicates their desire to trade off the extra cost of a zero-byte rea
 206                            // for reduced memory consumption when data is not immediately available.
 207                            // So, we will issue our own zero-byte read against the underlying stream to allow it to mak
 208                            // optimizations, such as deferring buffer allocation until data is actually available.
 0209                            await _connection.ReadAsync(buffer).ConfigureAwait(false);
 0210                        }
 211
 212                        // We're only here if we need more data to make forward progress.
 0213                        await FillAsync().ConfigureAwait(false);
 214
 215                        // Now that we have more, see if we can get any response data, and if
 216                        // we can we're done.
 0217                        if (buffer.Length == 0)
 0218                        {
 0219                            if (PeekChunkFromConnectionBuffer())
 0220                            {
 0221                                return 0;
 222                            }
 0223                        }
 224                        else
 0225                        {
 0226                            int bytesCopied = ReadChunksFromConnectionBuffer(buffer.Span, ctr);
 0227                            if (bytesCopied > 0)
 0228                            {
 0229                                return bytesCopied;
 230                            }
 0231                        }
 0232                    }
 233                }
 0234                catch (Exception exc) when (CancellationHelper.ShouldWrapInOperationCanceledException(exc, cancellationT
 0235                {
 0236                    throw CancellationHelper.CreateOperationCanceledException(exc, cancellationToken);
 237                }
 238                finally
 0239                {
 0240                    ctr.Dispose();
 0241                }
 0242            }
 243
 244            public override Task CopyToAsync(Stream destination, int bufferSize, CancellationToken cancellationToken)
 0245            {
 0246                ValidateCopyToArguments(destination, bufferSize);
 247
 0248                return
 0249                    cancellationToken.IsCancellationRequested ? Task.FromCanceled(cancellationToken) :
 0250                    _connection == null ? Task.CompletedTask :
 0251                    CopyToAsyncCore(destination, cancellationToken);
 0252            }
 253
 254            private async Task CopyToAsyncCore(Stream destination, CancellationToken cancellationToken)
 0255            {
 0256                CancellationTokenRegistration ctr = _connection!.RegisterCancellation(cancellationToken);
 257                try
 0258                {
 0259                    while (true)
 0260                    {
 0261                        while (true)
 0262                        {
 0263                            if (ReadChunkFromConnectionBuffer(int.MaxValue, ctr) is not ReadOnlyMemory<byte> bytesRead |
 0264                            {
 0265                                break;
 266                            }
 0267                            await destination.WriteAsync(bytesRead, cancellationToken).ConfigureAwait(false);
 0268                        }
 269
 0270                        if (_connection == null)
 0271                        {
 272                            // Fully consumed the response.
 0273                            return;
 274                        }
 275
 0276                        await FillAsync().ConfigureAwait(false);
 0277                    }
 278                }
 0279                catch (Exception exc) when (CancellationHelper.ShouldWrapInOperationCanceledException(exc, cancellationT
 0280                {
 0281                    throw CancellationHelper.CreateOperationCanceledException(exc, cancellationToken);
 282                }
 283                finally
 0284                {
 0285                    ctr.Dispose();
 0286                }
 0287            }
 288
 289            private bool PeekChunkFromConnectionBuffer()
 0290            {
 0291                return ReadChunkFromConnectionBuffer(maxBytesToRead: 0, cancellationRegistration: default).HasValue;
 0292            }
 293
 294            private int ReadChunksFromConnectionBuffer(Span<byte> buffer, CancellationTokenRegistration cancellationRegi
 0295            {
 0296                Debug.Assert(buffer.Length > 0);
 0297                int totalBytesRead = 0;
 0298                while (buffer.Length > 0)
 0299                {
 0300                    if (ReadChunkFromConnectionBuffer(buffer.Length, cancellationRegistration) is not ReadOnlyMemory<byt
 0301                    {
 0302                        break;
 303                    }
 304
 0305                    Debug.Assert(bytesRead.Length <= buffer.Length);
 0306                    totalBytesRead += bytesRead.Length;
 0307                    bytesRead.Span.CopyTo(buffer);
 0308                    buffer = buffer.Slice(bytesRead.Length);
 0309                }
 0310                return totalBytesRead;
 0311            }
 312
 313            private ReadOnlyMemory<byte>? ReadChunkFromConnectionBuffer(int maxBytesToRead, CancellationTokenRegistratio
 0314            {
 0315                Debug.Assert(_connection != null);
 316
 317                try
 0318                {
 319                    ReadOnlySpan<byte> currentLine;
 0320                    switch (_state)
 321                    {
 322                        case ParsingState.ExpectChunkHeader:
 0323                            Debug.Assert(_chunkBytesRemaining == 0, $"Expected {nameof(_chunkBytesRemaining)} == 0, got 
 324
 325                            // Read the chunk header line.
 0326                            if (!_connection.TryReadNextChunkedLine(out currentLine))
 0327                            {
 328                                // Could not get a whole line, so we can't parse the chunk header.
 0329                                return default;
 330                            }
 331
 332                            // Parse the hex value from it.
 0333                            if (!Utf8Parser.TryParse(currentLine, out ulong chunkSize, out int bytesConsumed, 'X'))
 0334                            {
 0335                                throw new HttpIOException(HttpRequestError.InvalidResponse, SR.Format(SR.net_http_invali
 336                            }
 0337                            _chunkBytesRemaining = chunkSize;
 338
 339                            // If there's a chunk extension after the chunk size, validate it.
 0340                            if (bytesConsumed != currentLine.Length)
 0341                            {
 0342                                ValidateChunkExtension(currentLine.Slice(bytesConsumed));
 0343                            }
 344
 345                            // Proceed to handle the chunk.  If there's data in it, go read it.
 346                            // Otherwise, finish handling the response.
 0347                            if (chunkSize > 0)
 0348                            {
 0349                                _state = ParsingState.ExpectChunkData;
 0350                                goto case ParsingState.ExpectChunkData;
 351                            }
 352                            else
 0353                            {
 0354                                _state = ParsingState.ConsumeTrailers;
 0355                                goto case ParsingState.ConsumeTrailers;
 356                            }
 357
 358                        case ParsingState.ExpectChunkData:
 0359                            Debug.Assert(_chunkBytesRemaining > 0);
 360
 0361                            ReadOnlyMemory<byte> connectionBuffer = _connection.RemainingBuffer;
 0362                            if (connectionBuffer.Length == 0)
 0363                            {
 0364                                return default;
 365                            }
 366
 0367                            int bytesToConsume = Math.Min(maxBytesToRead, (int)Math.Min((ulong)connectionBuffer.Length, 
 0368                            Debug.Assert(bytesToConsume > 0 || maxBytesToRead == 0);
 369
 0370                            _connection.ConsumeFromRemainingBuffer(bytesToConsume);
 0371                            _chunkBytesRemaining -= (ulong)bytesToConsume;
 0372                            if (_chunkBytesRemaining == 0)
 0373                            {
 0374                                _state = ParsingState.ExpectChunkTerminator;
 0375                            }
 376
 0377                            return connectionBuffer.Slice(0, bytesToConsume);
 378
 379                        case ParsingState.ExpectChunkTerminator:
 0380                            Debug.Assert(_chunkBytesRemaining == 0, $"Expected {nameof(_chunkBytesRemaining)} == 0, got 
 381
 0382                            if (!_connection.TryReadNextChunkedLine(out currentLine))
 0383                            {
 0384                                return default;
 385                            }
 386
 0387                            if (currentLine.Length != 0)
 0388                            {
 0389                                throw new HttpIOException(HttpRequestError.InvalidResponse, SR.Format(SR.net_http_invali
 390                            }
 391
 0392                            _state = ParsingState.ExpectChunkHeader;
 0393                            goto case ParsingState.ExpectChunkHeader;
 394
 395                        case ParsingState.ConsumeTrailers:
 0396                            Debug.Assert(_chunkBytesRemaining == 0, $"Expected {nameof(_chunkBytesRemaining)} == 0, got 
 397
 398                            // Consume the receive buffer. If the stream is disposed, pass a null response to avoid
 399                            // processing headers for a connection returned to the pool.
 0400                            if (_connection.ParseHeaders(IsDisposed ? null : _response, isFromTrailer: true))
 0401                            {
 402                                // Dispose of the registration and then check whether cancellation has been
 403                                // requested. This is necessary to make deterministic a race condition between
 404                                // cancellation being requested and unregistering from the token.  Otherwise,
 405                                // it's possible cancellation could be requested just before we unregister and
 406                                // we then return a connection to the pool that has been or will be disposed
 407                                // (e.g. if a timer is used and has already queued its callback but the
 408                                // callback hasn't yet run).
 0409                                cancellationRegistration.Dispose();
 0410                                CancellationHelper.ThrowIfCancellationRequested(cancellationRegistration.Token);
 411
 0412                                _state = ParsingState.Done;
 0413                                _connection.CompleteResponse();
 0414                                _connection = null;
 0415                            }
 416
 0417                            return default;
 418
 419                        default:
 420                        case ParsingState.Done: // shouldn't be called once we're done
 0421                            Debug.Fail($"Unexpected state: {_state}");
 422                            if (NetEventSource.Log.IsEnabled())
 423                            {
 424                                NetEventSource.Error(this, $"Unexpected state: {_state}");
 425                            }
 426
 427                            return default;
 428                    }
 429                }
 0430                catch (Exception)
 0431                {
 432                    // Ensure we don't try to read from the connection again (in particular, for draining)
 0433                    _connection!.Dispose();
 0434                    _connection = null;
 0435                    throw;
 436                }
 0437            }
 438
 439            private static void ValidateChunkExtension(ReadOnlySpan<byte> lineAfterChunkSize)
 0440            {
 441                // Until we see the ';' denoting the extension, the line after the chunk size
 442                // must contain only tabs and spaces.  After the ';', anything goes.
 0443                for (int i = 0; i < lineAfterChunkSize.Length; i++)
 0444                {
 0445                    byte c = lineAfterChunkSize[i];
 0446                    if (c == ';')
 0447                    {
 0448                        break;
 449                    }
 0450                    else if (c != ' ' && c != '\t') // not called out in the RFC, but WinHTTP allows it
 0451                    {
 0452                        throw new HttpIOException(HttpRequestError.InvalidResponse, SR.Format(SR.net_http_invalid_respon
 453                    }
 0454                }
 0455            }
 456
 457            private enum ParsingState : byte
 458            {
 459                ExpectChunkHeader,
 460                ExpectChunkData,
 461                ExpectChunkTerminator,
 462                ConsumeTrailers,
 463                Done
 464            }
 465
 0466            public override bool NeedsDrain => CanReadFromConnection;
 467
 468            public override async ValueTask<bool> DrainAsync(int maxDrainBytes)
 0469            {
 0470                Debug.Assert(_connection != null);
 471
 0472                CancellationTokenSource? cts = null;
 0473                CancellationTokenRegistration ctr = default;
 474                try
 0475                {
 0476                    int drainedBytes = 0;
 0477                    while (true)
 0478                    {
 0479                        drainedBytes += _connection.RemainingBuffer.Length;
 0480                        while (true)
 0481                        {
 0482                            if (ReadChunkFromConnectionBuffer(int.MaxValue, ctr) is not ReadOnlyMemory<byte> bytesRead |
 0483                            {
 0484                                break;
 485                            }
 0486                        }
 487
 488                        // When ReadChunkFromConnectionBuffer reads the final chunk, it will clear out _connection
 489                        // and return the connection to the pool.
 0490                        if (_connection == null)
 0491                        {
 0492                            return true;
 493                        }
 494
 0495                        if (drainedBytes >= maxDrainBytes)
 0496                        {
 0497                            return false;
 498                        }
 499
 0500                        if (cts == null) // only create the drain timer if we have to go async
 0501                        {
 0502                            TimeSpan drainTime = _connection._pool.Settings._maxResponseDrainTime;
 503
 0504                            if (drainTime == TimeSpan.Zero)
 0505                            {
 0506                                return false;
 507                            }
 508
 0509                            if (drainTime != Timeout.InfiniteTimeSpan)
 0510                            {
 0511                                cts = new CancellationTokenSource((int)drainTime.TotalMilliseconds);
 0512                                ctr = cts.Token.Register(static s => ((HttpConnection)s!).Dispose(), _connection);
 0513                            }
 0514                        }
 515
 0516                        await FillAsync().ConfigureAwait(false);
 0517                    }
 518                }
 519                finally
 0520                {
 0521                    ctr.Dispose();
 0522                    cts?.Dispose();
 0523                }
 0524            }
 525
 526            private void Fill()
 0527            {
 0528                Debug.Assert(_connection is not null);
 0529                ValueTask fillTask = _state == ParsingState.ConsumeTrailers
 0530                    ? _connection.FillForHeadersAsync(async: false)
 0531                    : _connection.FillAsync(async: false);
 0532                Debug.Assert(fillTask.IsCompleted);
 0533                fillTask.GetAwaiter().GetResult();
 0534            }
 535
 536            private ValueTask FillAsync()
 0537            {
 0538                Debug.Assert(_connection is not null);
 0539                return _state == ParsingState.ConsumeTrailers
 0540                    ? _connection.FillForHeadersAsync(async: true)
 0541                    : _connection.FillAsync(async: true);
 0542            }
 543        }
 544    }
 545}
 546

https://raw.githubusercontent.com/dotnet/runtime/811a7eabb75c42db53440e8ba3f60c07511cfd1f/src/libraries/System.Net.Http/src/System/Net/Http/SocketsHttpHandler/ChunkedEncodingWriteStream.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.Diagnostics;
 5using System.Threading;
 6using System.Threading.Tasks;
 7
 8namespace System.Net.Http
 9{
 10    internal sealed partial class HttpConnection : IDisposable
 11    {
 12        private sealed class ChunkedEncodingWriteStream : HttpContentWriteStream
 13        {
 014            private static readonly byte[] s_crlfBytes = "\r\n"u8.ToArray();
 015            private static readonly byte[] s_finalChunkBytes = "0\r\n\r\n"u8.ToArray();
 16
 017            public ChunkedEncodingWriteStream(HttpConnection connection) : base(connection)
 018            {
 019            }
 20
 21            public override void Write(ReadOnlySpan<byte> buffer)
 022            {
 023                BytesWritten += buffer.Length;
 24
 025                HttpConnection connection = GetConnectionOrThrow();
 026                Debug.Assert(connection._currentRequest != null);
 27
 028                if (buffer.Length == 0)
 029                {
 030                    connection.Flush();
 031                    return;
 32                }
 33
 34                // Write chunk length in hex followed by \r\n
 035                ValueTask writeTask = connection.WriteHexInt32Async(buffer.Length, async: false);
 036                Debug.Assert(writeTask.IsCompleted);
 037                writeTask.GetAwaiter().GetResult();
 038                connection.Write(s_crlfBytes);
 39
 40                // Write chunk contents followed by \r\n
 041                connection.Write(buffer);
 042                connection.Write(s_crlfBytes);
 043            }
 44
 45            public override ValueTask WriteAsync(ReadOnlyMemory<byte> buffer, CancellationToken ignored)
 046            {
 047                BytesWritten += buffer.Length;
 48
 049                HttpConnection connection = GetConnectionOrThrow();
 050                Debug.Assert(connection._currentRequest != null);
 51
 52                // The token is ignored because it's coming from SendAsync and the only operations
 53                // here are those that are already covered by the token having been registered with
 54                // to close the connection.
 55
 056                ValueTask task = buffer.Length == 0 ?
 057                    // Don't write if nothing was given, especially since we don't want to accidentally send a 0 chunk,
 058                    // which would indicate end of body.  Instead, just ensure no content is stuck in the buffer.
 059                    connection.FlushAsync(async: true) :
 060                    WriteChunkAsync(connection, buffer);
 61
 062                return task;
 63
 64                static async ValueTask WriteChunkAsync(HttpConnection connection, ReadOnlyMemory<byte> buffer)
 065                {
 66                    // Write chunk length in hex followed by \r\n
 067                    await connection.WriteHexInt32Async(buffer.Length, async: true).ConfigureAwait(false);
 068                    await connection.WriteAsync(s_crlfBytes).ConfigureAwait(false);
 69
 70                    // Write chunk contents followed by \r\n
 071                    await connection.WriteAsync(buffer).ConfigureAwait(false);
 072                    await connection.WriteAsync(s_crlfBytes).ConfigureAwait(false);
 073                }
 074            }
 75
 76            public override Task FinishAsync(bool async)
 077            {
 78                // Send 0 byte chunk to indicate end, then final CrLf
 079                HttpConnection connection = GetConnectionOrThrow();
 080                _connection = null;
 81
 082                if (async)
 083                {
 084                    return connection.WriteAsync(s_finalChunkBytes).AsTask();
 85                }
 86                else
 087                {
 088                    connection.Write(s_finalChunkBytes);
 089                    return Task.CompletedTask;
 90                }
 091            }
 92        }
 93    }
 94}
 95

https://raw.githubusercontent.com/dotnet/runtime/811a7eabb75c42db53440e8ba3f60c07511cfd1f/src/libraries/System.Net.Http/src/System/Net/Http/SocketsHttpHandler/ConnectionCloseReadStream.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.IO;
 5using System.Threading;
 6using System.Threading.Tasks;
 7
 8namespace System.Net.Http
 9{
 10    internal sealed partial class HttpConnection : IDisposable
 11    {
 12        private sealed class ConnectionCloseReadStream : HttpContentReadStream
 13        {
 014            public ConnectionCloseReadStream(HttpConnection connection) : base(connection)
 015            {
 016            }
 17
 18            public override int Read(Span<byte> buffer)
 019            {
 020                HttpConnection? connection = _connection;
 021                if (connection == null)
 022                {
 23                    // Response body fully consumed
 024                    return 0;
 25                }
 26
 027                int bytesRead = connection.Read(buffer);
 028                if (bytesRead == 0 && buffer.Length != 0)
 029                {
 30                    // We cannot reuse this connection, so close it.
 031                    _connection = null;
 032                    connection.Dispose();
 033                }
 34
 035                return bytesRead;
 036            }
 37
 38            public override async ValueTask<int> ReadAsync(Memory<byte> buffer, CancellationToken cancellationToken)
 039            {
 040                CancellationHelper.ThrowIfCancellationRequested(cancellationToken);
 41
 042                HttpConnection? connection = _connection;
 043                if (connection == null)
 044                {
 45                    // Response body fully consumed
 046                    return 0;
 47                }
 48
 049                ValueTask<int> readTask = connection.ReadAsync(buffer);
 50                int bytesRead;
 051                if (readTask.IsCompletedSuccessfully)
 052                {
 053                    bytesRead = readTask.Result;
 054                }
 55                else
 056                {
 057                    CancellationTokenRegistration ctr = connection.RegisterCancellation(cancellationToken);
 58                    try
 059                    {
 060                        bytesRead = await readTask.ConfigureAwait(false);
 061                    }
 062                    catch (Exception exc) when (CancellationHelper.ShouldWrapInOperationCanceledException(exc, cancellat
 063                    {
 064                        throw CancellationHelper.CreateOperationCanceledException(exc, cancellationToken);
 65                    }
 66                    finally
 067                    {
 068                        ctr.Dispose();
 069                    }
 070                }
 71
 072                if (bytesRead == 0 && buffer.Length != 0)
 073                {
 74                    // If cancellation is requested and tears down the connection, it could cause the read
 75                    // to return 0, which would otherwise signal the end of the data, but that would lead
 76                    // the caller to think that it actually received all of the data, rather than it ending
 77                    // early due to cancellation.  So we prioritize cancellation in this race condition, and
 78                    // if we read 0 bytes and then find that cancellation has requested, we assume cancellation
 79                    // was the cause and throw.
 080                    CancellationHelper.ThrowIfCancellationRequested(cancellationToken);
 81
 82                    // We cannot reuse this connection, so close it.
 083                    _connection = null;
 084                    connection.Dispose();
 085                }
 86
 087                return bytesRead;
 088            }
 89
 90            public override Task CopyToAsync(Stream destination, int bufferSize, CancellationToken cancellationToken)
 091            {
 092                ValidateCopyToArguments(destination, bufferSize);
 93
 094                if (cancellationToken.IsCancellationRequested)
 095                {
 096                    return Task.FromCanceled(cancellationToken);
 97                }
 98
 099                HttpConnection? connection = _connection;
 0100                if (connection == null)
 0101                {
 102                    // null if response body fully consumed
 0103                    return Task.CompletedTask;
 104                }
 105
 0106                Task copyTask = connection.CopyToUntilEofAsync(destination, async: true, bufferSize, cancellationToken);
 0107                if (copyTask.IsCompletedSuccessfully)
 0108                {
 0109                    Finish(connection);
 0110                    return Task.CompletedTask;
 111                }
 112
 0113                return CompleteCopyToAsync(copyTask, connection, cancellationToken);
 0114            }
 115
 116            private async Task CompleteCopyToAsync(Task copyTask, HttpConnection connection, CancellationToken cancellat
 0117            {
 0118                CancellationTokenRegistration ctr = connection.RegisterCancellation(cancellationToken);
 119                try
 0120                {
 0121                    await copyTask.ConfigureAwait(false);
 0122                }
 0123                catch (Exception exc) when (CancellationHelper.ShouldWrapInOperationCanceledException(exc, cancellationT
 0124                {
 0125                    throw CancellationHelper.CreateOperationCanceledException(exc, cancellationToken);
 126                }
 127                finally
 0128                {
 0129                    ctr.Dispose();
 0130                }
 131
 132                // If cancellation is requested and tears down the connection, it could cause the copy
 133                // to end early but think it ended successfully. So we prioritize cancellation in this
 134                // race condition, and if we find after the copy has completed that cancellation has
 135                // been requested, we assume the copy completed due to cancellation and throw.
 0136                CancellationHelper.ThrowIfCancellationRequested(cancellationToken);
 137
 0138                Finish(connection);
 0139            }
 140
 141            private void Finish(HttpConnection connection)
 0142            {
 143                // We cannot reuse this connection, so close it.
 0144                _connection = null;
 0145                connection.Dispose();
 0146            }
 147        }
 148    }
 149}
 150

https://raw.githubusercontent.com/dotnet/runtime/811a7eabb75c42db53440e8ba3f60c07511cfd1f/src/libraries/System.Net.Http/src/System/Net/Http/SocketsHttpHandler/ContentLengthReadStream.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.Diagnostics;
 5using System.IO;
 6using System.Threading;
 7using System.Threading.Tasks;
 8
 9namespace System.Net.Http
 10{
 11    internal sealed partial class HttpConnection : IDisposable
 12    {
 13        private sealed class ContentLengthReadStream : HttpContentReadStream
 14        {
 15            private ulong _contentBytesRemaining;
 16
 017            public ContentLengthReadStream(HttpConnection connection, ulong contentLength) : base(connection)
 018            {
 019                Debug.Assert(contentLength > 0, "Caller should have checked for 0.");
 020                _contentBytesRemaining = contentLength;
 021            }
 22
 23            public override int Read(Span<byte> buffer)
 024            {
 025                if (_connection == null)
 026                {
 27                    // Response body fully consumed
 028                    return 0;
 29                }
 30
 031                Debug.Assert(_contentBytesRemaining > 0);
 032                if ((ulong)buffer.Length > _contentBytesRemaining)
 033                {
 034                    buffer = buffer.Slice(0, (int)_contentBytesRemaining);
 035                }
 36
 037                int bytesRead = _connection.Read(buffer);
 038                if (bytesRead <= 0 && buffer.Length != 0)
 039                {
 40                    // Unexpected end of response stream.
 041                    throw new HttpIOException(HttpRequestError.ResponseEnded, SR.Format(SR.net_http_invalid_response_pre
 42                }
 43
 044                Debug.Assert((ulong)bytesRead <= _contentBytesRemaining);
 045                _contentBytesRemaining -= (ulong)bytesRead;
 46
 047                if (_contentBytesRemaining == 0)
 048                {
 49                    // End of response body
 050                    _connection.CompleteResponse();
 051                    _connection = null;
 052                }
 53
 054                return bytesRead;
 055            }
 56
 57            public override async ValueTask<int> ReadAsync(Memory<byte> buffer, CancellationToken cancellationToken)
 058            {
 059                CancellationHelper.ThrowIfCancellationRequested(cancellationToken);
 60
 061                if (_connection == null)
 062                {
 63                    // Response body fully consumed
 064                    return 0;
 65                }
 66
 067                Debug.Assert(_contentBytesRemaining > 0);
 68
 069                if ((ulong)buffer.Length > _contentBytesRemaining)
 070                {
 071                    buffer = buffer.Slice(0, (int)_contentBytesRemaining);
 072                }
 73
 074                ValueTask<int> readTask = _connection.ReadAsync(buffer);
 75                int bytesRead;
 076                if (readTask.IsCompletedSuccessfully)
 077                {
 078                    bytesRead = readTask.Result;
 079                }
 80                else
 081                {
 082                    CancellationTokenRegistration ctr = _connection.RegisterCancellation(cancellationToken);
 83                    try
 084                    {
 085                        bytesRead = await readTask.ConfigureAwait(false);
 086                    }
 087                    catch (Exception exc) when (CancellationHelper.ShouldWrapInOperationCanceledException(exc, cancellat
 088                    {
 089                        throw CancellationHelper.CreateOperationCanceledException(exc, cancellationToken);
 90                    }
 91                    finally
 092                    {
 093                        ctr.Dispose();
 094                    }
 095                }
 96
 097                if (bytesRead == 0 && buffer.Length != 0)
 098                {
 99                    // A cancellation request may have caused the EOF.
 0100                    CancellationHelper.ThrowIfCancellationRequested(cancellationToken);
 101
 102                    // Unexpected end of response stream.
 0103                    throw new HttpIOException(HttpRequestError.ResponseEnded, SR.Format(SR.net_http_invalid_response_pre
 104                }
 105
 0106                Debug.Assert((ulong)bytesRead <= _contentBytesRemaining);
 0107                _contentBytesRemaining -= (ulong)bytesRead;
 108
 0109                if (_contentBytesRemaining == 0)
 0110                {
 111                    // End of response body
 0112                    _connection.CompleteResponse();
 0113                    _connection = null;
 0114                }
 115
 0116                return bytesRead;
 0117            }
 118
 119            public override Task CopyToAsync(Stream destination, int bufferSize, CancellationToken cancellationToken)
 0120            {
 0121                ValidateCopyToArguments(destination, bufferSize);
 122
 0123                if (cancellationToken.IsCancellationRequested)
 0124                {
 0125                    return Task.FromCanceled(cancellationToken);
 126                }
 127
 0128                if (_connection == null)
 0129                {
 130                    // null if response body fully consumed
 0131                    return Task.CompletedTask;
 132                }
 133
 0134                Task copyTask = _connection.CopyToContentLengthAsync(destination, async: true, _contentBytesRemaining, b
 0135                if (copyTask.IsCompletedSuccessfully)
 0136                {
 0137                    Finish();
 0138                    return Task.CompletedTask;
 139                }
 140
 0141                return CompleteCopyToAsync(copyTask, cancellationToken);
 0142            }
 143
 144            private async Task CompleteCopyToAsync(Task copyTask, CancellationToken cancellationToken)
 0145            {
 0146                Debug.Assert(_connection != null);
 0147                CancellationTokenRegistration ctr = _connection.RegisterCancellation(cancellationToken);
 148                try
 0149                {
 0150                    await copyTask.ConfigureAwait(false);
 0151                }
 0152                catch (Exception exc) when (CancellationHelper.ShouldWrapInOperationCanceledException(exc, cancellationT
 0153                {
 0154                    throw CancellationHelper.CreateOperationCanceledException(exc, cancellationToken);
 155                }
 156                finally
 0157                {
 0158                    ctr.Dispose();
 0159                }
 160
 0161                Finish();
 0162            }
 163
 164            private void Finish()
 0165            {
 0166                _contentBytesRemaining = 0;
 0167                _connection!.CompleteResponse();
 0168                _connection = null;
 0169            }
 170
 171            // Based on ReadChunkFromConnectionBuffer; perhaps we should refactor into a common routine.
 172            private ReadOnlyMemory<byte> ReadFromConnectionBuffer(int maxBytesToRead)
 0173            {
 0174                Debug.Assert(maxBytesToRead > 0);
 0175                Debug.Assert(_contentBytesRemaining > 0);
 0176                Debug.Assert(_connection != null);
 177
 0178                ReadOnlyMemory<byte> connectionBuffer = _connection.RemainingBuffer;
 0179                if (connectionBuffer.Length == 0)
 0180                {
 0181                    return default;
 182                }
 183
 0184                int bytesToConsume = Math.Min(maxBytesToRead, (int)Math.Min((ulong)connectionBuffer.Length, _contentByte
 0185                Debug.Assert(bytesToConsume > 0);
 186
 0187                _connection.ConsumeFromRemainingBuffer(bytesToConsume);
 0188                _contentBytesRemaining -= (ulong)bytesToConsume;
 189
 0190                return connectionBuffer.Slice(0, bytesToConsume);
 0191            }
 192
 0193            public override bool NeedsDrain => CanReadFromConnection;
 194
 195            public override async ValueTask<bool> DrainAsync(int maxDrainBytes)
 0196            {
 0197                Debug.Assert(_connection != null);
 0198                Debug.Assert(_contentBytesRemaining > 0);
 199
 0200                ReadFromConnectionBuffer(int.MaxValue);
 0201                if (_contentBytesRemaining == 0)
 0202                {
 0203                    Finish();
 0204                    return true;
 205                }
 206
 0207                if (_contentBytesRemaining > (ulong)maxDrainBytes)
 0208                {
 0209                    return false;
 210                }
 211
 0212                CancellationTokenSource? cts = null;
 0213                CancellationTokenRegistration ctr = default;
 0214                TimeSpan drainTime = _connection._pool.Settings._maxResponseDrainTime;
 215
 0216                if (drainTime == TimeSpan.Zero)
 0217                {
 0218                    return false;
 219                }
 220
 0221                if (drainTime != Timeout.InfiniteTimeSpan)
 0222                {
 0223                    cts = new CancellationTokenSource((int)drainTime.TotalMilliseconds);
 0224                    ctr = cts.Token.Register(static s => ((HttpConnection)s!).Dispose(), _connection);
 0225                }
 226
 227                try
 0228                {
 0229                    while (true)
 0230                    {
 0231                        await _connection.FillAsync(async: true).ConfigureAwait(false);
 0232                        ReadFromConnectionBuffer(int.MaxValue);
 0233                        if (_contentBytesRemaining == 0)
 0234                        {
 235                            // Dispose of the registration and then check whether cancellation has been
 236                            // requested. This is necessary to make deterministic a race condition between
 237                            // cancellation being requested and unregistering from the token.  Otherwise,
 238                            // it's possible cancellation could be requested just before we unregister and
 239                            // we then return a connection to the pool that has been or will be disposed
 240                            // (e.g. if a timer is used and has already queued its callback but the
 241                            // callback hasn't yet run).
 0242                            ctr.Dispose();
 0243                            CancellationHelper.ThrowIfCancellationRequested(ctr.Token);
 244
 0245                            Finish();
 0246                            return true;
 247                        }
 0248                    }
 249                }
 250                finally
 0251                {
 0252                    ctr.Dispose();
 0253                    cts?.Dispose();
 0254                }
 0255            }
 256        }
 257    }
 258}
 259

https://raw.githubusercontent.com/dotnet/runtime/811a7eabb75c42db53440e8ba3f60c07511cfd1f/src/libraries/System.Net.Http/src/System/Net/Http/SocketsHttpHandler/ContentLengthWriteStream.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.Diagnostics;
 5using System.Runtime.ExceptionServices;
 6using System.Threading;
 7using System.Threading.Tasks;
 8
 9namespace System.Net.Http
 10{
 11    internal sealed partial class HttpConnection : IDisposable
 12    {
 13        private sealed class ContentLengthWriteStream : HttpContentWriteStream
 14        {
 15            private readonly long _contentLength;
 16
 17            public ContentLengthWriteStream(HttpConnection connection, long contentLength)
 018                : base(connection)
 019            {
 020                _contentLength = contentLength;
 021            }
 22
 23            public override void Write(ReadOnlySpan<byte> buffer)
 024            {
 025                BytesWritten += buffer.Length;
 26
 027                if (BytesWritten > _contentLength)
 028                {
 029                    throw new HttpRequestException(SR.net_http_content_write_larger_than_content_length);
 30                }
 31
 032                HttpConnection connection = GetConnectionOrThrow();
 033                Debug.Assert(connection._currentRequest != null);
 034                connection.Write(buffer);
 035            }
 36
 37            public override ValueTask WriteAsync(ReadOnlyMemory<byte> buffer, CancellationToken ignored) // token ignore
 038            {
 039                BytesWritten += buffer.Length;
 40
 041                if (BytesWritten > _contentLength)
 042                {
 043                    return ValueTask.FromException(ExceptionDispatchInfo.SetCurrentStackTrace(new HttpRequestException(S
 44                }
 45
 046                HttpConnection connection = GetConnectionOrThrow();
 047                Debug.Assert(connection._currentRequest != null);
 048                return connection.WriteAsync(buffer);
 049            }
 50
 51            public override Task FinishAsync(bool async)
 052            {
 053                if (BytesWritten != _contentLength)
 054                {
 055                    return Task.FromException(ExceptionDispatchInfo.SetCurrentStackTrace(new HttpRequestException(SR.For
 56                }
 57
 058                _connection = null;
 059                return Task.CompletedTask;
 060            }
 61        }
 62    }
 63}
 64

https://raw.githubusercontent.com/dotnet/runtime/811a7eabb75c42db53440e8ba3f60c07511cfd1f/src/libraries/System.Net.Http/src/System/Net/Http/SocketsHttpHandler/HttpConnection.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.Buffers;
 5using System.Buffers.Text;
 6using System.Diagnostics;
 7using System.Globalization;
 8using System.IO;
 9using System.Net.Http.Headers;
 10using System.Net.Sockets;
 11using System.Runtime.CompilerServices;
 12using System.Text;
 13using System.Threading;
 14using System.Threading.Tasks;
 15
 16namespace System.Net.Http
 17{
 18    internal sealed partial class HttpConnection : HttpConnectionBase
 19    {
 20        /// <summary>Default size of the read buffer used for the connection.</summary>
 21        private const int InitialReadBufferSize =
 22#if DEBUG
 23            10;
 24#else
 25            4096;
 26#endif
 27        /// <summary>Default size of the write buffer used for the connection.</summary>
 28        private const int InitialWriteBufferSize = InitialReadBufferSize;
 29        /// <summary>
 30        /// Size after which we'll close the connection rather than send the payload in response
 31        /// to final error status code sent by the server when using Expect: 100-continue.
 32        /// </summary>
 33        private const int Expect100ErrorSendThreshold = 1024;
 34        /// <summary>How long a chunk indicator is allowed to be.</summary>
 35        /// <remarks>
 36        /// While most chunks indicators will contain no more than ulong.MaxValue.ToString("X").Length characters,
 37        /// "chunk extensions" are allowed. We place a limit on how long a line can be to avoid OOM issues if an
 38        /// infinite chunk length is sent.  This value is arbitrary and can be changed as needed.
 39        /// </remarks>
 40        private const int MaxChunkBytesAllowed = 16 * 1024;
 41
 042        private static readonly ulong s_http10Bytes = BitConverter.ToUInt64("HTTP/1.0"u8);
 043        private static readonly ulong s_http11Bytes = BitConverter.ToUInt64("HTTP/1.1"u8);
 44
 45        internal readonly Stream _stream;
 46        private readonly TransportContext? _transportContext;
 47
 48        private HttpRequestMessage? _currentRequest;
 49        private ArrayBuffer _writeBuffer;
 50        private int _allowedReadLineBytes;
 51
 52        /// <summary>Reusable array used to get the values for each header being written to the wire.</summary>
 53        [ThreadStatic]
 54        private static string[]? t_headerValues;
 55
 56        private const int ReadAheadTask_NotStarted = 0;
 57        private const int ReadAheadTask_Started = 1;
 58        private const int ReadAheadTask_CompletionReserved = 2;
 59        private const int ReadAheadTask_Completed = 3;
 60        private int _readAheadTaskStatus;
 61        private ValueTask<int> _readAheadTask;
 62        private ArrayBuffer _readBuffer;
 63
 64        private int _keepAliveTimeoutSeconds; // 0 == no timeout
 65        private bool _inUse;
 66        private bool _detachedFromPool;
 67        private bool _canRetry;
 68        private bool _connectionClose; // Connection: close was seen on last response
 69
 70        private volatile bool _disposed;
 71
 72        public HttpConnection(
 73            HttpConnectionPool pool,
 74            Stream stream,
 75            TransportContext? transportContext,
 76            Activity? connectionSetupActivity,
 77            IPEndPoint? remoteEndPoint,
 78            long connectionId)
 079            : base(pool, connectionId, connectionSetupActivity, remoteEndPoint)
 080        {
 081            Debug.Assert(stream != null);
 82
 083            _stream = stream;
 84
 085            _transportContext = transportContext;
 86
 087            _writeBuffer = new ArrayBuffer(InitialWriteBufferSize, usePool: false);
 088            _readBuffer = new ArrayBuffer(InitialReadBufferSize, usePool: false);
 89
 090            if (NetEventSource.Log.IsEnabled()) TraceConnection(_stream);
 091        }
 92
 093        ~HttpConnection() => Dispose(disposing: false);
 94
 095        public override void Dispose() => Dispose(disposing: true);
 96
 97        private void Dispose(bool disposing)
 098        {
 99            // Ensure we're only disposed once.  Dispose could be called concurrently, for example,
 100            // if the request and the response were running concurrently and both incurred an exception.
 0101            if (!Interlocked.Exchange(ref _disposed, true))
 0102            {
 0103                if (NetEventSource.Log.IsEnabled()) Trace("Connection closing.");
 104
 0105                MarkConnectionAsClosed();
 106
 0107                if (!_detachedFromPool)
 0108                {
 0109                    _pool.InvalidateHttp11Connection(this, disposing);
 0110                }
 111
 0112                if (disposing)
 0113                {
 0114                    GC.SuppressFinalize(this);
 0115                    _stream.Dispose();
 0116                }
 0117            }
 0118        }
 119
 120        private bool ReadAheadTaskHasStarted =>
 0121            _readAheadTaskStatus != ReadAheadTask_NotStarted;
 122
 123        /// <summary>Prepare an idle connection to be used for a new request.
 124        /// The caller MUST call SendAsync afterwards if this method returns true, or dispose the connection if it retur
 125        /// <param name="async">Indicates whether the coming request will be sync or async.</param>
 126        /// <returns>True if connection can be used, false if it is invalid due to a timeout or receiving EOF or unexpec
 127        public bool PrepareForReuse(bool async)
 0128        {
 0129            if (CheckKeepAliveTimeoutExceeded())
 0130            {
 0131                return false;
 132            }
 133
 134            // We may already have a read-ahead task if we did a previous scavenge and haven't used the connection since
 135            // If the read-ahead task is completed, then we've received either EOF or erroneous data the connection, so 
 0136            if (ReadAheadTaskHasStarted)
 0137            {
 0138                Debug.Assert(_readAheadTaskStatus is ReadAheadTask_Started or ReadAheadTask_Completed);
 139
 0140                return Interlocked.Exchange(ref _readAheadTaskStatus, ReadAheadTask_CompletionReserved) == ReadAheadTask
 141            }
 142
 143            // Check to see if we've received anything on the connection; if we have, that's
 144            // either erroneous data (we shouldn't have received anything yet) or the connection
 145            // has been closed; either way, we can't use it.
 0146            if (!async && _stream is NetworkStream networkStream)
 0147            {
 148                // Directly poll the socket rather than doing an async read, so that we can
 149                // issue an appropriate sync read when we actually need it.
 150                try
 0151                {
 0152                    return !networkStream.Socket.Poll(0, SelectMode.SelectRead);
 153                }
 0154                catch (Exception e) when (e is SocketException || e is ObjectDisposedException)
 0155                {
 156                    // Poll can throw when used on a closed socket.
 0157                    return false;
 158                }
 159            }
 160            else
 0161            {
 0162                Debug.Assert(_readAheadTaskStatus == ReadAheadTask_NotStarted);
 0163                _readAheadTaskStatus = ReadAheadTask_CompletionReserved;
 164
 165                // Perform an async read on the stream, since we're going to need to read from it
 166                // anyway, and in doing so we can avoid the extra syscall.
 167                try
 0168                {
 169#pragma warning disable CA2012 // we're very careful to ensure the ValueTask is only consumed once, even though it's sto
 0170                    _readAheadTask = _stream.ReadAsync(_readBuffer.AvailableMemory);
 171#pragma warning restore CA2012
 172
 173                    // If the read-ahead task already completed, we can't reuse the connection.
 174                    // We're still responsible for observing potential exceptions thrown by the read-ahead task to avoid
 0175                    if (_readAheadTask.IsCompleted)
 0176                    {
 0177                        LogExceptions(_readAheadTask.AsTask());
 0178                        return false;
 179                    }
 180
 0181                    return true;
 182                }
 0183                catch (Exception error)
 0184                {
 185                    // If reading throws, eat the error and don't reuse the connection.
 0186                    if (NetEventSource.Log.IsEnabled()) Trace($"Error performing read ahead: {error}");
 0187                    return false;
 188                }
 189            }
 0190        }
 191
 192        /// <summary>Takes ownership of the scavenging task completion if it was started.
 193        /// The caller MUST call either SendAsync or return the completion ownership afterwards if this method returns t
 194        public bool TryOwnScavengingTaskCompletion()
 0195        {
 0196            Debug.Assert(_readAheadTaskStatus != ReadAheadTask_CompletionReserved);
 197
 0198            return !ReadAheadTaskHasStarted
 0199                || Interlocked.Exchange(ref _readAheadTaskStatus, ReadAheadTask_CompletionReserved) == ReadAheadTask_Sta
 0200        }
 201
 202        /// <summary>Returns ownership of the scavenging task completion if it was started.
 203        /// The caller MUST Dispose the connection afterwards if this method returns false.</summary>
 204        public bool TryReturnScavengingTaskCompletionOwnership()
 0205        {
 0206            Debug.Assert(_readAheadTaskStatus != ReadAheadTask_Started);
 207
 0208            if (!ReadAheadTaskHasStarted ||
 0209                Interlocked.Exchange(ref _readAheadTaskStatus, ReadAheadTask_Started) == ReadAheadTask_CompletionReserve
 0210            {
 0211                return true;
 212            }
 213
 214            // The read-ahead task has started, and we failed to transition back to Started.
 215            // This means that the read-ahead task has completed, and we can't reuse the connection. The caller must dis
 216            // We're still responsible for observing potential exceptions thrown by the read-ahead task to avoid leaking
 0217            LogExceptions(_readAheadTask.AsTask());
 0218            return false;
 0219        }
 220
 221        /// <summary>Check whether a currently idle connection is still usable, or should be scavenged.</summary>
 222        /// <returns>True if connection can be used, false if it is invalid due to a timeout or receiving EOF or unexpec
 223        public override bool CheckUsabilityOnScavenge()
 0224        {
 0225            if (CheckKeepAliveTimeoutExceeded())
 0226            {
 0227                return false;
 228            }
 229
 230            // We may already have a read-ahead task if we did a previous scavenge and haven't used the connection since
 0231            if (!ReadAheadTaskHasStarted)
 0232            {
 0233                Debug.Assert(_readAheadTask == default);
 234
 0235                _readAheadTaskStatus = ReadAheadTask_Started;
 236
 237#pragma warning disable CA2012 // we're very careful to ensure the ValueTask is only consumed once, even though it's sto
 0238                _readAheadTask = ReadAheadWithZeroByteReadAsync();
 239#pragma warning restore CA2012
 0240            }
 241
 242            // If the read-ahead task is completed, then we've received either EOF or erroneous data the connection, so 
 0243            return !_readAheadTask.IsCompleted;
 244
 245            async ValueTask<int> ReadAheadWithZeroByteReadAsync()
 0246            {
 0247                Debug.Assert(_readAheadTask == default);
 0248                Debug.Assert(_readBuffer.ActiveLength == 0);
 249
 250                try
 0251                {
 252                    // Issue a zero-byte read.
 253                    // If the underlying stream supports it, this will not complete until the stream has data available,
 254                    // which will avoid pinning the connection's read buffer (and possibly allow us to release it to the
 255                    // If not, it will complete immediately.
 0256                    await _stream.ReadAsync(Memory<byte>.Empty).ConfigureAwait(false);
 257
 258                    // We don't know for sure that the stream actually has data available, so we need to issue a real re
 0259                    int read = await _stream.ReadAsync(_readBuffer.AvailableMemory).ConfigureAwait(false);
 260
 261                    // PrepareForReuse will check TryOwnReadAheadTaskCompletion before calling into SendAsync.
 262                    // If we can own the completion from within the read-ahead task, it means that PrepareForReuse hasn'
 263                    // In that case we've received EOF/erroneous data before we sent the request headers, and the connec
 0264                    if (TransitionToCompletedAndTryOwnCompletion())
 0265                    {
 0266                        if (NetEventSource.Log.IsEnabled()) Trace("Read-ahead task observed data before the request was 
 0267                    }
 268
 0269                    return read;
 270                }
 0271                catch (Exception error) when (TransitionToCompletedAndTryOwnCompletion())
 0272                {
 0273                    if (NetEventSource.Log.IsEnabled()) Trace($"Error performing read ahead: {error}");
 274
 0275                    return 0;
 276                }
 277
 278                bool TransitionToCompletedAndTryOwnCompletion()
 0279                {
 0280                    Debug.Assert(_readAheadTaskStatus is ReadAheadTask_Started or ReadAheadTask_CompletionReserved);
 281
 0282                    return Interlocked.Exchange(ref _readAheadTaskStatus, ReadAheadTask_Completed) == ReadAheadTask_Star
 0283                }
 0284            }
 0285        }
 286
 287        private bool CheckKeepAliveTimeoutExceeded()
 0288        {
 289            // We intentionally honor the Keep-Alive timeout on all HTTP/1.X versions, not just 1.0. This is to maximize
 290            // servers that use a lower idle timeout than the client, but give us a hint in the form of a Keep-Alive tim
 291            // If _keepAliveTimeoutSeconds is 0, no timeout has been set.
 0292            return _keepAliveTimeoutSeconds != 0 &&
 0293                GetIdleTicks(Environment.TickCount64) >= _keepAliveTimeoutSeconds * 1000;
 0294        }
 295
 0296        public TransportContext? TransportContext => _transportContext;
 297
 0298        public HttpConnectionKind Kind => _pool.Kind;
 299
 0300        private int ReadBufferSize => _readBuffer.Capacity;
 301
 0302        private ReadOnlyMemory<byte> RemainingBuffer => _readBuffer.ActiveMemory;
 303
 304        private void ConsumeFromRemainingBuffer(int bytesToConsume)
 0305        {
 0306            Debug.Assert(bytesToConsume <= _readBuffer.ActiveLength);
 0307            _readBuffer.Discard(bytesToConsume);
 0308        }
 309
 310        private void WriteHeaders(HttpRequestMessage request)
 0311        {
 0312            Debug.Assert(request.RequestUri is not null);
 313
 314            // Write the request line
 0315            WriteBytes(request.Method.Http1EncodedBytes);
 316
 0317            if (request.Method.IsConnect)
 0318            {
 319                // RFC 7231 #section-4.3.6.
 320                // Write only CONNECT foo.com:345 HTTP/1.1
 0321                if (!request.HasHeaders || request.Headers.Host is not string host)
 0322                {
 0323                    throw new HttpRequestException(SR.net_http_request_no_host);
 324                }
 325
 0326                WriteAsciiString(host);
 0327            }
 328            else
 0329            {
 0330                if (Kind == HttpConnectionKind.Proxy)
 0331                {
 332                    // Proxied requests contain full URL
 0333                    Debug.Assert(request.RequestUri.Scheme == Uri.UriSchemeHttp);
 0334                    WriteBytes("http://"u8);
 0335                    WriteHost(request.RequestUri);
 0336                }
 337
 0338                WriteAsciiString(request.RequestUri.PathAndQuery);
 0339            }
 340
 341            // Fall back to 1.1 for all versions other than 1.0
 0342            Debug.Assert(request.Version.Major >= 0 && request.Version.Minor >= 0); // guaranteed by Version class
 0343            bool isHttp10 = request.Version.Minor == 0 && request.Version.Major == 1;
 0344            WriteBytes(isHttp10 ? " HTTP/1.0\r\n"u8 : " HTTP/1.1\r\n"u8);
 345
 346            // Write special additional headers.  If a host isn't in the headers list, then a Host header
 347            // wasn't set, so as it's required by HTTP 1.1 spec, send one based on the Request Uri.
 0348            if (!request.HasHeaders || request.Headers.Host is null)
 0349            {
 0350                if (_pool.HostHeaderLineBytes is byte[] hostHeaderLineBytes)
 0351                {
 0352                    Debug.Assert(Kind != HttpConnectionKind.Proxy);
 0353                    WriteBytes(hostHeaderLineBytes);
 0354                }
 355                else
 0356                {
 0357                    Debug.Assert(Kind == HttpConnectionKind.Proxy);
 0358                    WriteBytes(KnownHeaders.Host.AsciiBytesWithColonSpace);
 0359                    WriteHost(request.RequestUri);
 0360                    WriteCRLF();
 0361                }
 0362            }
 363
 364            // Determine cookies to send
 0365            string? cookiesFromContainer = null;
 0366            if (_pool.Settings._useCookies)
 0367            {
 0368                cookiesFromContainer = _pool.Settings._cookieContainer!.GetCookieHeader(request.RequestUri);
 0369                if (cookiesFromContainer == "")
 0370                {
 0371                    cookiesFromContainer = null;
 0372                }
 0373            }
 374
 375            // Write request headers
 0376            if (request.HasHeaders || cookiesFromContainer is not null)
 0377            {
 0378                WriteHeaderCollection(request.Headers, cookiesFromContainer);
 0379            }
 380
 381            // Write content headers
 0382            if (request.Content is HttpContent content)
 0383            {
 0384                WriteHeaderCollection(content.Headers);
 0385            }
 386            else
 0387            {
 388                // Write out Content-Length: 0 header to indicate no body,
 389                // unless this is a method that never has a body.
 0390                if (request.Method.MustHaveRequestBody)
 0391                {
 0392                    WriteBytes("Content-Length: 0\r\n"u8);
 0393                }
 0394            }
 395
 396            // CRLF for end of headers.
 0397            WriteCRLF();
 398
 399            void WriteHost(Uri requestUri)
 0400            {
 401                // Uri.IdnHost is missing '[', ']' characters around IPv6 address
 402                // and it also contains ScopeID for Link-Local addresses
 0403                string host = requestUri.HostNameType == UriHostNameType.IPv6 ? requestUri.Host : requestUri.IdnHost;
 0404                WriteAsciiString(host);
 405
 0406                if (!requestUri.IsDefaultPort)
 0407                {
 0408                    _writeBuffer.EnsureAvailableSpace(6);
 0409                    Span<byte> buffer = _writeBuffer.AvailableSpan;
 0410                    buffer[0] = (byte)':';
 0411                    bool success = ((uint)requestUri.Port).TryFormat(buffer.Slice(1), out int bytesWritten);
 0412                    Debug.Assert(success);
 0413                    _writeBuffer.Commit(bytesWritten + 1);
 0414                }
 0415            }
 0416        }
 417
 418        private void WriteHeaderCollection(HttpHeaders headers, string? cookiesFromContainer = null)
 0419        {
 0420            Debug.Assert(_currentRequest is not null);
 421
 0422            HeaderEncodingSelector<HttpRequestMessage>? encodingSelector = _pool.Settings._requestHeaderEncodingSelector
 0423            ref string[]? headerValues = ref t_headerValues;
 424
 0425            foreach (HeaderEntry header in headers.GetEntries())
 0426            {
 0427                if (header.Key.KnownHeader is KnownHeader knownHeader)
 0428                {
 0429                    WriteBytes(knownHeader.AsciiBytesWithColonSpace);
 0430                }
 431                else
 0432                {
 0433                    WriteAsciiString(header.Key.Name);
 0434                    WriteBytes(": "u8);
 0435                }
 436
 0437                int headerValuesCount = HttpHeaders.GetStoreValuesIntoStringArray(header.Key, header.Value, ref headerVa
 0438                Debug.Assert(headerValuesCount > 0, "No values for header??");
 439
 0440                Encoding? valueEncoding = encodingSelector?.Invoke(header.Key.Name, _currentRequest);
 441
 0442                WriteString(headerValues[0], valueEncoding);
 443
 0444                if (cookiesFromContainer is not null && header.Key.Equals(KnownHeaders.Cookie))
 0445                {
 0446                    WriteBytes("; "u8); // Cookies use "; " as the separator
 0447                    WriteString(cookiesFromContainer, valueEncoding);
 0448                    cookiesFromContainer = null;
 0449                }
 450
 451                // Some headers such as User-Agent and Server use space as a separator (see: ProductInfoHeaderParser)
 0452                if (headerValuesCount > 1)
 0453                {
 0454                    byte[] separator = header.Key.SeparatorBytes;
 455
 0456                    for (int i = 1; i < headerValuesCount; i++)
 0457                    {
 0458                        WriteBytes(separator);
 0459                        WriteString(headerValues[i], valueEncoding);
 0460                    }
 0461                }
 462
 0463                WriteCRLF();
 0464            }
 465
 0466            if (cookiesFromContainer is not null)
 0467            {
 0468                WriteBytes(KnownHeaders.Cookie.AsciiBytesWithColonSpace);
 0469                WriteString(cookiesFromContainer, encodingSelector?.Invoke(HttpKnownHeaderNames.Cookie, _currentRequest)
 0470                WriteCRLF();
 0471            }
 0472        }
 473
 474        private void WriteCRLF()
 0475        {
 0476            _writeBuffer.EnsureAvailableSpace(2);
 0477            Span<byte> buffer = _writeBuffer.AvailableSpan;
 0478            buffer[1] = (byte)'\n';
 0479            buffer[0] = (byte)'\r';
 0480            _writeBuffer.Commit(2);
 0481        }
 482
 483        private void WriteBytes(ReadOnlySpan<byte> bytes)
 0484        {
 0485            _writeBuffer.EnsureAvailableSpace(bytes.Length);
 0486            bytes.CopyTo(_writeBuffer.AvailableSpan);
 0487            _writeBuffer.Commit(bytes.Length);
 0488        }
 489
 490        private void WriteAsciiString(string s)
 0491        {
 0492            Debug.Assert(Ascii.IsValid(s));
 493
 0494            _writeBuffer.EnsureAvailableSpace(s.Length);
 495
 0496            OperationStatus status = Ascii.FromUtf16(s, _writeBuffer.AvailableSpan, out int bytesWritten);
 0497            Debug.Assert(status == OperationStatus.Done);
 0498            Debug.Assert(bytesWritten == s.Length);
 499
 0500            _writeBuffer.Commit(bytesWritten);
 0501        }
 502
 503        private void WriteString(string s, Encoding? encoding)
 0504        {
 0505            if (encoding is null)
 0506            {
 0507                _writeBuffer.EnsureAvailableSpace(s.Length);
 0508                Span<byte> buffer = _writeBuffer.AvailableSpan;
 509
 0510                OperationStatus status = Ascii.FromUtf16(s, buffer, out int bytesWritten);
 511
 0512                if (status == OperationStatus.InvalidData)
 0513                {
 0514                    ThrowForInvalidCharEncoding();
 0515                }
 516
 0517                Debug.Assert(status == OperationStatus.Done);
 0518                Debug.Assert(bytesWritten == s.Length);
 519
 0520                _writeBuffer.Commit(s.Length);
 0521            }
 522            else
 0523            {
 0524                _writeBuffer.EnsureAvailableSpace(encoding.GetMaxByteCount(s.Length));
 0525                int length = encoding.GetBytes(s, _writeBuffer.AvailableSpan);
 0526                _writeBuffer.Commit(length);
 0527            }
 528
 529            static void ThrowForInvalidCharEncoding() =>
 0530                throw new HttpRequestException(SR.net_http_request_invalid_char_encoding);
 0531        }
 532
 533        public async Task<HttpResponseMessage> SendAsync(HttpRequestMessage request, bool async, CancellationToken cance
 0534        {
 0535            request.ConnectionId = Id;
 536
 0537            Debug.Assert(_currentRequest == null, $"Expected null {nameof(_currentRequest)}.");
 0538            Debug.Assert(_readBuffer.ActiveLength == 0, "Unexpected data in read buffer");
 0539            Debug.Assert(_readAheadTaskStatus != ReadAheadTask_Started,
 0540                "The caller should have called PrepareForReuse or TryOwnScavengingTaskCompletion if the connection was i
 541
 0542            MarkConnectionAsNotIdle();
 543
 0544            TaskCompletionSource<bool>? allowExpect100ToContinue = null;
 0545            Task? sendRequestContentTask = null;
 546
 0547            _currentRequest = request;
 548
 0549            _canRetry = false;
 550
 551            // Send the request.
 0552            if (NetEventSource.Log.IsEnabled()) Trace($"Sending request: {request}");
 0553            if (ConnectionSetupActivity is not null) ConnectionSetupDistributedTracing.AddConnectionLinkToRequestActivit
 0554            CancellationTokenRegistration cancellationRegistration = RegisterCancellation(cancellationToken);
 555            try
 0556            {
 0557                if (HttpTelemetry.Log.IsEnabled()) HttpTelemetry.Log.RequestHeadersStart(Id);
 558
 0559                WriteHeaders(request);
 560
 0561                if (HttpTelemetry.Log.IsEnabled()) HttpTelemetry.Log.RequestHeadersStop();
 562
 0563                if (request.Content == null)
 0564                {
 565                    // We have nothing more to send, so flush out any headers we haven't yet sent.
 0566                    await FlushAsync(async).ConfigureAwait(false);
 0567                }
 568                else
 0569                {
 0570                    bool hasExpectContinueHeader = request.HasHeaders && request.Headers.ExpectContinue == true;
 0571                    if (NetEventSource.Log.IsEnabled()) Trace($"Request content is not null, start processing it. hasExp
 572
 573                    // Send the body if there is one.  We prefer to serialize the sending of the content before
 574                    // we try to receive any response, but if ExpectContinue has been set, we allow the sending
 575                    // to run concurrently until we receive the final status line, at which point we wait for it.
 0576                    if (!hasExpectContinueHeader)
 0577                    {
 0578                        await SendRequestContentAsync(request, CreateRequestContentStream(request), async, cancellationT
 0579                    }
 580                    else
 0581                    {
 582                        // We're sending an Expect: 100-continue header. We need to flush headers so that the server rec
 583                        // all of them, and we need to do so before initiating the send, as once we do that, it effectiv
 584                        // owns the right to write, and we don't want to concurrently be accessing the write buffer.
 0585                        await FlushAsync(async).ConfigureAwait(false);
 586
 587                        // Create a TCS we'll use to block the request content from being sent, and create a timer that'
 588                        // as a fail-safe to unblock the request content if we don't hear back from the server in a time
 589                        // Then kick off the request.  The TCS' result indicates whether content should be sent or not.
 0590                        allowExpect100ToContinue = new TaskCompletionSource<bool>();
 0591                        var expect100Timer = new Timer(
 0592                            static s => ((TaskCompletionSource<bool>)s!).TrySetResult(true),
 0593                            allowExpect100ToContinue, _pool.Settings._expect100ContinueTimeout, Timeout.InfiniteTimeSpan
 594#pragma warning disable CA2025
 0595                        sendRequestContentTask = SendRequestContentWithExpect100ContinueAsync(
 0596                            request, allowExpect100ToContinue.Task, CreateRequestContentStream(request), expect100Timer,
 597#pragma warning restore
 0598                    }
 0599                }
 600
 601                // Start to read response.
 0602                _allowedReadLineBytes = _pool.Settings.MaxResponseHeadersByteLength;
 603
 604                // We should not have any buffered data here; if there was, it should have been treated as an error
 605                // by the previous request handling.  (Note we do not support HTTP pipelining.)
 0606                Debug.Assert(_readBuffer.ActiveLength == 0);
 607
 608                // When the connection was taken out of the pool, a pre-emptive read was performed
 609                // into the read buffer. We need to consume that read prior to issuing another read.
 0610                if (ReadAheadTaskHasStarted)
 0611                {
 612                    // If the read-ahead task completed synchronously, it would have claimed ownership of its completion
 613                    // meaning that PrepareForReuse would have failed, and we wouldn't have called SendAsync.
 614                    // The task therefore shouldn't be 'default', as it's representing an async operation that had to yi
 0615                    Debug.Assert(_readAheadTask != default);
 0616                    Debug.Assert(_readAheadTaskStatus is ReadAheadTask_CompletionReserved or ReadAheadTask_Completed);
 617
 618                    // Handle the pre-emptive read.  For the async==false case, hopefully the read has
 619                    // already completed and this will be a nop, but if it hasn't, the caller will be forced to block
 620                    // waiting for the async operation to complete.  We will only hit this case for proxied HTTPS
 621                    // requests that use a pooled connection, as in that case we don't have a Socket we
 622                    // can poll and are forced to issue an async read.
 0623                    ValueTask<int> vt = _readAheadTask;
 0624                    _readAheadTask = default;
 625
 626                    int bytesRead;
 0627                    if (vt.IsCompleted)
 0628                    {
 0629                        bytesRead = vt.Result;
 0630                    }
 631                    else
 0632                    {
 0633                        if (NetEventSource.Log.IsEnabled() && !async)
 0634                        {
 0635                            Trace($"Pre-emptive read completed asynchronously for a synchronous request.");
 0636                        }
 637
 0638                        bytesRead = await vt.ConfigureAwait(false);
 0639                    }
 640
 0641                    _readBuffer.Commit(bytesRead);
 642
 0643                    if (NetEventSource.Log.IsEnabled()) Trace($"Received {bytesRead} bytes.");
 644
 0645                    _readAheadTaskStatus = ReadAheadTask_NotStarted;
 0646                }
 647                else
 0648                {
 649                    // No read-ahead, so issue a read ourselves. We will check below for EOF.
 0650                    await InitialFillAsync(async).ConfigureAwait(false);
 0651                }
 652
 0653                if (_readBuffer.ActiveLength == 0)
 0654                {
 655                    // The server shutdown the connection on their end, likely because of an idle timeout.
 656                    // If we haven't started sending the request body yet (or there is no request body),
 657                    // then we allow the request to be retried.
 0658                    if (request.Content is null || allowExpect100ToContinue is not null)
 0659                    {
 0660                        _canRetry = true;
 0661                    }
 662
 0663                    throw new HttpIOException(HttpRequestError.ResponseEnded, SR.net_http_invalid_response_premature_eof
 664                }
 665
 666
 667                // Parse the response status line.
 0668                var response = new HttpResponseMessage() { RequestMessage = request, Content = new HttpConnectionRespons
 669
 0670                while (!ParseStatusLine(response))
 0671                {
 0672                    await FillForHeadersAsync(async).ConfigureAwait(false);
 0673                }
 674
 0675                if (HttpTelemetry.Log.IsEnabled()) HttpTelemetry.Log.ResponseHeadersStart();
 676
 677                // Multiple 1xx responses handling.
 678                // RFC 7231: A client MUST be able to parse one or more 1xx responses received prior to a final response
 679                // even if the client does not expect one. A user agent MAY ignore unexpected 1xx responses.
 680                // In .NET Core, apart from 100 Continue, and 101 Switching Protocols, we will treat all other 1xx respo
 681                // as unknown, and will discard them.
 0682                while ((uint)(response.StatusCode - 100) <= 199 - 100)
 0683                {
 684                    // If other 1xx responses come before an expected 100 continue, we will wait for the 100 response be
 685                    // sending request body (if any).
 0686                    if (allowExpect100ToContinue != null && response.StatusCode == HttpStatusCode.Continue)
 0687                    {
 0688                        allowExpect100ToContinue.TrySetResult(true);
 0689                        allowExpect100ToContinue = null;
 0690                    }
 0691                    else if (response.StatusCode == HttpStatusCode.SwitchingProtocols)
 0692                    {
 693                        // 101 Upgrade is a final response as it's used to switch protocols with WebSockets handshake.
 694                        // Will return a response object with status 101 and a raw connection stream later.
 695                        // RFC 7230: If a server receives both an Upgrade and an Expect header field with the "100-conti
 696                        // the server MUST send a 100 (Continue) response before sending a 101 (Switching Protocols) res
 697                        // If server doesn't follow RFC, we treat 101 as a final response and stop waiting for 100 conti
 698                        // never sends a 100-continue. The request body will be sent after expect100Timer expires.
 0699                        break;
 700                    }
 701
 702                    // In case read hangs which eventually leads to connection timeout.
 0703                    if (NetEventSource.Log.IsEnabled()) Trace($"Current {response.StatusCode} response is an interim res
 704
 705                    // Discard headers that come with the interim 1xx responses.
 0706                    while (!ParseHeaders(response: null, isFromTrailer: false))
 0707                    {
 0708                        await FillForHeadersAsync(async).ConfigureAwait(false);
 0709                    }
 710
 711                    // Parse the status line for next response.
 0712                    while (!ParseStatusLine(response))
 0713                    {
 0714                        await FillForHeadersAsync(async).ConfigureAwait(false);
 0715                    }
 0716                }
 717
 718                // Parse the response headers.  Logic after this point depends on being able to examine headers in the r
 0719                while (!ParseHeaders(response, isFromTrailer: false))
 0720                {
 0721                    await FillForHeadersAsync(async).ConfigureAwait(false);
 0722                }
 723
 0724                if (HttpTelemetry.Log.IsEnabled()) HttpTelemetry.Log.ResponseHeadersStop((int)response.StatusCode);
 725
 0726                if (allowExpect100ToContinue != null)
 0727                {
 728                    // If we sent an Expect: 100-continue header, and didn't receive a 100-continue. Handle the final re
 729                    // Note that the developer may have added an Expect: 100-continue header even if there is no Content
 0730                    if ((int)response.StatusCode >= 300 &&
 0731                        request.Content != null &&
 0732                        (request.Content.Headers.ContentLength == null || request.Content.Headers.ContentLength.GetValue
 0733                        !AuthenticationHelper.IsSessionAuthenticationChallenge(response))
 0734                    {
 735                        // For error final status codes, try to avoid sending the payload if its size is unknown or if i
 736                        // If we already sent a header detailing the size of the payload, if we then don't send that pay
 737                        // for it and assume that the next request on the connection is actually this request's payload.
 738                        // to be closed.  However, we may have also lost a race condition with the Expect: 100-continue 
 739                        // we've already started sending the payload (we weren't able to cancel it), then we don't need 
 740                        // We also must not clone connection if we do NTLM or Negotiate authentication.
 0741                        allowExpect100ToContinue.TrySetResult(false);
 742
 0743                        if (!allowExpect100ToContinue.Task.Result) // if Result is true, the timeout already expired and
 0744                        {
 0745                            _connectionClose = true;
 0746                        }
 0747                    }
 748                    else
 0749                    {
 750                        // For any success status codes, for errors when the request content length is known to be small
 751                        // or for session-based authentication challenges, send the payload
 752                        // (if there is one... if there isn't, Content is null and thus allowExpect100ToContinue is also
 0753                        allowExpect100ToContinue.TrySetResult(true);
 0754                    }
 0755                }
 756
 757                // Determine whether we need to force close the connection when the request/response has completed.
 0758                if (response.Headers.ConnectionClose.GetValueOrDefault())
 0759                {
 0760                    _connectionClose = true;
 0761                }
 762
 763                // Now that we've received our final status line, wait for the request content to fully send.
 764                // In most common scenarios, the server won't send back a response until all of the request
 765                // content has been received, so this task should generally already be complete.
 0766                if (sendRequestContentTask != null)
 0767                {
 0768                    Task sendTask = sendRequestContentTask;
 0769                    sendRequestContentTask = null;
 0770                    await sendTask.ConfigureAwait(false);
 0771                }
 772
 773                // Now we are sure that the request was fully sent.
 0774                if (NetEventSource.Log.IsEnabled()) Trace("Request is fully sent.");
 775
 776                // We're about to create the response stream, at which point responsibility for canceling
 777                // the remainder of the response lies with the stream.  Thus we dispose of our registration
 778                // here (if an exception has occurred or does occur while creating/returning the stream,
 779                // we'll still dispose of it in the catch below as part of Dispose'ing the connection).
 0780                cancellationRegistration.Dispose();
 0781                CancellationHelper.ThrowIfCancellationRequested(cancellationToken); // in case cancellation may have dis
 782
 783                // Create the response stream.
 784                Stream responseStream;
 0785                if (request.Method.IsConnect && response.IsSuccessStatusCode)
 0786                {
 787                    // Successful response to CONNECT does not have body.
 788                    // What ever comes next should be opaque.
 0789                    responseStream = new RawConnectionStream(this);
 790
 791                    // Don't put connection back to the pool if we upgraded to tunnel.
 792                    // We cannot use it for normal HTTP requests any more.
 0793                    _connectionClose = true;
 794
 0795                    _pool.InvalidateHttp11Connection(this);
 0796                    _detachedFromPool = true;
 0797                }
 0798                else if (request.Method.IsHead || response.StatusCode is HttpStatusCode.NoContent or HttpStatusCode.NotM
 0799                {
 0800                    responseStream = EmptyReadStream.Instance;
 0801                    CompleteResponse();
 0802                }
 0803                else if (response.StatusCode == HttpStatusCode.SwitchingProtocols)
 0804                {
 0805                    responseStream = new RawConnectionStream(this);
 806
 807                    // Don't put connection back to the pool if we switched protocols.
 808                    // We cannot use it for normal HTTP requests any more.
 0809                    _connectionClose = true;
 810
 0811                    _pool.InvalidateHttp11Connection(this);
 0812                    _detachedFromPool = true;
 0813                }
 0814                else if (response.Headers.TransferEncodingChunked == true)
 0815                {
 0816                    responseStream = new ChunkedEncodingReadStream(this, response);
 0817                }
 0818                else if (response.Content.Headers.ContentLength != null)
 0819                {
 0820                    long contentLength = response.Content.Headers.ContentLength.GetValueOrDefault();
 0821                    if (contentLength <= 0)
 0822                    {
 0823                        responseStream = EmptyReadStream.Instance;
 0824                        CompleteResponse();
 0825                    }
 826                    else
 0827                    {
 0828                        responseStream = new ContentLengthReadStream(this, (ulong)contentLength);
 0829                    }
 0830                }
 831                else
 0832                {
 0833                    responseStream = new ConnectionCloseReadStream(this);
 0834                }
 0835                ((HttpConnectionResponseContent)response.Content).SetStream(responseStream);
 836
 0837                if (NetEventSource.Log.IsEnabled()) Trace($"Received response: {response}");
 838
 839                // Process Set-Cookie headers.
 0840                if (_pool.Settings._useCookies)
 0841                {
 0842                    CookieHelper.ProcessReceivedCookies(response, _pool.Settings._cookieContainer!);
 0843                }
 844
 0845                return response;
 846            }
 0847            catch (Exception error)
 0848            {
 849                // Clean up the cancellation registration in case we're still registered.
 0850                cancellationRegistration.Dispose();
 851
 852                // Make sure to complete the allowExpect100ToContinue task if it exists.
 0853                if (allowExpect100ToContinue is not null && !allowExpect100ToContinue.TrySetResult(false))
 0854                {
 855                    // allowExpect100ToContinue was already signaled and we may have started sending the request body.
 0856                    _canRetry = false;
 0857                }
 858
 0859                if (_readAheadTask != default)
 0860                {
 0861                    Debug.Assert(_readAheadTaskStatus is ReadAheadTask_CompletionReserved or ReadAheadTask_Completed);
 862
 0863                    LogExceptions(_readAheadTask.AsTask());
 0864                }
 865
 0866                if (NetEventSource.Log.IsEnabled()) Trace($"Error sending request: {error}");
 867
 868                // In the rare case where Expect: 100-continue was used and then processing
 869                // of the response headers encountered an error such that we weren't able to
 870                // wait for the sending to complete, it's possible the sending also encountered
 871                // an exception or potentially is still going and will encounter an exception
 872                // (we're about to Dispose for the connection). In such cases, we don't want any
 873                // exception in that sending task to become unobserved and raise alarm bells, so we
 874                // hook up a continuation that will log it.
 0875                if (sendRequestContentTask != null && !sendRequestContentTask.IsCompletedSuccessfully)
 0876                {
 877                    // In case the connection is disposed, it's most probable that
 878                    // expect100Continue timer expired and request content sending failed.
 879                    // We're awaiting the task to propagate the exception in this case.
 0880                    if (_disposed)
 0881                    {
 882                        try
 0883                        {
 0884                            await sendRequestContentTask.ConfigureAwait(false);
 0885                        }
 886                        // Map the exception the same way as we normally do.
 0887                        catch (Exception ex) when (MapSendException(ex, cancellationToken, out Exception mappedEx))
 0888                        {
 0889                            throw mappedEx;
 890                        }
 0891                    }
 0892                    LogExceptions(sendRequestContentTask);
 0893                }
 894
 895                // Now clean up the connection.
 0896                Dispose();
 897
 898                // At this point, we're going to throw an exception; we just need to
 899                // determine which exception to throw.
 0900                if (MapSendException(error, cancellationToken, out Exception mappedException))
 0901                {
 0902                    throw mappedException;
 903                }
 904                // Otherwise, just allow the original exception to propagate.
 0905                throw;
 906            }
 0907        }
 908
 909        private bool MapSendException(Exception exception, CancellationToken cancellationToken, out Exception mappedExce
 0910        {
 0911            if (CancellationHelper.ShouldWrapInOperationCanceledException(exception, cancellationToken))
 0912            {
 913                // Cancellation was requested, so assume that the failure is due to
 914                // the cancellation request. This is a bit unorthodox, as usually we'd
 915                // prioritize a non-OperationCanceledException over a cancellation
 916                // request to avoid losing potentially pertinent information.  But given
 917                // the cancellation design where we tear down the underlying connection upon
 918                // a cancellation request, which can then result in a myriad of different
 919                // exceptions (argument exceptions, object disposed exceptions, socket exceptions,
 920                // etc.), as a middle ground we treat it as cancellation, but still propagate the
 921                // original information as the inner exception, for diagnostic purposes.
 0922                mappedException = CancellationHelper.CreateOperationCanceledException(exception, cancellationToken);
 0923                return true;
 924            }
 925
 0926            if (exception is InvalidOperationException)
 0927            {
 928                // For consistency with other handlers we wrap the exception in an HttpRequestException.
 0929                mappedException = new HttpRequestException(SR.net_http_client_execution_error, exception);
 0930                return true;
 931            }
 932
 0933            if (exception is IOException ioe)
 0934            {
 935                // For consistency with other handlers we wrap the exception in an HttpRequestException.
 936                // If the request is retryable, indicate that on the exception.
 0937                HttpRequestError error = ioe is HttpIOException httpIoe ? httpIoe.HttpRequestError : HttpRequestError.Un
 0938                mappedException = new HttpRequestException(error, SR.net_http_client_execution_error, ioe, _canRetry ? R
 0939                return true;
 940            }
 941
 942            // Otherwise, just allow the original exception to propagate.
 0943            mappedException = exception;
 0944            return false;
 0945        }
 946
 947        private HttpContentWriteStream CreateRequestContentStream(HttpRequestMessage request)
 0948        {
 0949            Debug.Assert(request.Content is not null);
 0950            bool requestTransferEncodingChunked = request.HasHeaders && request.Headers.TransferEncodingChunked == true;
 0951            HttpContentWriteStream requestContentStream = requestTransferEncodingChunked ? (HttpContentWriteStream)
 0952                new ChunkedEncodingWriteStream(this) :
 0953                new ContentLengthWriteStream(this, request.Content.Headers.ContentLength.GetValueOrDefault());
 0954            return requestContentStream;
 0955        }
 956
 957        private CancellationTokenRegistration RegisterCancellation(CancellationToken cancellationToken)
 0958        {
 959            // Cancellation design:
 960            // - We register with the SendAsync CancellationToken for the duration of the SendAsync operation.
 961            // - We register with the Read/Write/CopyToAsync methods on the response stream for each such individual ope
 962            // - The registration disposes of the connection, tearing it down and causing any pending operations to wake
 963            // - Because such a tear down can result in a variety of different exception types, we check for a cancellat
 964            //   request and prioritize that over other exceptions, wrapping the actual exception as an inner of an OCE.
 0965            return cancellationToken.Register(static s =>
 0966            {
 0967                var connection = (HttpConnection)s!;
 0968                if (NetEventSource.Log.IsEnabled()) connection.Trace("Cancellation requested. Disposing of the connectio
 0969                connection.Dispose();
 0970            }, this);
 0971        }
 972
 973        private async ValueTask SendRequestContentAsync(HttpRequestMessage request, HttpContentWriteStream stream, bool 
 0974        {
 0975            Debug.Assert(stream.BytesWritten == 0);
 0976            if (HttpTelemetry.Log.IsEnabled()) HttpTelemetry.Log.RequestContentStart();
 977
 978            // Copy all of the data to the server.
 0979            if (async)
 0980            {
 0981                await request.Content!.CopyToAsync(stream, _transportContext, cancellationToken).ConfigureAwait(false);
 0982            }
 983            else
 0984            {
 0985                request.Content!.CopyTo(stream, _transportContext, cancellationToken);
 0986            }
 987
 988            // Finish the content; with a chunked upload, this includes writing the terminating chunk.
 0989            await stream.FinishAsync(async).ConfigureAwait(false);
 990
 991            // Flush any content that might still be buffered.
 0992            await FlushAsync(async).ConfigureAwait(false);
 993
 0994            if (HttpTelemetry.Log.IsEnabled()) HttpTelemetry.Log.RequestContentStop(stream.BytesWritten);
 995
 0996            if (NetEventSource.Log.IsEnabled()) Trace("Finished sending request content.");
 0997        }
 998
 999        private async Task SendRequestContentWithExpect100ContinueAsync(
 1000            HttpRequestMessage request, Task<bool> allowExpect100ToContinueTask,
 1001            HttpContentWriteStream stream, Timer expect100Timer, bool async, CancellationToken cancellationToken)
 01002        {
 1003            // Wait until we receive a trigger notification that it's ok to continue sending content.
 1004            // This will come either when the timer fires or when we receive a response status line from the server.
 01005            bool sendRequestContent = await allowExpect100ToContinueTask.ConfigureAwait(false);
 1006
 1007            // Clean up the timer; it's no longer needed.
 01008            expect100Timer.Dispose();
 1009
 1010            // Send the content if we're supposed to.  Otherwise, we're done.
 01011            if (sendRequestContent)
 01012            {
 01013                if (NetEventSource.Log.IsEnabled()) Trace($"Sending request content for Expect: 100-continue.");
 1014                try
 01015                {
 01016                    await SendRequestContentAsync(request, stream, async, cancellationToken).ConfigureAwait(false);
 01017                }
 01018                catch
 01019                {
 1020                    // Tear down the connection if called from the timer thread because caller's thread will wait for se
 1021                    // or till HttpClient.Timeout tear the connection itself.
 01022                    Dispose();
 01023                    throw;
 1024                }
 01025            }
 1026            else
 01027            {
 01028                if (NetEventSource.Log.IsEnabled()) Trace($"Canceling request content for Expect: 100-continue.");
 01029            }
 01030        }
 1031
 1032        private bool ParseStatusLine(HttpResponseMessage response)
 01033        {
 01034            Span<byte> buffer = _readBuffer.ActiveSpan;
 1035
 01036            int lineFeedIndex = buffer.IndexOf((byte)'\n');
 01037            if (lineFeedIndex >= 0)
 01038            {
 01039                int bytesConsumed = lineFeedIndex + 1;
 01040                _readBuffer.Discard(bytesConsumed);
 01041                _allowedReadLineBytes -= bytesConsumed;
 1042
 01043                int carriageReturnIndex = lineFeedIndex - 1;
 01044                int length = (uint)carriageReturnIndex < (uint)buffer.Length && buffer[carriageReturnIndex] == '\r'
 01045                    ? carriageReturnIndex
 01046                    : lineFeedIndex;
 1047
 01048                ParseStatusLineCore(buffer.Slice(0, length), response);
 01049                return true;
 1050            }
 1051            else
 01052            {
 01053                if (_allowedReadLineBytes <= buffer.Length)
 01054                {
 01055                    ThrowExceededAllowedReadLineBytes();
 01056                }
 01057                return false;
 1058            }
 01059        }
 1060
 1061        private static void ParseStatusLineCore(Span<byte> line, HttpResponseMessage response)
 01062        {
 1063            // We sent the request version as either 1.0 or 1.1.
 1064            // We expect a response version of the form 1.X, where X is a single digit as per RFC.
 1065
 1066            // Validate the beginning of the status line and set the response version.
 1067            const int MinStatusLineLength = 12; // "HTTP/1.x 123"
 01068            if (line.Length < MinStatusLineLength || line[8] != ' ')
 01069            {
 01070                throw new HttpRequestException(HttpRequestError.InvalidResponse, SR.Format(SR.net_http_invalid_response_
 1071            }
 1072
 01073            ulong first8Bytes = BitConverter.ToUInt64(line);
 01074            if (first8Bytes == s_http11Bytes)
 01075            {
 01076                response.SetVersionWithoutValidation(HttpVersion.Version11);
 01077            }
 01078            else if (first8Bytes == s_http10Bytes)
 01079            {
 01080                response.SetVersionWithoutValidation(HttpVersion.Version10);
 01081            }
 1082            else
 01083            {
 01084                byte minorVersion = line[7];
 01085                if (IsDigit(minorVersion) && line.StartsWith("HTTP/1."u8))
 01086                {
 01087                    response.SetVersionWithoutValidation(new Version(1, minorVersion - '0'));
 01088                }
 1089                else
 01090                {
 01091                    throw new HttpRequestException(HttpRequestError.InvalidResponse, SR.Format(SR.net_http_invalid_respo
 1092                }
 01093            }
 1094
 1095            // Set the status code
 01096            byte status1 = line[9], status2 = line[10], status3 = line[11];
 01097            if (!IsDigit(status1) || !IsDigit(status2) || !IsDigit(status3))
 01098            {
 01099                throw new HttpRequestException(HttpRequestError.InvalidResponse, SR.Format(SR.net_http_invalid_response_
 1100            }
 01101            response.SetStatusCodeWithoutValidation((HttpStatusCode)(100 * (status1 - '0') + 10 * (status2 - '0') + (sta
 1102
 1103            // Parse (optional) reason phrase
 01104            if (line.Length == MinStatusLineLength)
 01105            {
 01106                response.SetReasonPhraseWithoutValidation(string.Empty);
 01107            }
 01108            else if (line[MinStatusLineLength] == ' ')
 01109            {
 01110                ReadOnlySpan<byte> reasonBytes = line.Slice(MinStatusLineLength + 1);
 01111                string? knownReasonPhrase = HttpStatusDescription.Get(response.StatusCode);
 01112                if (knownReasonPhrase != null && Ascii.Equals(reasonBytes, knownReasonPhrase))
 01113                {
 01114                    response.SetReasonPhraseWithoutValidation(knownReasonPhrase);
 01115                }
 1116                else
 01117                {
 1118                    try
 01119                    {
 01120                        response.ReasonPhrase = HttpRuleParser.DefaultHttpEncoding.GetString(reasonBytes);
 01121                    }
 01122                    catch (FormatException formatEx)
 01123                    {
 01124                        throw new HttpRequestException(HttpRequestError.InvalidResponse, SR.Format(SR.net_http_invalid_r
 1125                    }
 01126                }
 01127            }
 1128            else
 01129            {
 01130                throw new HttpRequestException(HttpRequestError.InvalidResponse, SR.Format(SR.net_http_invalid_response_
 1131            }
 01132        }
 1133
 1134        private bool ParseHeaders(HttpResponseMessage? response, bool isFromTrailer)
 01135        {
 01136            Span<byte> buffer = _readBuffer.ActiveSpan;
 1137
 01138            (bool finished, int bytesConsumed) = ParseHeadersCore(buffer, response, isFromTrailer);
 1139
 01140            int bytesScanned = finished ? bytesConsumed : buffer.Length;
 01141            if (_allowedReadLineBytes < bytesScanned)
 01142            {
 01143                ThrowExceededAllowedReadLineBytes();
 01144            }
 1145
 01146            _readBuffer.Discard(bytesConsumed);
 01147            _allowedReadLineBytes -= bytesConsumed;
 01148            Debug.Assert(_allowedReadLineBytes >= 0);
 1149
 01150            return finished;
 01151        }
 1152
 1153        private (bool finished, int bytesConsumed) ParseHeadersCore(Span<byte> buffer, HttpResponseMessage? response, bo
 01154        {
 01155            int originalBufferLength = buffer.Length;
 1156
 01157            while (true)
 01158            {
 01159                int colIdx = buffer.IndexOfAny((byte)':', (byte)'\n');
 01160                if (colIdx < 0)
 01161                {
 01162                    return (finished: false, bytesConsumed: originalBufferLength - buffer.Length);
 1163                }
 1164
 01165                if (buffer[colIdx] == '\n')
 01166                {
 01167                    if ((colIdx == 1 && buffer[0] == '\r') || colIdx == 0)
 01168                    {
 01169                        return (finished: true, bytesConsumed: originalBufferLength - buffer.Length + colIdx + 1);
 1170                    }
 1171
 01172                    ThrowForInvalidHeaderLine(buffer, colIdx);
 01173                }
 1174
 01175                int valueStartIdx = colIdx + 1;
 01176                if ((uint)valueStartIdx >= (uint)buffer.Length)
 01177                {
 01178                    return (finished: false, bytesConsumed: originalBufferLength - buffer.Length);
 1179                }
 1180
 1181                // Iterate over the value and handle any line folds (new lines followed by SP/HTAB).
 1182                // valueIterator refers to the remainder of the buffer that we can still scan for new lines.
 01183                Span<byte> valueIterator = buffer.Slice(valueStartIdx);
 1184
 01185                while (true)
 01186                {
 01187                    int lfIdx = valueIterator.IndexOf((byte)'\n');
 01188                    if ((uint)lfIdx >= (uint)valueIterator.Length)
 01189                    {
 01190                        return (finished: false, bytesConsumed: originalBufferLength - buffer.Length);
 1191                    }
 1192
 01193                    int crIdx = lfIdx - 1;
 01194                    int crOrLfIdx = (uint)crIdx < (uint)valueIterator.Length && valueIterator[crIdx] == '\r'
 01195                        ? crIdx
 01196                        : lfIdx;
 1197
 01198                    int spIdx = lfIdx + 1;
 01199                    if ((uint)spIdx >= (uint)valueIterator.Length)
 01200                    {
 01201                        return (finished: false, bytesConsumed: originalBufferLength - buffer.Length);
 1202                    }
 1203
 01204                    if (valueIterator[spIdx] is not (byte)'\t' and not (byte)' ')
 01205                    {
 1206                        // Found the end of the header value.
 1207
 01208                        if (response is not null)
 01209                        {
 01210                            ReadOnlySpan<byte> headerName = buffer.Slice(0, valueStartIdx - 1);
 01211                            ReadOnlySpan<byte> headerValue = buffer.Slice(valueStartIdx, buffer.Length - valueIterator.L
 01212                            AddResponseHeader(headerName, headerValue, response, isFromTrailer);
 01213                        }
 1214
 01215                        buffer = buffer.Slice(buffer.Length - valueIterator.Length + spIdx);
 01216                        break;
 1217                    }
 1218
 1219                    // Found an obs-fold (CRLFHT/CRLFSP).
 1220                    // Replace the CRLF with SPSP and keep looking for the final newline.
 01221                    valueIterator[crOrLfIdx] = (byte)' ';
 01222                    valueIterator[lfIdx] = (byte)' ';
 1223
 01224                    valueIterator = valueIterator.Slice(spIdx + 1);
 01225                }
 01226            }
 1227
 1228            static void ThrowForInvalidHeaderLine(ReadOnlySpan<byte> buffer, int newLineIndex) =>
 01229                throw new HttpRequestException(HttpRequestError.InvalidResponse, SR.Format(SR.net_http_invalid_response_
 01230        }
 1231
 1232        private void AddResponseHeader(ReadOnlySpan<byte> name, ReadOnlySpan<byte> value, HttpResponseMessage response, 
 01233        {
 1234            // Skip trailing whitespace and check for empty length.
 01235            while (true)
 01236            {
 01237                int spIdx = name.Length - 1;
 1238
 01239                if ((uint)spIdx < (uint)name.Length)
 01240                {
 01241                    if (name[spIdx] != ' ')
 01242                    {
 1243                        // hot path
 01244                        break;
 1245                    }
 1246
 01247                    name = name.Slice(0, spIdx);
 01248                }
 1249                else
 01250                {
 01251                    ThrowForEmptyHeaderName();
 01252                }
 01253            }
 1254
 1255            // Skip leading OWS for value.
 1256            // hot path: loop body runs only once.
 01257            while (value.Length != 0 && value[0] is (byte)' ' or (byte)'\t')
 01258            {
 01259                value = value.Slice(1);
 01260            }
 1261
 1262            // Skip trailing OWS for value.
 01263            while (true)
 01264            {
 01265                int spIdx = value.Length - 1;
 1266
 01267                if ((uint)spIdx >= (uint)value.Length || !(value[spIdx] is (byte)' ' or (byte)'\t'))
 01268                {
 1269                    // hot path
 01270                    break;
 1271                }
 1272
 01273                value = value.Slice(0, spIdx);
 01274            }
 1275
 01276            if (!HeaderDescriptor.TryGet(name, out HeaderDescriptor descriptor))
 01277            {
 01278                ThrowForInvalidHeaderName(name);
 01279            }
 1280
 01281            Encoding? valueEncoding = _pool.Settings._responseHeaderEncodingSelector?.Invoke(descriptor.Name, _currentRe
 1282
 01283            HttpHeaderType headerType = descriptor.HeaderType;
 1284
 1285            // Request headers returned on the response must be treated as custom headers.
 01286            if ((headerType & HttpHeaderType.Request) != 0)
 01287            {
 01288                descriptor = descriptor.AsCustomHeader();
 01289            }
 1290
 1291            string headerValue;
 1292            HttpHeaders headers;
 1293
 01294            if (isFromTrailer)
 01295            {
 01296                if ((headerType & HttpHeaderType.NonTrailing) != 0)
 01297                {
 1298                    // Disallowed trailer fields.
 1299                    // A recipient MUST ignore fields that are forbidden to be sent in a trailer.
 01300                    return;
 1301                }
 1302
 01303                headerValue = descriptor.GetHeaderValue(value, valueEncoding);
 01304                headers = response.TrailingHeaders;
 01305            }
 01306            else if ((headerType & HttpHeaderType.Content) != 0)
 01307            {
 01308                headerValue = descriptor.GetHeaderValue(value, valueEncoding);
 01309                headers = response.Content!.Headers;
 01310            }
 1311            else
 01312            {
 01313                headerValue = GetResponseHeaderValueWithCaching(descriptor, value, valueEncoding);
 01314                headers = response.Headers;
 1315
 01316                if (descriptor.Equals(KnownHeaders.KeepAlive))
 01317                {
 1318                    // We are intentionally going against RFC to honor the Keep-Alive header even if
 1319                    // we haven't received a Keep-Alive connection token to maximize compat with servers.
 01320                    ProcessKeepAliveHeader(headerValue);
 01321                }
 01322            }
 1323
 01324            bool added = headers.TryAddWithoutValidation(descriptor, headerValue);
 01325            Debug.Assert(added);
 1326
 1327            static void ThrowForEmptyHeaderName() =>
 01328                throw new HttpRequestException(HttpRequestError.InvalidResponse, SR.Format(SR.net_http_invalid_response_
 1329
 1330            static void ThrowForInvalidHeaderName(ReadOnlySpan<byte> name) =>
 01331                throw new HttpRequestException(HttpRequestError.InvalidResponse, SR.Format(SR.net_http_invalid_response_
 01332        }
 1333
 1334        private void ThrowExceededAllowedReadLineBytes() =>
 01335            throw new HttpRequestException(HttpRequestError.ConfigurationLimitExceeded, SR.Format(SR.net_http_response_h
 1336
 1337        private void ProcessKeepAliveHeader(string keepAlive)
 01338        {
 01339            var parsedValues = new UnvalidatedObjectCollection<NameValueHeaderValue>();
 1340
 01341            if (NameValueHeaderValue.GetNameValueListLength(keepAlive, 0, ',', parsedValues) == keepAlive.Length)
 01342            {
 01343                foreach (NameValueHeaderValue nameValue in parsedValues)
 01344                {
 1345                    // The HTTP/1.1 spec does not define any parameters for the Keep-Alive header, so we are using the d
 01346                    if (string.Equals(nameValue.Name, "timeout", StringComparison.OrdinalIgnoreCase))
 01347                    {
 01348                        if (!string.IsNullOrEmpty(nameValue.Value) &&
 01349                            HeaderUtilities.TryParseInt32(nameValue.Value, out int timeout) &&
 01350                            timeout >= 0)
 01351                        {
 1352                            // Some servers are very strict with closing the connection exactly at the timeout.
 1353                            // Avoid using the connection if it is about to exceed the timeout to avoid resulting reques
 1354                            const int OffsetSeconds = 1;
 1355
 01356                            if (timeout <= OffsetSeconds)
 01357                            {
 01358                                _connectionClose = true;
 01359                            }
 1360                            else
 01361                            {
 01362                                _keepAliveTimeoutSeconds = timeout - OffsetSeconds;
 01363                            }
 01364                        }
 01365                    }
 01366                    else if (string.Equals(nameValue.Name, "max", StringComparison.OrdinalIgnoreCase))
 01367                    {
 01368                        if (nameValue.Value == "0")
 01369                        {
 01370                            _connectionClose = true;
 01371                        }
 01372                    }
 01373                }
 01374            }
 01375        }
 1376
 1377        private void WriteToBuffer(ReadOnlySpan<byte> source)
 01378        {
 01379            Debug.Assert(source.Length <= _writeBuffer.AvailableLength);
 01380            source.CopyTo(_writeBuffer.AvailableSpan);
 01381            _writeBuffer.Commit(source.Length);
 01382        }
 1383
 1384        private void Write(ReadOnlySpan<byte> source)
 01385        {
 01386            int remaining = _writeBuffer.AvailableLength;
 1387
 01388            if (source.Length <= remaining)
 01389            {
 1390                // Fits in current write buffer.  Just copy and return.
 01391                WriteToBuffer(source);
 01392                return;
 1393            }
 1394
 01395            if (_writeBuffer.ActiveLength != 0)
 01396            {
 1397                // Fit what we can in the current write buffer and flush it.
 01398                WriteToBuffer(source.Slice(0, remaining));
 01399                source = source.Slice(remaining);
 01400                Flush();
 01401            }
 1402
 01403            if (source.Length >= _writeBuffer.Capacity)
 01404            {
 1405                // Large write.  No sense buffering this.  Write directly to stream.
 01406                WriteToStream(source);
 01407            }
 1408            else
 01409            {
 1410                // Copy remainder into buffer
 01411                WriteToBuffer(source);
 01412            }
 01413        }
 1414
 1415        private ValueTask WriteAsync(ReadOnlyMemory<byte> source)
 01416        {
 01417            int remaining = _writeBuffer.AvailableLength;
 1418
 01419            if (source.Length <= remaining)
 01420            {
 1421                // Fits in current write buffer.  Just copy and return.
 01422                WriteToBuffer(source.Span);
 01423                return default;
 1424            }
 1425
 01426            if (_writeBuffer.ActiveLength != 0)
 01427            {
 1428                // Fit what we can in the current write buffer and flush it.
 01429                WriteToBuffer(source.Span.Slice(0, remaining));
 01430                source = source.Slice(remaining);
 1431
 01432                ValueTask flushTask = FlushAsync(async: true);
 1433
 01434                if (flushTask.IsCompletedSuccessfully)
 01435                {
 01436                    flushTask.GetAwaiter().GetResult();
 1437
 01438                    if (source.Length <= _writeBuffer.Capacity)
 01439                    {
 01440                        WriteToBuffer(source.Span);
 01441                        return default;
 1442                    }
 1443
 1444                    // Fall-through to WriteToStreamAsync
 01445                }
 1446                else
 01447                {
 01448                    return AwaitFlushAndWriteAsync(flushTask, source);
 1449                }
 01450            }
 1451
 1452            // Large write.  No sense buffering this.  Write directly to stream.
 01453            return WriteToStreamAsync(source, async: true);
 1454
 1455            async ValueTask AwaitFlushAndWriteAsync(ValueTask flushTask, ReadOnlyMemory<byte> source)
 01456            {
 01457                await flushTask.ConfigureAwait(false);
 1458
 01459                if (source.Length <= _writeBuffer.Capacity)
 01460                {
 01461                    WriteToBuffer(source.Span);
 01462                }
 1463                else
 01464                {
 01465                    await WriteToStreamAsync(source, async: true).ConfigureAwait(false);
 01466                }
 01467            }
 01468        }
 1469
 1470        private void WriteWithoutBuffering(ReadOnlySpan<byte> source)
 01471        {
 01472            if (_writeBuffer.ActiveLength != 0)
 01473            {
 01474                if (source.Length <= _writeBuffer.AvailableLength)
 01475                {
 1476                    // There's something already in the write buffer, but the content
 1477                    // we're writing can also fit after it in the write buffer.  Copy
 1478                    // the content to the write buffer and then flush it, so that we
 1479                    // can do a single send rather than two.
 01480                    WriteToBuffer(source);
 01481                    Flush();
 01482                    return;
 1483                }
 1484
 1485                // There's data in the write buffer and the data we're writing doesn't fit after it.
 1486                // Do two writes, one to flush the buffer and then another to write the supplied content.
 01487                Flush();
 01488            }
 1489
 01490            WriteToStream(source);
 01491        }
 1492
 1493        private ValueTask WriteWithoutBufferingAsync(ReadOnlyMemory<byte> source, bool async)
 01494        {
 01495            if (_writeBuffer.ActiveLength == 0)
 01496            {
 1497                // There's nothing in the write buffer we need to flush.
 1498                // Just write the supplied data out to the stream.
 01499                return WriteToStreamAsync(source, async);
 1500            }
 1501
 01502            if (source.Length <= _writeBuffer.AvailableLength)
 01503            {
 1504                // There's something already in the write buffer, but the content
 1505                // we're writing can also fit after it in the write buffer.  Copy
 1506                // the content to the write buffer and then flush it, so that we
 1507                // can do a single send rather than two.
 01508                WriteToBuffer(source.Span);
 01509                return FlushAsync(async);
 1510            }
 1511
 1512            // There's data in the write buffer and the data we're writing doesn't fit after it.
 1513            // Do two writes, one to flush the buffer and then another to write the supplied content.
 01514            return FlushThenWriteWithoutBufferingAsync(source, async);
 01515        }
 1516
 1517        private async ValueTask FlushThenWriteWithoutBufferingAsync(ReadOnlyMemory<byte> source, bool async)
 01518        {
 01519            await FlushAsync(async).ConfigureAwait(false);
 01520            await WriteToStreamAsync(source, async).ConfigureAwait(false);
 01521        }
 1522
 1523        private unsafe ValueTask WriteHexInt32Async(int value, bool async)
 01524        {
 1525            // Try to format into our output buffer directly.
 01526            if (value.TryFormat(_writeBuffer.AvailableSpan, out int bytesWritten, "X"))
 01527            {
 01528                _writeBuffer.Commit(bytesWritten);
 01529                return default;
 1530            }
 1531
 1532            // If we don't have enough room, do it the slow way.
 01533            if (async)
 01534            {
 01535                Span<byte> temp = stackalloc byte[8]; // max length of Int32 as hex
 01536                bool formatted = value.TryFormat(temp, out bytesWritten, "X");
 01537                Debug.Assert(formatted);
 01538                return WriteAsync(temp.Slice(0, bytesWritten).ToArray());
 1539            }
 1540            else
 01541            {
 1542                // We should have enough capacity to write any hex-encoded int after flushing the buffer.
 01543                Debug.Assert(_writeBuffer.Capacity >= 8);
 1544
 01545                Flush();
 01546                return WriteHexInt32Async(value, async: false);
 1547            }
 01548        }
 1549
 1550        private void Flush()
 01551        {
 01552            ReadOnlySpan<byte> bytes = _writeBuffer.ActiveSpan;
 01553            if (bytes.Length > 0)
 01554            {
 01555                _writeBuffer.DiscardAll();
 01556                WriteToStream(bytes);
 01557            }
 01558        }
 1559
 1560        private ValueTask FlushAsync(bool async)
 01561        {
 01562            ReadOnlyMemory<byte> bytes = _writeBuffer.ActiveMemory;
 01563            if (bytes.Length > 0)
 01564            {
 01565                _writeBuffer.DiscardAll();
 01566                return WriteToStreamAsync(bytes, async);
 1567            }
 01568            return default;
 01569        }
 1570
 1571        private void WriteToStream(ReadOnlySpan<byte> source)
 01572        {
 01573            if (NetEventSource.Log.IsEnabled()) Trace($"Writing {source.Length} bytes.");
 01574            _stream.Write(source);
 01575        }
 1576
 1577        private ValueTask WriteToStreamAsync(ReadOnlyMemory<byte> source, bool async)
 01578        {
 01579            if (NetEventSource.Log.IsEnabled()) Trace($"Writing {source.Length} bytes.");
 1580
 01581            if (async)
 01582            {
 01583                return _stream.WriteAsync(source);
 1584            }
 1585            else
 01586            {
 01587                _stream.Write(source.Span);
 01588                return default;
 1589            }
 01590        }
 1591
 1592        private bool TryReadNextChunkedLine(out ReadOnlySpan<byte> line)
 01593        {
 01594            ReadOnlySpan<byte> buffer = _readBuffer.ActiveReadOnlySpan;
 1595
 1596            // Unlike the status line and headers, the chunked encoding grammar (RFC 9112 7.1)
 1597            // requires that each line be terminated by a CRLF. Interpreting a lone LF as a line
 1598            // terminator, or allowing a bare CR within the line, is not permitted here.
 01599            int index = buffer.IndexOfAny((byte)'\r', (byte)'\n');
 01600            if ((uint)index >= (uint)buffer.Length)
 01601            {
 1602                // We haven't found a CR or LF yet, so we don't have a complete line.
 01603                if (buffer.Length < MaxChunkBytesAllowed)
 01604                {
 01605                    line = default;
 01606                    return false;
 1607                }
 01608            }
 01609            else if (buffer[index] == '\n')
 01610            {
 1611                // We found an LF that is not preceded by a CR.
 01612                throw new HttpIOException(HttpRequestError.InvalidResponse, SR.net_http_invalid_response_chunk_line_endi
 1613            }
 1614            else
 01615            {
 1616                // We found a CR. It must be immediately followed by an LF.
 01617                int lineFeedIndex = index + 1;
 01618                if ((uint)lineFeedIndex < (uint)buffer.Length)
 01619                {
 01620                    if (buffer[lineFeedIndex] != '\n')
 01621                    {
 1622                        // We found a bare CR that is not part of a CRLF sequence.
 01623                        throw new HttpIOException(HttpRequestError.InvalidResponse, SR.net_http_invalid_response_chunk_l
 1624                    }
 1625
 01626                    int bytesConsumed = lineFeedIndex + 1;
 01627                    if (bytesConsumed <= MaxChunkBytesAllowed)
 01628                    {
 01629                        _readBuffer.Discard(bytesConsumed);
 1630
 01631                        line = buffer.Slice(0, index);
 01632                        return true;
 1633                    }
 01634                }
 01635                else if (buffer.Length < MaxChunkBytesAllowed)
 01636                {
 1637                    // We have the CR but haven't received the following byte yet.
 01638                    line = default;
 01639                    return false;
 1640                }
 01641            }
 1642
 1643            // We either didn't find a line terminator within the allowed number of bytes, or the
 1644            // line (including its CRLF) is longer than we're willing to buffer.
 01645            throw new HttpRequestException(SR.net_http_chunk_too_large);
 01646        }
 1647
 1648        // Does not throw on EOF. Also assumes there is no buffered data.
 1649        private async ValueTask InitialFillAsync(bool async)
 01650        {
 01651            Debug.Assert(!ReadAheadTaskHasStarted);
 01652            Debug.Assert(_readBuffer.AvailableLength == _readBuffer.Capacity);
 01653            Debug.Assert(_readBuffer.AvailableLength >= InitialReadBufferSize);
 1654
 01655            int bytesRead = async ?
 01656                await _stream.ReadAsync(_readBuffer.AvailableMemory).ConfigureAwait(false) :
 01657                _stream.Read(_readBuffer.AvailableSpan);
 1658
 01659            _readBuffer.Commit(bytesRead);
 1660
 01661            if (NetEventSource.Log.IsEnabled()) Trace($"Received {bytesRead} bytes.");
 01662        }
 1663
 1664        // Throws IOException on EOF.  This is only called when we expect more data.
 1665        private async ValueTask FillAsync(bool async)
 01666        {
 01667            Debug.Assert(_readAheadTask == default);
 1668
 01669            _readBuffer.EnsureAvailableSpace(1);
 1670
 01671            int bytesRead = async ?
 01672                await _stream.ReadAsync(_readBuffer.AvailableMemory).ConfigureAwait(false) :
 01673                _stream.Read(_readBuffer.AvailableSpan);
 1674
 01675            _readBuffer.Commit(bytesRead);
 1676
 01677            if (NetEventSource.Log.IsEnabled()) Trace($"Received {bytesRead} bytes.");
 01678            if (bytesRead == 0)
 01679            {
 01680                throw new HttpIOException(HttpRequestError.ResponseEnded, SR.net_http_invalid_response_premature_eof);
 1681            }
 01682        }
 1683
 1684        private ValueTask FillForHeadersAsync(bool async)
 01685        {
 1686            // If the start offset is 0, it means we haven't consumed any data since the last FillAsync.
 1687            // If so, read until we either find the next new line or we hit the MaxResponseHeadersLength limit.
 01688            return _readBuffer.ActiveStartOffset == 0
 01689                ? ReadUntilEndOfHeaderAsync(async)
 01690                : FillAsync(async);
 1691
 1692            // This method guarantees that the next call to ParseHeaders will consume at least one header.
 1693            // This is the slow path, but guarantees O(n) worst-case parsing complexity.
 1694            async ValueTask ReadUntilEndOfHeaderAsync(bool async)
 01695            {
 01696                int searchOffset = _readBuffer.ActiveLength;
 01697                if (searchOffset > 0)
 01698                {
 1699                    // The last character we've buffered could be a new line,
 1700                    // we just haven't checked the byte following it to see if it's a space or tab.
 01701                    searchOffset--;
 01702                }
 1703
 01704                while (true)
 01705                {
 01706                    await FillAsync(async).ConfigureAwait(false);
 01707                    Debug.Assert(_readBuffer.ActiveStartOffset == 0);
 01708                    Debug.Assert(_readBuffer.ActiveLength > searchOffset);
 1709
 1710                    // There's no need to search the whole buffer, only look through the new bytes we just read.
 01711                    if (TryFindEndOfLine(_readBuffer.ActiveReadOnlySpan.Slice(searchOffset), out int offset))
 01712                    {
 01713                        break;
 1714                    }
 1715
 01716                    searchOffset += offset;
 1717
 01718                    int readLength = _readBuffer.ActiveLength;
 01719                    if (searchOffset != readLength)
 01720                    {
 01721                        Debug.Assert(searchOffset == readLength - 1 && _readBuffer.ActiveReadOnlySpan[searchOffset] == '
 01722                        if (readLength <= 2)
 01723                        {
 1724                            // There are no headers - we start off with a new line.
 1725                            // This is reachable from ChunkedEncodingReadStream if the buffers allign just right and the
 01726                            break;
 1727                        }
 01728                    }
 1729
 01730                    if (readLength >= _allowedReadLineBytes)
 01731                    {
 01732                        ThrowExceededAllowedReadLineBytes();
 01733                    }
 01734                }
 1735
 1736                static bool TryFindEndOfLine(ReadOnlySpan<byte> buffer, out int searchOffset)
 01737                {
 01738                    Debug.Assert(buffer.Length > 0);
 1739
 01740                    int originalBufferLength = buffer.Length;
 1741
 01742                    while (true)
 01743                    {
 01744                        int newLineOffset = buffer.IndexOf((byte)'\n');
 01745                        if (newLineOffset < 0)
 01746                        {
 01747                            searchOffset = originalBufferLength;
 01748                            return false;
 1749                        }
 1750
 01751                        int tabOrSpaceIndex = newLineOffset + 1;
 01752                        if (tabOrSpaceIndex == buffer.Length)
 01753                        {
 1754                            // The new line is the last character, read again to make sure it doesn't continue with spac
 01755                            searchOffset = originalBufferLength - 1;
 01756                            return false;
 1757                        }
 1758
 01759                        if (buffer[tabOrSpaceIndex] is not (byte)'\t' and not (byte)' ')
 01760                        {
 01761                            searchOffset = 0;
 01762                            return true;
 1763                        }
 1764
 01765                        buffer = buffer.Slice(tabOrSpaceIndex + 1);
 01766                    }
 01767                }
 01768            }
 01769        }
 1770
 1771        private int ReadFromBuffer(Span<byte> buffer)
 01772        {
 01773            ReadOnlySpan<byte> available = _readBuffer.ActiveSpan;
 01774            int toCopy = Math.Min(available.Length, buffer.Length);
 1775
 01776            available.Slice(0, toCopy).CopyTo(buffer);
 01777            _readBuffer.Discard(toCopy);
 1778
 01779            return toCopy;
 01780        }
 1781
 1782        private int Read(Span<byte> destination)
 01783        {
 1784            // This is called when reading the response body.
 1785
 01786            if (_readBuffer.ActiveLength > 0)
 01787            {
 1788                // We have data in the read buffer.  Return it to the caller.
 01789                return ReadFromBuffer(destination);
 1790            }
 1791
 1792            // No data in read buffer.
 1793            // Do an unbuffered read directly against the underlying stream.
 01794            Debug.Assert(_readAheadTask == default, "Read ahead task should have been consumed as part of the headers.")
 01795            int count = _stream.Read(destination);
 01796            if (NetEventSource.Log.IsEnabled()) Trace($"Received {count} bytes.");
 01797            return count;
 01798        }
 1799
 1800        private ValueTask<int> ReadAsync(Memory<byte> destination)
 01801        {
 1802            // This is called when reading the response body.
 1803
 01804            if (_readBuffer.ActiveLength > 0)
 01805            {
 1806                // We have data in the read buffer.  Return it to the caller.
 01807                return new ValueTask<int>(ReadFromBuffer(destination.Span));
 1808            }
 1809
 1810            // No data in read buffer.
 1811            // Do an unbuffered read directly against the underlying stream.
 01812            Debug.Assert(_readAheadTask == default, "Read ahead task should have been consumed as part of the headers.")
 1813
 01814            return NetEventSource.Log.IsEnabled()
 01815                ? ReadAndLogBytesReadAsync(destination)
 01816                : _stream.ReadAsync(destination);
 1817
 1818            async ValueTask<int> ReadAndLogBytesReadAsync(Memory<byte> destination)
 01819            {
 01820                int count = await _stream.ReadAsync(destination).ConfigureAwait(false);
 01821                if (NetEventSource.Log.IsEnabled()) Trace($"Received {count} bytes.");
 01822                return count;
 01823            }
 01824        }
 1825
 1826        private int ReadBuffered(Span<byte> destination)
 01827        {
 1828            // This is called when reading the response body.
 1829
 01830            if (_readBuffer.ActiveLength == 0)
 01831            {
 1832                // Do a buffered read directly against the underlying stream.
 01833                Debug.Assert(_readAheadTask == default, "Read ahead task should have been consumed as part of the header
 1834
 01835                if (destination.Length == 0)
 01836                {
 01837                    return _stream.Read(Array.Empty<byte>());
 1838                }
 1839
 01840                Debug.Assert(_readBuffer.AvailableLength == _readBuffer.Capacity);
 01841                int bytesRead = _stream.Read(_readBuffer.AvailableSpan);
 01842                _readBuffer.Commit(bytesRead);
 1843
 01844                if (NetEventSource.Log.IsEnabled()) Trace($"Received {bytesRead} bytes.");
 01845            }
 1846
 1847            // Hand back as much data as we can fit.
 01848            return ReadFromBuffer(destination);
 01849        }
 1850
 1851        private ValueTask<int> ReadBufferedAsync(Memory<byte> destination)
 01852        {
 1853            // If the caller provided buffer, and thus the amount of data desired to be read,
 1854            // is larger than the internal buffer, there's no point going through the internal
 1855            // buffer, so just do an unbuffered read.
 1856            // Also avoid avoid using the internal buffer if the user requested a zero-byte read to allow
 1857            // underlying streams to efficiently handle such a read (e.g. SslStream defering buffer allocation).
 01858            return destination.Length >= _readBuffer.Capacity || destination.Length == 0 ?
 01859                ReadAsync(destination) :
 01860                ReadBufferedAsyncCore(destination);
 01861        }
 1862
 1863        [AsyncMethodBuilder(typeof(PoolingAsyncValueTaskMethodBuilder<>))]
 1864        [RuntimeAsyncMethodGeneration(false)]
 1865        private async ValueTask<int> ReadBufferedAsyncCore(Memory<byte> destination)
 01866        {
 1867            // This is called when reading the response body.
 1868
 01869            if (_readBuffer.ActiveLength == 0)
 01870            {
 1871                // Do a buffered read directly against the underlying stream.
 01872                Debug.Assert(_readAheadTask == default, "Read ahead task should have been consumed as part of the header
 1873
 01874                Debug.Assert(_readBuffer.AvailableLength == _readBuffer.Capacity);
 01875                int bytesRead = await _stream.ReadAsync(_readBuffer.AvailableMemory).ConfigureAwait(false);
 01876                _readBuffer.Commit(bytesRead);
 1877
 01878                if (NetEventSource.Log.IsEnabled()) Trace($"Received {bytesRead} bytes.");
 01879            }
 1880
 1881            // Hand back as much data as we can fit.
 01882            return ReadFromBuffer(destination.Span);
 01883        }
 1884
 1885        private ValueTask CopyFromBufferAsync(Stream destination, bool async, int count, CancellationToken cancellationT
 01886        {
 01887            Debug.Assert(count <= _readBuffer.ActiveLength);
 1888
 01889            if (NetEventSource.Log.IsEnabled()) Trace($"Copying {count} bytes to stream.");
 1890
 01891            ReadOnlyMemory<byte> source = _readBuffer.ActiveMemory.Slice(0, count);
 01892            _readBuffer.Discard(count);
 1893
 01894            if (async)
 01895            {
 01896                return destination.WriteAsync(source, cancellationToken);
 1897            }
 1898            else
 01899            {
 01900                destination.Write(source.Span);
 01901                return default;
 1902            }
 01903        }
 1904
 1905        private Task CopyToUntilEofAsync(Stream destination, bool async, int bufferSize, CancellationToken cancellationT
 01906        {
 01907            Debug.Assert(destination != null);
 1908
 01909            if (_readBuffer.ActiveLength > 0)
 01910            {
 01911                return CopyToUntilEofWithExistingBufferedDataAsync(destination, async, bufferSize, cancellationToken);
 1912            }
 1913
 01914            if (async)
 01915            {
 01916                return _stream.CopyToAsync(destination, bufferSize, cancellationToken);
 1917            }
 1918
 01919            _stream.CopyTo(destination, bufferSize);
 01920            return Task.CompletedTask;
 01921        }
 1922
 1923        private async Task CopyToUntilEofWithExistingBufferedDataAsync(Stream destination, bool async, int bufferSize, C
 01924        {
 01925            int remaining = _readBuffer.ActiveLength;
 01926            Debug.Assert(remaining > 0);
 1927
 01928            await CopyFromBufferAsync(destination, async, remaining, cancellationToken).ConfigureAwait(false);
 1929
 01930            if (async)
 01931            {
 01932                await _stream.CopyToAsync(destination, bufferSize, cancellationToken).ConfigureAwait(false);
 01933            }
 1934            else
 01935            {
 01936                _stream.CopyTo(destination, bufferSize);
 01937            }
 01938        }
 1939
 1940        // Copy *exactly* [length] bytes into destination; throws on end of stream.
 1941        private async Task CopyToContentLengthAsync(Stream destination, bool async, ulong length, int bufferSize, Cancel
 01942        {
 01943            Debug.Assert(destination != null);
 01944            Debug.Assert(length > 0);
 1945
 1946            // Copy any data left in the connection's buffer to the destination.
 01947            int remaining = _readBuffer.ActiveLength;
 01948            if (remaining > 0)
 01949            {
 01950                if ((ulong)remaining > length)
 01951                {
 01952                    remaining = (int)length;
 01953                }
 01954                await CopyFromBufferAsync(destination, async, remaining, cancellationToken).ConfigureAwait(false);
 1955
 01956                length -= (ulong)remaining;
 01957                if (length == 0)
 01958                {
 01959                    return;
 1960                }
 1961
 01962                Debug.Assert(_readBuffer.ActiveLength == 0, "HttpConnection's buffer should have been empty.");
 01963            }
 1964
 1965            // Repeatedly read into HttpConnection's buffer and write that buffer to the destination
 1966            // stream. If after doing so, we find that we filled the whole connection's buffer (which
 1967            // is sized mainly for HTTP headers rather than large payloads), grow the connection's
 1968            // read buffer to the requested buffer size to use for the remainder of the operation. We
 1969            // use a temporary buffer from the ArrayPool so that the connection doesn't hog large
 1970            // buffers from the pool for extended durations, especially if it's going to sit in the
 1971            // connection pool for a prolonged period.
 01972            byte[]? origReadBuffer = null;
 1973            try
 01974            {
 01975                while (true)
 01976                {
 01977                    await FillAsync(async).ConfigureAwait(false);
 1978
 01979                    remaining = (int)Math.Min((ulong)_readBuffer.ActiveLength, length);
 01980                    await CopyFromBufferAsync(destination, async, remaining, cancellationToken).ConfigureAwait(false);
 1981
 01982                    length -= (ulong)remaining;
 01983                    if (length == 0)
 01984                    {
 01985                        return;
 1986                    }
 1987
 1988                    // If we haven't yet grown the buffer (if we previously grew it, then it's sufficiently large), and
 1989                    // if we filled the read buffer while doing the last read (which is at least one indication that the
 1990                    // data arrival rate is fast enough to warrant a larger buffer), and if the buffer size we'd want is
 1991                    // larger than the one we already have, then grow the connection's read buffer to that size.
 01992                    if (origReadBuffer is null)
 01993                    {
 01994                        int currentCapacity = _readBuffer.Capacity;
 01995                        if (remaining == currentCapacity)
 01996                        {
 01997                            int desiredBufferSize = (int)Math.Min((ulong)bufferSize, length);
 01998                            if (desiredBufferSize > currentCapacity)
 01999                            {
 02000                                origReadBuffer = _readBuffer.DangerousGetUnderlyingBuffer();
 02001                                byte[] pooledBuffer = ArrayPool<byte>.Shared.Rent(desiredBufferSize);
 02002                                _readBuffer = new ArrayBuffer(pooledBuffer);
 02003                            }
 02004                        }
 02005                    }
 02006                }
 2007            }
 2008            finally
 02009            {
 02010                if (origReadBuffer is not null)
 02011                {
 02012                    Debug.Assert(origReadBuffer.Length > 0);
 2013
 2014                    // We don't care how much remaining data there was, just if there was any.
 2015                    // Subsequent code is going to check whether the receive buffer is empty
 2016                    // and then force the connection closed if it's not.
 02017                    bool anyDataAvailable = _readBuffer.ActiveLength > 0;
 2018
 02019                    byte[] pooledBuffer = _readBuffer.DangerousGetUnderlyingBuffer();
 02020                    _readBuffer = new ArrayBuffer(origReadBuffer);
 02021                    ArrayPool<byte>.Shared.Return(pooledBuffer);
 2022
 02023                    if (anyDataAvailable)
 02024                    {
 02025                        _readBuffer.Commit(1);
 02026                    }
 02027                }
 02028            }
 02029        }
 2030
 2031        internal void Acquire()
 02032        {
 02033            Debug.Assert(_currentRequest == null);
 02034            Debug.Assert(!_inUse);
 2035
 02036            _inUse = true;
 02037        }
 2038
 2039        internal void Release()
 02040        {
 02041            Debug.Assert(_inUse);
 2042
 02043            _inUse = false;
 2044
 2045            // If the last request already completed (because the response had no content), return the connection to the
 2046            // Otherwise, it will be returned when the response has been consumed and CompleteResponse below is called.
 02047            if (_currentRequest == null)
 02048            {
 02049                ReturnConnectionToPool();
 02050            }
 02051        }
 2052
 2053        /// <summary>
 2054        /// Detach the connection from the pool, so it is no longer counted against the connection limit.
 2055        /// This is used when we are creating a replacement connection for NT auth challenges.
 2056        /// </summary>
 2057        internal void DetachFromPool()
 02058        {
 02059            Debug.Assert(_inUse);
 2060
 02061            _detachedFromPool = true;
 02062        }
 2063
 2064        private void CompleteResponse()
 02065        {
 02066            Debug.Assert(_currentRequest != null, "Expected the connection to be associated with a request.");
 02067            Debug.Assert(_writeBuffer.ActiveLength == 0, "Everything in write buffer should have been flushed.");
 2068
 2069            // Disassociate the connection from a request.
 02070            _currentRequest = null;
 2071
 2072            // If we have extraneous data in the read buffer, don't reuse the connection;
 2073            // otherwise we'd interpret this as part of the next response. Plus, we may
 2074            // have been using a temporary buffer to read this erroneous data, and thus
 2075            // may not even have it any more.
 02076            if (_readBuffer.ActiveLength > 0)
 02077            {
 02078                if (NetEventSource.Log.IsEnabled())
 02079                {
 02080                    Trace("Unexpected data on connection after response read.");
 02081                }
 2082
 02083                _readBuffer.DiscardAll();
 02084                _connectionClose = true;
 02085            }
 2086
 2087            // If the connection is no longer in use (i.e. for NT authentication), then we can
 2088            // return it to the pool now; otherwise, it will be returned by the Release method later.
 2089            // The cancellation logic in HTTP/1.1 response stream reading methods is prone to race conditions
 2090            // where CancellationTokenRegistration callbacks may dispose the connection without the disposal
 2091            // leading to an actual cancellation of the response reading methods by an OperationCanceledException.
 2092            // To guard against these cases, it is necessary to check if the connection is disposed before
 2093            // attempting to return it to the pool.
 02094            if (!_inUse && !_disposed)
 02095            {
 02096                ReturnConnectionToPool();
 02097            }
 02098        }
 2099
 2100        public async ValueTask DrainResponseAsync(HttpResponseMessage response, CancellationToken cancellationToken)
 02101        {
 02102            Debug.Assert(_inUse);
 2103
 02104            if (_connectionClose)
 02105            {
 02106                throw new HttpRequestException(HttpRequestError.UserAuthenticationError, SR.net_http_authconnectionfailu
 2107            }
 2108
 02109            Debug.Assert(response.Content != null);
 02110            Stream stream = response.Content.ReadAsStream(cancellationToken);
 02111            HttpContentReadStream? responseStream = stream as HttpContentReadStream;
 2112
 02113            Debug.Assert(responseStream != null || stream is EmptyReadStream);
 2114
 02115            if (responseStream != null && responseStream.NeedsDrain)
 02116            {
 02117                Debug.Assert(response.RequestMessage == _currentRequest);
 2118
 02119                if (!await responseStream.DrainAsync(_pool.Settings._maxResponseDrainSize).ConfigureAwait(false) ||
 02120                    _connectionClose)       // Draining may have set this
 02121                {
 02122                    throw new HttpRequestException(HttpRequestError.UserAuthenticationError, SR.net_http_authconnectionf
 2123                }
 02124            }
 2125
 02126            Debug.Assert(_currentRequest == null);
 2127
 02128            response.Dispose();
 02129        }
 2130
 2131        private void ReturnConnectionToPool()
 02132        {
 02133            Debug.Assert(!_disposed, "Connection should not be disposed.");
 02134            Debug.Assert(_currentRequest == null, "Connection should no longer be associated with a request.");
 02135            Debug.Assert(_readAheadTask == default, "Expected a previous initial read to already be consumed.");
 02136            Debug.Assert(_readAheadTaskStatus == ReadAheadTask_NotStarted, "Expected SendAsync to reset the read-ahead t
 02137            Debug.Assert(_readBuffer.ActiveLength == 0, "Unexpected data in connection read buffer.");
 2138
 2139            // If we decided not to reuse the connection (either because the server sent Connection: close,
 2140            // or there was some other problem while processing the request that makes the connection unusable),
 2141            // don't put the connection back in the pool.
 02142            if (_connectionClose)
 02143            {
 02144                if (NetEventSource.Log.IsEnabled())
 02145                {
 02146                    Trace("Connection will not be reused.");
 02147                }
 2148
 2149                // We're not putting the connection back in the pool. Dispose it.
 02150                Dispose();
 02151            }
 2152            else
 02153            {
 02154                Debug.Assert(!_detachedFromPool, "Should not be detached from pool unless _connectionClose is true");
 2155
 2156                // Put connection back in the pool.
 02157                _pool.RecycleHttp11Connection(this);
 02158            }
 02159        }
 2160
 02161        public sealed override string ToString() => $"{nameof(HttpConnection)}({_pool})"; // Description for diagnostic 
 2162
 2163        public sealed override void Trace(string message, [CallerMemberName] string? memberName = null) =>
 02164            NetEventSource.Log.HandlerMessage(
 02165                _pool?.GetHashCode() ?? 0,           // pool ID
 02166                GetHashCode(),                       // connection ID
 02167                _currentRequest?.GetHashCode() ?? 0, // request ID
 02168                memberName,                          // method name
 02169                message);                            // message
 2170    }
 2171}
 2172

https://raw.githubusercontent.com/dotnet/runtime/811a7eabb75c42db53440e8ba3f60c07511cfd1f/src/libraries/System.Net.Http/src/System/Net/Http/SocketsHttpHandler/HttpContentReadStream.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.Diagnostics;
 5using System.Threading;
 6using System.Threading.Tasks;
 7
 8namespace System.Net.Http
 9{
 10    internal sealed partial class HttpConnection
 11    {
 12        internal abstract class HttpContentReadStream : HttpContentStream
 13        {
 14            private bool _disposed;
 15
 016            public HttpContentReadStream(HttpConnection connection) : base(connection)
 017            {
 018            }
 19
 020            public sealed override bool CanRead => !_disposed;
 021            public sealed override bool CanWrite => false;
 22
 023            public sealed override void Write(ReadOnlySpan<byte> buffer) => throw new NotSupportedException(SR.net_http_
 24
 025            public sealed override ValueTask WriteAsync(ReadOnlyMemory<byte> destination, CancellationToken cancellation
 26
 027            public virtual bool NeedsDrain => false;
 28
 029            protected bool IsDisposed => _disposed;
 30
 31            protected bool CanReadFromConnection
 32            {
 33                get
 034                {
 35                    // _connection == null typically means that we have finished reading the response.
 36                    // Cancellation may lead to a state where a disposed _connection is not null.
 037                    HttpConnection? connection = _connection;
 038                    return connection != null && !connection._disposed;
 039                }
 40            }
 41
 42            public virtual ValueTask<bool> DrainAsync(int maxDrainBytes)
 043            {
 044                Debug.Fail($"DrainAsync should not be called for this response stream: {GetType()}");
 45                return new ValueTask<bool>(false);
 46            }
 47
 48            protected override void Dispose(bool disposing)
 049            {
 50                // Only attempt draining if we haven't started draining due to disposal; otherwise
 51                // multiple calls to Dispose (which happens frequently when someone disposes of the
 52                // response stream and response content) will kick off multiple concurrent draining
 53                // operations. Also don't delegate to the base if Dispose has already been called,
 54                // as doing so will end up disposing of the connection before we're done draining.
 055                if (Interlocked.Exchange(ref _disposed, true))
 056                {
 057                    return;
 58                }
 59
 060                if (disposing && NeedsDrain)
 061                {
 62                    // Start the asynchronous drain.
 63                    // It may complete synchronously, in which case the connection will be put back in the pool synchron
 64                    // Skip the call to base.Dispose -- it will be deferred until DrainOnDisposeAsync finishes.
 065                    _ = DrainOnDisposeAsync();
 066                    return;
 67                }
 68
 069                base.Dispose(disposing);
 070            }
 71
 72            private async Task DrainOnDisposeAsync()
 073            {
 074                HttpConnection? connection = _connection;        // Will be null after drain succeeds
 075                Debug.Assert(connection != null);
 76                try
 077                {
 078                    bool drained = await DrainAsync(connection._pool.Settings._maxResponseDrainSize).ConfigureAwait(fals
 79
 080                    if (NetEventSource.Log.IsEnabled())
 081                    {
 082                        connection.Trace(drained ?
 083                            "Connection drain succeeded" :
 084                            $"Connection drain failed when MaxResponseDrainSize={connection._pool.Settings._maxResponseD
 085                    }
 086                }
 087                catch (Exception e)
 088                {
 089                    if (NetEventSource.Log.IsEnabled())
 090                    {
 091                        connection.Trace($"Connection drain failed due to exception: {e}");
 092                    }
 93
 94                    // Eat any exceptions and just Dispose.
 095                }
 96
 097                base.Dispose(true);
 098            }
 99        }
 100    }
 101}
 102

https://raw.githubusercontent.com/dotnet/runtime/811a7eabb75c42db53440e8ba3f60c07511cfd1f/src/libraries/System.Net.Http/src/System/Net/Http/SocketsHttpHandler/HttpContentWriteStream.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.Diagnostics;
 5using System.IO;
 6using System.Threading;
 7using System.Threading.Tasks;
 8
 9namespace System.Net.Http
 10{
 11    internal sealed partial class HttpConnection : IDisposable
 12    {
 13        private abstract class HttpContentWriteStream : HttpContentStream
 14        {
 015            public long BytesWritten { get; protected set; }
 16
 017            public HttpContentWriteStream(HttpConnection connection) : base(connection) =>
 018                Debug.Assert(connection != null);
 19
 020            public sealed override bool CanRead => false;
 021            public sealed override bool CanWrite => _connection != null;
 22
 23            public sealed override void Flush() =>
 024                _connection?.Flush();
 25
 26            public sealed override Task FlushAsync(CancellationToken ignored)
 027            {
 028                HttpConnection? connection = _connection;
 029                return connection != null ?
 030                    connection.FlushAsync(async: true).AsTask() :
 031                    default!;
 032            }
 33
 034            public sealed override int Read(Span<byte> buffer) => throw new NotSupportedException();
 35
 036            public sealed override ValueTask<int> ReadAsync(Memory<byte> buffer, CancellationToken cancellationToken) =>
 37
 038            public sealed override Task CopyToAsync(Stream destination, int bufferSize, CancellationToken cancellationTo
 39
 40            public abstract Task FinishAsync(bool async);
 41        }
 42    }
 43}
 44

https://raw.githubusercontent.com/dotnet/runtime/811a7eabb75c42db53440e8ba3f60c07511cfd1f/src/libraries/System.Net.Http/src/System/Net/Http/SocketsHttpHandler/RawConnectionStream.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.IO;
 5using System.Runtime.CompilerServices;
 6using System.Runtime.ExceptionServices;
 7using System.Threading;
 8using System.Threading.Tasks;
 9
 10namespace System.Net.Http
 11{
 12    internal sealed partial class HttpConnection : IDisposable
 13    {
 14        private sealed class RawConnectionStream : HttpContentStream
 15        {
 016            public RawConnectionStream(HttpConnection connection) : base(connection)
 017            {
 018                if (NetEventSource.Log.IsEnabled()) NetEventSource.Info(this);
 019            }
 20
 021            public sealed override bool CanRead => _connection != null;
 022            public sealed override bool CanWrite => _connection != null;
 23
 24            public override int Read(Span<byte> buffer)
 025            {
 026                HttpConnection? connection = _connection;
 027                if (connection == null)
 028                {
 29                    // Response body fully consumed or the caller didn't ask for any data
 030                    return 0;
 31                }
 32
 033                int bytesRead = connection.ReadBuffered(buffer);
 034                if (bytesRead == 0 && buffer.Length != 0)
 035                {
 36                    // We cannot reuse this connection, so close it.
 037                    _connection = null;
 038                    connection.Dispose();
 039                }
 40
 041                return bytesRead;
 042            }
 43
 44            [AsyncMethodBuilder(typeof(PoolingAsyncValueTaskMethodBuilder<>))]
 45            [RuntimeAsyncMethodGeneration(false)]
 46            public override async ValueTask<int> ReadAsync(Memory<byte> buffer, CancellationToken cancellationToken)
 047            {
 048                CancellationHelper.ThrowIfCancellationRequested(cancellationToken);
 49
 050                HttpConnection? connection = _connection;
 051                if (connection == null)
 052                {
 53                    // Response body fully consumed
 054                    return 0;
 55                }
 56
 057                ValueTask<int> readTask = connection.ReadBufferedAsync(buffer);
 58                int bytesRead;
 059                if (readTask.IsCompletedSuccessfully)
 060                {
 061                    bytesRead = readTask.Result;
 062                }
 63                else
 064                {
 065                    CancellationTokenRegistration ctr = connection.RegisterCancellation(cancellationToken);
 66                    try
 067                    {
 068                        bytesRead = await readTask.ConfigureAwait(false);
 069                    }
 070                    catch (Exception exc) when (CancellationHelper.ShouldWrapInOperationCanceledException(exc, cancellat
 071                    {
 072                        throw CancellationHelper.CreateOperationCanceledException(exc, cancellationToken);
 73                    }
 74                    finally
 075                    {
 076                        ctr.Dispose();
 077                    }
 078                }
 79
 080                if (bytesRead == 0 && buffer.Length != 0)
 081                {
 82                    // A cancellation request may have caused the EOF.
 083                    CancellationHelper.ThrowIfCancellationRequested(cancellationToken);
 84
 85                    // We cannot reuse this connection, so close it.
 086                    _connection = null;
 087                    connection.Dispose();
 088                }
 89
 090                return bytesRead;
 091            }
 92
 93            public override Task CopyToAsync(Stream destination, int bufferSize, CancellationToken cancellationToken)
 094            {
 095                ValidateCopyToArguments(destination, bufferSize);
 96
 097                if (cancellationToken.IsCancellationRequested)
 098                {
 099                    return Task.FromCanceled(cancellationToken);
 100                }
 101
 0102                HttpConnection? connection = _connection;
 0103                if (connection == null)
 0104                {
 105                    // null if response body fully consumed
 0106                    return Task.CompletedTask;
 107                }
 108
 0109                Task copyTask = connection.CopyToUntilEofAsync(destination, async: true, bufferSize, cancellationToken);
 0110                if (copyTask.IsCompletedSuccessfully)
 0111                {
 0112                    Finish(connection);
 0113                    return Task.CompletedTask;
 114                }
 115
 0116                return CompleteCopyToAsync(copyTask, connection, cancellationToken);
 0117            }
 118
 119            private async Task CompleteCopyToAsync(Task copyTask, HttpConnection connection, CancellationToken cancellat
 0120            {
 0121                CancellationTokenRegistration ctr = connection.RegisterCancellation(cancellationToken);
 122                try
 0123                {
 0124                    await copyTask.ConfigureAwait(false);
 0125                }
 0126                catch (Exception exc) when (CancellationHelper.ShouldWrapInOperationCanceledException(exc, cancellationT
 0127                {
 0128                    throw CancellationHelper.CreateOperationCanceledException(exc, cancellationToken);
 129                }
 130                finally
 0131                {
 0132                    ctr.Dispose();
 0133                }
 134
 135                // If cancellation is requested and tears down the connection, it could cause the copy
 136                // to end early but think it ended successfully. So we prioritize cancellation in this
 137                // race condition, and if we find after the copy has completed that cancellation has
 138                // been requested, we assume the copy completed due to cancellation and throw.
 0139                CancellationHelper.ThrowIfCancellationRequested(cancellationToken);
 140
 0141                Finish(connection);
 0142            }
 143
 144            private void Finish(HttpConnection connection)
 0145            {
 146                // We cannot reuse this connection, so close it.
 0147                connection.Dispose();
 0148                _connection = null;
 0149            }
 150
 151            public override void Write(ReadOnlySpan<byte> buffer)
 0152            {
 0153                HttpConnection? connection = _connection;
 0154                if (connection == null)
 0155                {
 0156                    throw new IOException(SR.ObjectDisposed_StreamClosed);
 157                }
 158
 0159                if (buffer.Length != 0)
 0160                {
 0161                    connection.WriteWithoutBuffering(buffer);
 0162                }
 0163            }
 164
 165            public override ValueTask WriteAsync(ReadOnlyMemory<byte> buffer, CancellationToken cancellationToken)
 0166            {
 0167                if (cancellationToken.IsCancellationRequested)
 0168                {
 0169                    return ValueTask.FromCanceled(cancellationToken);
 170                }
 171
 0172                HttpConnection? connection = _connection;
 0173                if (connection == null)
 0174                {
 0175                    return ValueTask.FromException(ExceptionDispatchInfo.SetCurrentStackTrace(new IOException(SR.ObjectD
 176                }
 177
 0178                if (buffer.Length == 0)
 0179                {
 0180                    return default;
 181                }
 182
 0183                ValueTask writeTask = connection.WriteWithoutBufferingAsync(buffer, async: true);
 0184                return writeTask.IsCompleted ?
 0185                    writeTask :
 0186                    new ValueTask(WaitWithConnectionCancellationAsync(writeTask, connection, cancellationToken));
 0187            }
 188
 0189            public override void Flush() => _connection?.Flush();
 190
 191            public override Task FlushAsync(CancellationToken cancellationToken)
 0192            {
 0193                if (cancellationToken.IsCancellationRequested)
 0194                {
 0195                    return Task.FromCanceled(cancellationToken);
 196                }
 197
 0198                HttpConnection? connection = _connection;
 0199                if (connection == null)
 0200                {
 0201                    return Task.CompletedTask;
 202                }
 203
 0204                ValueTask flushTask = connection.FlushAsync(async: true);
 0205                return flushTask.IsCompleted ?
 0206                    flushTask.AsTask() :
 0207                    WaitWithConnectionCancellationAsync(flushTask, connection, cancellationToken);
 0208            }
 209
 210            private static async Task WaitWithConnectionCancellationAsync(ValueTask task, HttpConnection connection, Can
 0211            {
 0212                CancellationTokenRegistration ctr = connection.RegisterCancellation(cancellationToken);
 213                try
 0214                {
 0215                    await task.ConfigureAwait(false);
 0216                }
 0217                catch (Exception exc) when (CancellationHelper.ShouldWrapInOperationCanceledException(exc, cancellationT
 0218                {
 0219                    throw CancellationHelper.CreateOperationCanceledException(exc, cancellationToken);
 220                }
 221                finally
 0222                {
 0223                    ctr.Dispose();
 0224                }
 0225            }
 226        }
 227    }
 228}
 229

Methods/Properties

.ctor(System.Net.Http.HttpConnection,System.Net.Http.HttpResponseMessage)
Read(System.Span`1<System.Byte>)
ReadAsync(System.Memory`1<System.Byte>,System.Threading.CancellationToken)
ReadAsyncCore(System.Memory`1<System.Byte>,System.Threading.CancellationToken)
CopyToAsync(System.IO.Stream,System.Int32,System.Threading.CancellationToken)
CopyToAsyncCore(System.IO.Stream,System.Threading.CancellationToken)
PeekChunkFromConnectionBuffer()
ReadChunksFromConnectionBuffer(System.Span`1<System.Byte>,System.Threading.CancellationTokenRegistration)
ReadChunkFromConnectionBuffer(System.Int32,System.Threading.CancellationTokenRegistration)
ValidateChunkExtension(System.ReadOnlySpan`1<System.Byte>)
NeedsDrain()
DrainAsync(System.Int32)
Fill()
FillAsync()
.cctor()
.ctor(System.Net.Http.HttpConnection)
Write(System.ReadOnlySpan`1<System.Byte>)
WriteAsync(System.ReadOnlyMemory`1<System.Byte>,System.Threading.CancellationToken)
WriteChunkAsync(System.Net.Http.HttpConnection,System.ReadOnlyMemory`1<System.Byte>)
FinishAsync(System.Boolean)
.ctor(System.Net.Http.HttpConnection)
Read(System.Span`1<System.Byte>)
ReadAsync(System.Memory`1<System.Byte>,System.Threading.CancellationToken)
CopyToAsync(System.IO.Stream,System.Int32,System.Threading.CancellationToken)
CompleteCopyToAsync(System.Threading.Tasks.Task,System.Net.Http.HttpConnection,System.Threading.CancellationToken)
Finish(System.Net.Http.HttpConnection)
.ctor(System.Net.Http.HttpConnection,System.UInt64)
Read(System.Span`1<System.Byte>)
ReadAsync(System.Memory`1<System.Byte>,System.Threading.CancellationToken)
CopyToAsync(System.IO.Stream,System.Int32,System.Threading.CancellationToken)
CompleteCopyToAsync(System.Threading.Tasks.Task,System.Threading.CancellationToken)
Finish()
ReadFromConnectionBuffer(System.Int32)
NeedsDrain()
DrainAsync(System.Int32)
.ctor(System.Net.Http.HttpConnection,System.Int64)
Write(System.ReadOnlySpan`1<System.Byte>)
WriteAsync(System.ReadOnlyMemory`1<System.Byte>,System.Threading.CancellationToken)
FinishAsync(System.Boolean)
.cctor()
.ctor(System.Net.Http.HttpConnectionPool,System.IO.Stream,System.Net.TransportContext,System.Diagnostics.Activity,System.Net.IPEndPoint,System.Int64)
Finalize()
Dispose()
Dispose(System.Boolean)
ReadAheadTaskHasStarted()
PrepareForReuse(System.Boolean)
TryOwnScavengingTaskCompletion()
TryReturnScavengingTaskCompletionOwnership()
CheckUsabilityOnScavenge()
ReadAheadWithZeroByteReadAsync()
TransitionToCompletedAndTryOwnCompletion()
CheckKeepAliveTimeoutExceeded()
TransportContext()
Kind()
ReadBufferSize()
RemainingBuffer()
ConsumeFromRemainingBuffer(System.Int32)
WriteHeaders(System.Net.Http.HttpRequestMessage)
WriteHost(System.Uri)
WriteHeaderCollection(System.Net.Http.Headers.HttpHeaders,System.String)
WriteCRLF()
WriteBytes(System.ReadOnlySpan`1<System.Byte>)
WriteAsciiString(System.String)
WriteString(System.String,System.Text.Encoding)
ThrowForInvalidCharEncoding()
SendAsync(System.Net.Http.HttpRequestMessage,System.Boolean,System.Threading.CancellationToken)
MapSendException(System.Exception,System.Threading.CancellationToken,System.Exception&)
CreateRequestContentStream(System.Net.Http.HttpRequestMessage)
RegisterCancellation(System.Threading.CancellationToken)
SendRequestContentAsync(System.Net.Http.HttpRequestMessage,System.Net.Http.HttpConnection/HttpContentWriteStream,System.Boolean,System.Threading.CancellationToken)
SendRequestContentWithExpect100ContinueAsync(System.Net.Http.HttpRequestMessage,System.Threading.Tasks.Task`1<System.Boolean>,System.Net.Http.HttpConnection/HttpContentWriteStream,System.Threading.Timer,System.Boolean,System.Threading.CancellationToken)
ParseStatusLine(System.Net.Http.HttpResponseMessage)
ParseStatusLineCore(System.Span`1<System.Byte>,System.Net.Http.HttpResponseMessage)
ParseHeaders(System.Net.Http.HttpResponseMessage,System.Boolean)
ParseHeadersCore(System.Span`1<System.Byte>,System.Net.Http.HttpResponseMessage,System.Boolean)
ThrowForInvalidHeaderLine(System.ReadOnlySpan`1<System.Byte>,System.Int32)
AddResponseHeader(System.ReadOnlySpan`1<System.Byte>,System.ReadOnlySpan`1<System.Byte>,System.Net.Http.HttpResponseMessage,System.Boolean)
ThrowForEmptyHeaderName()
ThrowForInvalidHeaderName(System.ReadOnlySpan`1<System.Byte>)
ThrowExceededAllowedReadLineBytes()
ProcessKeepAliveHeader(System.String)
WriteToBuffer(System.ReadOnlySpan`1<System.Byte>)
Write(System.ReadOnlySpan`1<System.Byte>)
WriteAsync(System.ReadOnlyMemory`1<System.Byte>)
AwaitFlushAndWriteAsync(System.Threading.Tasks.ValueTask,System.ReadOnlyMemory`1<System.Byte>)
WriteWithoutBuffering(System.ReadOnlySpan`1<System.Byte>)
WriteWithoutBufferingAsync(System.ReadOnlyMemory`1<System.Byte>,System.Boolean)
FlushThenWriteWithoutBufferingAsync(System.ReadOnlyMemory`1<System.Byte>,System.Boolean)
WriteHexInt32Async(System.Int32,System.Boolean)
Flush()
FlushAsync(System.Boolean)
WriteToStream(System.ReadOnlySpan`1<System.Byte>)
WriteToStreamAsync(System.ReadOnlyMemory`1<System.Byte>,System.Boolean)
TryReadNextChunkedLine(System.ReadOnlySpan`1<System.Byte>&)
InitialFillAsync(System.Boolean)
FillAsync(System.Boolean)
FillForHeadersAsync(System.Boolean)
ReadUntilEndOfHeaderAsync(System.Boolean)
TryFindEndOfLine(System.ReadOnlySpan`1<System.Byte>,System.Int32&)
ReadFromBuffer(System.Span`1<System.Byte>)
Read(System.Span`1<System.Byte>)
ReadAsync(System.Memory`1<System.Byte>)
ReadAndLogBytesReadAsync(System.Memory`1<System.Byte>)
ReadBuffered(System.Span`1<System.Byte>)
ReadBufferedAsync(System.Memory`1<System.Byte>)
ReadBufferedAsyncCore()
CopyFromBufferAsync(System.IO.Stream,System.Boolean,System.Int32,System.Threading.CancellationToken)
CopyToUntilEofAsync(System.IO.Stream,System.Boolean,System.Int32,System.Threading.CancellationToken)
CopyToUntilEofWithExistingBufferedDataAsync(System.IO.Stream,System.Boolean,System.Int32,System.Threading.CancellationToken)
CopyToContentLengthAsync(System.IO.Stream,System.Boolean,System.UInt64,System.Int32,System.Threading.CancellationToken)
Acquire()
Release()
DetachFromPool()
CompleteResponse()
DrainResponseAsync(System.Net.Http.HttpResponseMessage,System.Threading.CancellationToken)
ReturnConnectionToPool()
ToString()
Trace(System.String,System.String)
.ctor(System.Net.Http.HttpConnection)
CanRead()
CanWrite()
Write(System.ReadOnlySpan`1<System.Byte>)
WriteAsync(System.ReadOnlyMemory`1<System.Byte>,System.Threading.CancellationToken)
NeedsDrain()
IsDisposed()
CanReadFromConnection()
DrainAsync(System.Int32)
Dispose(System.Boolean)
DrainOnDisposeAsync()
BytesWritten()
.ctor(System.Net.Http.HttpConnection)
CanRead()
CanWrite()
Flush()
FlushAsync(System.Threading.CancellationToken)
Read(System.Span`1<System.Byte>)
ReadAsync(System.Memory`1<System.Byte>,System.Threading.CancellationToken)
CopyToAsync(System.IO.Stream,System.Int32,System.Threading.CancellationToken)
.ctor(System.Net.Http.HttpConnection)
CanRead()
CanWrite()
Read(System.Span`1<System.Byte>)
ReadAsync()
CopyToAsync(System.IO.Stream,System.Int32,System.Threading.CancellationToken)
CompleteCopyToAsync(System.Threading.Tasks.Task,System.Net.Http.HttpConnection,System.Threading.CancellationToken)
Finish(System.Net.Http.HttpConnection)
Write(System.ReadOnlySpan`1<System.Byte>)
WriteAsync(System.ReadOnlyMemory`1<System.Byte>,System.Threading.CancellationToken)
Flush()
FlushAsync(System.Threading.CancellationToken)
WaitWithConnectionCancellationAsync(System.Threading.Tasks.ValueTask,System.Net.Http.HttpConnection,System.Threading.CancellationToken)