From fc38a134b97c384a5beb83b4158abfc6ed446256 Mon Sep 17 00:00:00 2001 From: justcoding121 Date: Sun, 4 Oct 2026 16:07:03 +0530 Subject: [PATCH 1/2] Clear the open Sonar code smells without changing per-frame HTTP/2 send behavior. --- .../Services/IInspectorPathPicker.cs | 2 +- .../Services/SessionArchive.cs | 72 +++++----- .../Services/SessionBodyDiskCache.cs | 9 +- .../Services/SessionStore.cs | 59 +++++---- .../Services/SessionStreamBuffer.cs | 74 ++++++----- .../ViewModels/AutoResponderViewModel.cs | 6 +- .../MainWindowViewModel.Sessions.cs | 35 ++--- .../ViewModels/MapRemoteViewModel.cs | 6 +- .../Views/ExcludedHostsWindow.axaml.cs | 2 - .../Handlers/H1TerminateFastForward.cs | 6 +- .../Handlers/Http11ToHttp2BridgeHandler.cs | 1 - .../Handlers/Http2ToHttp11BridgeHandler.cs | 7 +- .../Http/HttpHeaderHygiene.cs | 9 +- .../Http2/Http2DeferredOutboundData.cs | 2 +- .../Http2/Http2Helper.Copy.Headers.cs | 2 +- .../Http2/Http2Helper.Copy.Relay.cs | 5 +- .../Http2/Http2Helper.Copy.cs | 27 ++-- .../Http2/Http2Helper.Send.cs | 14 +- src/Titanium.Web.Proxy/Http2/Http2Helper.cs | 2 + .../Http2/Http2OriginConnection.cs | 93 ++++++++----- .../Http3/H3H1QpackResponseReader.cs | 45 +++++-- src/Titanium.Web.Proxy/Http3/Http3Frame.cs | 18 +-- .../Http3/Http3OriginBridge.Http2.cs | 123 ++++++++++-------- .../Http3/Http3OriginBridge.Tcp.cs | 7 +- .../Http3/Http3RequestStream.cs | 13 +- .../DeferredInspectorUiCoverageTests.cs | 3 +- .../LiveSessionUpdateUxTests.cs | 4 +- .../Http2ContinuationAndAbuseBudgetTests.cs | 6 +- .../HandlerAndProtocolHelperCoverageTests.cs | 4 +- .../Http2BufferedSendOrderingTests.cs | 3 +- .../SonarGateCoverageBumpTests.cs | 6 + .../SonarNewCodeCoverageTests.cs | 10 +- 32 files changed, 376 insertions(+), 299 deletions(-) diff --git a/src/Titanium.Inspector/Services/IInspectorPathPicker.cs b/src/Titanium.Inspector/Services/IInspectorPathPicker.cs index c71b99534..2389752b6 100644 --- a/src/Titanium.Inspector/Services/IInspectorPathPicker.cs +++ b/src/Titanium.Inspector/Services/IInspectorPathPicker.cs @@ -111,7 +111,7 @@ public Task> PickOpenPathsAsync(string title, string filte internal readonly record struct StoragePickAttempt(bool DialogShown, string? Path) { public IReadOnlyList Paths { get; init; } = - string.IsNullOrEmpty(Path) ? Array.Empty() : [Path!]; + string.IsNullOrEmpty(Path) ? Array.Empty() : [Path]; } internal static class InspectorPathPickerHelpers diff --git a/src/Titanium.Inspector/Services/SessionArchive.cs b/src/Titanium.Inspector/Services/SessionArchive.cs index f90bff937..cf005ea70 100644 --- a/src/Titanium.Inspector/Services/SessionArchive.cs +++ b/src/Titanium.Inspector/Services/SessionArchive.cs @@ -235,7 +235,7 @@ public static JsonObject ToHarEntry(SessionSnapshot s) { try { - if (!TryReadHarRequest(entry, out var method, out var url, out var host, out var reqHeaders, out var reqBody, out var contentType, out var reqBytes)) + if (!TryReadHarRequest(entry, out var request)) { return null; } @@ -246,18 +246,18 @@ public static JsonObject ToHarEntry(SessionSnapshot s) var snap = new SessionSnapshot { Id = id, - Method = method, - Url = url, - Host = host, + Method = request.Method, + Url = request.Url, + Host = request.Host, StartedUtc = started, StatusCode = status, - RequestHeadersText = reqHeaders, + RequestHeadersText = request.Headers, ResponseHeadersText = respHeaders, - RequestBodyText = reqBody, + RequestBodyText = request.Body, ResponseBodyText = respBody, - RequestBodyBytes = reqBytes, + RequestBodyBytes = request.BodyBytes, ResponseBodyBytes = respBytes, - ContentType = contentType ?? respMime, + ContentType = request.ContentType ?? respMime, DurationMs = durationMs, TtfbMs = ttfbMs, BodySize = respBytes?.LongLength ?? respBody?.Length, @@ -386,6 +386,14 @@ private static void ApplyInspectorExtension(JsonElement entry, SessionSnapshot s return; } + ApplyInspectorIdentity(ext, snap); + ApplyInspectorCapture(ext, snap); + ApplyInspectorFlags(ext, snap); + ApplyInspectorBodies(ext, snap); + } + + private static void ApplyInspectorIdentity(JsonElement ext, SessionSnapshot snap) + { if (ext.TryGetProperty("id", out var idEl) && idEl.TryGetInt64(out var savedId) && savedId > 0) { snap.Id = savedId; @@ -415,7 +423,10 @@ private static void ApplyInspectorExtension(JsonElement entry, SessionSnapshot s { snap.ProcessName = pn.GetString(); } + } + private static void ApplyInspectorCapture(JsonElement ext, SessionSnapshot snap) + { if (ext.TryGetProperty("sentBytes", out var sent) && sent.TryGetInt64(out var sentBytes)) { snap.SentBytes = sentBytes; @@ -447,7 +458,10 @@ private static void ApplyInspectorExtension(JsonElement entry, SessionSnapshot s { snap.ResponseBodyOriginalSize = respOrig; } + } + private static void ApplyInspectorFlags(JsonElement ext, SessionSnapshot snap) + { snap.IsWebSocket = ReadBool(ext, "isWebSocket") ?? snap.IsWebSocket; snap.IsGrpc = ReadBool(ext, "isGrpc") ?? snap.IsGrpc; snap.IsTranscoded = ReadBool(ext, "isTranscoded") ?? snap.IsTranscoded; @@ -468,7 +482,10 @@ private static void ApplyInspectorExtension(JsonElement entry, SessionSnapshot s snap.UpstreamPath = ReadString(ext, "upstreamPath") ?? snap.UpstreamPath; snap.UpstreamContentType = ReadString(ext, "upstreamContentType") ?? snap.UpstreamContentType; snap.ProtobufDecodedText = ReadString(ext, "protobufDecodedText") ?? snap.ProtobufDecodedText; + } + private static void ApplyInspectorBodies(JsonElement ext, SessionSnapshot snap) + { // Prefer _inspector body payloads when present (lossless for Inspect). if (TryDecodeBase64(ext, "requestBodyBase64", out var reqBytes)) { @@ -539,7 +556,7 @@ private static bool TryDecodeBase64(JsonElement ext, string name, out byte[] byt } } - private static IReadOnlyList? DeserializeList(JsonElement ext, string name) + private static List? DeserializeList(JsonElement ext, string name) { if (!ext.TryGetProperty(name, out var el) || el.ValueKind != JsonValueKind.Array) { @@ -556,37 +573,33 @@ private static bool TryDecodeBase64(JsonElement ext, string name, out byte[] byt } } - private static bool TryReadHarRequest( - JsonElement entry, - out string method, - out string url, - out string? host, - out string reqHeaders, - out string? reqBody, - out string? contentType, - out byte[]? reqBytes) - { - method = "GET"; - url = ""; - host = null; - reqHeaders = ""; - reqBody = null; - contentType = null; - reqBytes = null; + private readonly record struct HarRequestFields( + string Method, + string Url, + string? Host, + string Headers, + string? Body, + string? ContentType, + byte[]? BodyBytes); + private static bool TryReadHarRequest(JsonElement entry, out HarRequestFields request) + { + request = new HarRequestFields("GET", "", null, "", null, null, null); if (!entry.TryGetProperty("request", out var req)) { return false; } - method = req.TryGetProperty("method", out var m) ? m.GetString() ?? "GET" : "GET"; - url = req.TryGetProperty("url", out var u) ? u.GetString() ?? "" : ""; + var method = req.TryGetProperty("method", out var m) ? m.GetString() ?? "GET" : "GET"; + var url = req.TryGetProperty("url", out var u) ? u.GetString() ?? "" : ""; + string? host = null; if (Uri.TryCreate(url, UriKind.Absolute, out var uri)) { host = uri.Host; } - reqHeaders = FormatHarHeaders(req); + string? reqBody = null; + string? contentType = null; if (req.TryGetProperty("postData", out var post) && post.ValueKind == JsonValueKind.Object) { if (post.TryGetProperty("text", out var pt)) @@ -600,6 +613,7 @@ private static bool TryReadHarRequest( } } + request = new HarRequestFields(method, url, host, FormatHarHeaders(req), reqBody, contentType, null); return true; } diff --git a/src/Titanium.Inspector/Services/SessionBodyDiskCache.cs b/src/Titanium.Inspector/Services/SessionBodyDiskCache.cs index 0b0ddeaab..fc7179bc4 100644 --- a/src/Titanium.Inspector/Services/SessionBodyDiskCache.cs +++ b/src/Titanium.Inspector/Services/SessionBodyDiskCache.cs @@ -695,7 +695,7 @@ private void RemoveIndexEntryLocked(string path) } } - private void QueueCleanup(IReadOnlyList paths) + private void QueueCleanup(List paths) { if (paths.Count == 0) { @@ -704,12 +704,9 @@ private void QueueCleanup(IReadOnlyList paths) lock (_cleanupGate) { - foreach (var path in paths) + foreach (var path in paths.Where(static path => !string.IsNullOrEmpty(path))) { - if (!string.IsNullOrEmpty(path)) - { - _cleanupPaths.Enqueue(path); - } + _cleanupPaths.Enqueue(path); } _cleanupTask = _cleanupTask.ContinueWith( diff --git a/src/Titanium.Inspector/Services/SessionStore.cs b/src/Titanium.Inspector/Services/SessionStore.cs index a398905bf..380c88605 100644 --- a/src/Titanium.Inspector/Services/SessionStore.cs +++ b/src/Titanium.Inspector/Services/SessionStore.cs @@ -147,14 +147,12 @@ public long? PinnedSessionId { var prevId = _pinnedSessionId; _pinnedSessionId = value; - if (prevId is long oldId && oldId != value && _byId.TryGetValue(oldId, out previous)) + // Drop RAM bodies for the previous selection when the file already exists. + if (prevId is long oldId && oldId != value && _byId.TryGetValue(oldId, out previous) + && previous.BodiesOnDisk) { - // Drop RAM bodies for the previous selection when the file already exists. - if (previous.BodiesOnDisk) - { - ClearBodyFields(previous); - RecalcInMemoryBodyBytesLocked(); - } + ClearBodyFields(previous); + RecalcInMemoryBodyBytesLocked(); } if (value is long newId && _byId.TryGetValue(newId, out next)) @@ -338,7 +336,8 @@ public async Task FlushDiskCleanupAsync(TimeSpan? timeout = null) public async Task EnsureBodiesLoadedAsync(SessionSnapshot snapshot, CancellationToken ct = default) { ObjectDisposedException.ThrowIf(_disposed, this); - if (_disk is null) + var disk = _disk; + if (disk is null) { return; } @@ -355,15 +354,34 @@ public async Task EnsureBodiesLoadedAsync(SessionSnapshot snapshot, Cancellation } // Fast path: file already gone and nothing pending — do not wait ~1s. - if (!_disk.FileExists(snapshot.Id) && Volatile.Read(ref _pendingSpills) == 0) + if (!disk.FileExists(snapshot.Id) && Volatile.Read(ref _pendingSpills) == 0) { snapshot.BodiesMissingFromDisk = true; return; } - // Under heavy capture, the spill writer may still be draining thousands of HARs. - // Keep waiting while work is queued; only mark missing once the channel is idle - // and the file is still absent (bounded by ct / ~2 minutes). + if (await WaitForSpilledBodyAsync(snapshot, disk, ct).ConfigureAwait(false)) + { + return; + } + + if (!HasInMemoryBodies(snapshot) && + Volatile.Read(ref _pendingSpills) == 0 && + !disk.FileExists(snapshot.Id)) + { + snapshot.BodiesMissingFromDisk = true; + } + } + + /// + /// Under heavy capture, the spill writer may still be draining thousands of HARs. + /// Keep waiting while work is queued; only mark missing once the channel is idle + /// and the file is still absent (bounded by ct / ~2 minutes). + /// Returns true when the caller should stop (body loaded or marked missing). + /// + private async Task WaitForSpilledBodyAsync( + SessionSnapshot snapshot, SessionBodyDiskCache disk, CancellationToken ct) + { var deadline = DateTime.UtcNow + TimeSpan.FromMinutes(2); while (DateTime.UtcNow < deadline) { @@ -373,33 +391,28 @@ public async Task EnsureBodiesLoadedAsync(SessionSnapshot snapshot, Cancellation if (HasInMemoryBodies(snapshot)) { snapshot.BodiesMissingFromDisk = false; - return; + return true; } - if (_disk.TryLoad(snapshot)) + if (disk.TryLoad(snapshot)) { // Keep BodiesOnDisk=true so deselect can unload without rewriting the file. snapshot.BodiesMissingFromDisk = false; RecalcInMemoryBodyBytesLocked(); - return; + return true; } } - if (Volatile.Read(ref _pendingSpills) == 0 && !_disk.FileExists(snapshot.Id)) + if (Volatile.Read(ref _pendingSpills) == 0 && !disk.FileExists(snapshot.Id)) { snapshot.BodiesMissingFromDisk = true; - return; + return true; } await Task.Delay(25, ct).ConfigureAwait(false); } - if (!HasInMemoryBodies(snapshot) && - Volatile.Read(ref _pendingSpills) == 0 && - !_disk.FileExists(snapshot.Id)) - { - snapshot.BodiesMissingFromDisk = true; - } + return false; } public async Task EnsureBodiesLoadedAsync(IEnumerable snapshots, CancellationToken ct = default) diff --git a/src/Titanium.Inspector/Services/SessionStreamBuffer.cs b/src/Titanium.Inspector/Services/SessionStreamBuffer.cs index 865943dc8..9b240a297 100644 --- a/src/Titanium.Inspector/Services/SessionStreamBuffer.cs +++ b/src/Titanium.Inspector/Services/SessionStreamBuffer.cs @@ -55,48 +55,58 @@ private async Task ReadLoopAsync() await foreach (var snapshot in _channel.Reader.ReadAllAsync()) { batch.Add(snapshot); - var deadline = Environment.TickCount64 + _batchWindowMs; - while (batch.Count < _batchMax && _channel.Reader.TryRead(out var more)) + await CoalesceBatchAsync(batch).ConfigureAwait(false); + SessionsBatchAdded?.Invoke(batch.ToList()); + RaiseSessionAdded(batch); + batch.Clear(); + } + } + + private async Task CoalesceBatchAsync(List batch) + { + var deadline = Environment.TickCount64 + _batchWindowMs; + while (batch.Count < _batchMax && _channel.Reader.TryRead(out var more)) + { + batch.Add(more); + } + + // Wait briefly for more arrivals when the channel is momentarily empty. + while (batch.Count < _batchMax && Environment.TickCount64 < deadline) + { + var remaining = (int)(deadline - Environment.TickCount64); + if (remaining <= 0) { - batch.Add(more); + break; } - // Wait briefly for more arrivals when the channel is momentarily empty. - while (batch.Count < _batchMax && Environment.TickCount64 < deadline) + using var delayCts = new CancellationTokenSource(remaining); + try { - var remaining = (int)(deadline - Environment.TickCount64); - if (remaining <= 0) - { - break; - } - - using var delayCts = new CancellationTokenSource(remaining); - try - { - if (await _channel.Reader.WaitToReadAsync(delayCts.Token).ConfigureAwait(false)) - { - while (batch.Count < _batchMax && _channel.Reader.TryRead(out var more)) - { - batch.Add(more); - } - } - } - catch (OperationCanceledException) + if (await _channel.Reader.WaitToReadAsync(delayCts.Token).ConfigureAwait(false)) { - break; + DrainAvailable(batch); } } - - SessionsBatchAdded?.Invoke(batch.ToList()); - if (SessionAdded is { } singleHandler) + catch (OperationCanceledException) { - foreach (var s in batch) - { - singleHandler(s); - } + break; } + } + } - batch.Clear(); + private void DrainAvailable(List batch) + { + while (batch.Count < _batchMax && _channel.Reader.TryRead(out var more)) + { + batch.Add(more); + } + } + + private void RaiseSessionAdded(List batch) + { + foreach (var snapshot in batch) + { + SessionAdded?.Invoke(snapshot); } } } diff --git a/src/Titanium.Inspector/ViewModels/AutoResponderViewModel.cs b/src/Titanium.Inspector/ViewModels/AutoResponderViewModel.cs index 2b3ce523d..5d7690bee 100644 --- a/src/Titanium.Inspector/ViewModels/AutoResponderViewModel.cs +++ b/src/Titanium.Inspector/ViewModels/AutoResponderViewModel.cs @@ -158,6 +158,9 @@ public bool TryMatch(string url, string? requestBody, out AutoResponderRule? mat return false; } + public bool TryMatch(string url, out AutoResponderRule? matched) + => TryMatch(url, requestBody: null, out matched); + private AutoResponderRule[] SnapshotRules() { lock (_rulesGate) @@ -168,9 +171,6 @@ private AutoResponderRule[] SnapshotRules() } } - public bool TryMatch(string url, out AutoResponderRule? matched) - => TryMatch(url, requestBody: null, out matched); - public bool TryRespond(SessionSnapshot session, out AutoResponderRule? matched) => TryMatch(session.Url, session.RequestBodyText, out matched); diff --git a/src/Titanium.Inspector/ViewModels/MainWindowViewModel.Sessions.cs b/src/Titanium.Inspector/ViewModels/MainWindowViewModel.Sessions.cs index 84117b783..1495290dd 100644 --- a/src/Titanium.Inspector/ViewModels/MainWindowViewModel.Sessions.cs +++ b/src/Titanium.Inspector/ViewModels/MainWindowViewModel.Sessions.cs @@ -603,15 +603,13 @@ private void OnSessionsBatchAdded(IReadOnlyList batch) return; } - foreach (var snapshot in batch) + // SessionUpdated may have raced ahead of this batched capture add and already + // inserted the row; never append the same snapshot twice. + foreach (var snapshot in batch.Where(snapshot => + SessionSearch.Matches(snapshot, SearchQuery, bodyMatcher: null) && + Sessions.IndexOf(snapshot) < 0)) { - // SessionUpdated may have raced ahead of this batched capture add and already - // inserted the row; never append the same snapshot twice. - if (SessionSearch.Matches(snapshot, SearchQuery, bodyMatcher: null) && - Sessions.IndexOf(snapshot) < 0) - { - Sessions.Add(snapshot); - } + Sessions.Add(snapshot); } RefreshSessionCountText(); @@ -820,15 +818,7 @@ private void RemoveVisibleByIds(HashSet ids) return; } - var removeCount = 0; - foreach (var session in Sessions) - { - if (ids.Contains(session.Id)) - { - removeCount++; - } - } - + var removeCount = Sessions.Count(session => ids.Contains(session.Id)); if (removeCount == 0) { return; @@ -842,16 +832,7 @@ private void RemoveVisibleByIds(HashSet ids) if (removeCount >= BulkGridEditThreshold && removeCount * 2 >= Sessions.Count) { - var keep = new List(Sessions.Count - removeCount); - foreach (var session in Sessions) - { - if (!ids.Contains(session.Id)) - { - keep.Add(session); - } - } - - ReplaceVisibleSessions(keep); + ReplaceVisibleSessions(Sessions.Where(session => !ids.Contains(session.Id)).ToList()); return; } diff --git a/src/Titanium.Inspector/ViewModels/MapRemoteViewModel.cs b/src/Titanium.Inspector/ViewModels/MapRemoteViewModel.cs index 2a621d96d..9ca09ea34 100644 --- a/src/Titanium.Inspector/ViewModels/MapRemoteViewModel.cs +++ b/src/Titanium.Inspector/ViewModels/MapRemoteViewModel.cs @@ -157,6 +157,9 @@ public bool TryRewrite(string url, string? requestBody, out string? rewritten, o return false; } + public bool TryRewrite(string url, out string? rewritten, out MapRemoteRule? matched) + => TryRewrite(url, requestBody: null, out rewritten, out matched); + private MapRemoteRule[] SnapshotRules() { lock (_rulesGate) @@ -167,9 +170,6 @@ private MapRemoteRule[] SnapshotRules() } } - public bool TryRewrite(string url, out string? rewritten, out MapRemoteRule? matched) - => TryRewrite(url, requestBody: null, out rewritten, out matched); - /// /// Rewrites using wildcard capture from /// into . A single trailing * in both match and target diff --git a/src/Titanium.Inspector/Views/ExcludedHostsWindow.axaml.cs b/src/Titanium.Inspector/Views/ExcludedHostsWindow.axaml.cs index 23b1ccbd3..f5eb82acb 100644 --- a/src/Titanium.Inspector/Views/ExcludedHostsWindow.axaml.cs +++ b/src/Titanium.Inspector/Views/ExcludedHostsWindow.axaml.cs @@ -10,7 +10,6 @@ namespace Titanium.Inspector.Views; public partial class ExcludedHostsWindow : Window { private readonly SettingsService _settings; - private readonly bool _readOnly; private readonly Action? _onSaved; private readonly InterceptionService? _interception; private bool _saved; @@ -27,7 +26,6 @@ public ExcludedHostsWindow( InterceptionService? interception = null) { _settings = settings; - _readOnly = readOnly; _onSaved = onSaved; _interception = interception; InitializeComponent(); diff --git a/src/Titanium.Web.Proxy/Handlers/H1TerminateFastForward.cs b/src/Titanium.Web.Proxy/Handlers/H1TerminateFastForward.cs index d2de15548..0ba7da5e7 100644 --- a/src/Titanium.Web.Proxy/Handlers/H1TerminateFastForward.cs +++ b/src/Titanium.Web.Proxy/Handlers/H1TerminateFastForward.cs @@ -26,7 +26,7 @@ namespace Titanium.Web.Proxy; /// Signals that H1 terminate-lite cannot finish this exchange (typically a 1xx interim response) /// and the caller should fall through to the full session path without treating it as a failure. /// -internal sealed class H1TerminateLiteFallbackException : Exception +internal sealed class H1TerminateLiteFallbackException : Exception // NOSONAR S3871 -- Internal control-flow signal for the terminate-lite fallback, not a public error contract. { public H1TerminateLiteFallbackException(string message) : base(message) { @@ -211,7 +211,7 @@ internal bool CanUseH1TerminateLite(ProxyEndPoint endPoint, Request request, boo if (response.StatusCode is >= 100 and <= 199) { // Full session path has the interim 1xx loop; keep the origin connection for it. - closeConnection = false; + // closeConnection is already false, so the fallback catch releases the socket for reuse. throw new H1TerminateLiteFallbackException( "H1 terminate lite does not handle interim 1xx responses."); } @@ -424,7 +424,7 @@ await connection.Stream.CopyBodyAsync(response, false, clientStream, Transformat var response = http.Response; if (response.StatusCode is >= 100 and <= 199) { - closeConnection = false; + // closeConnection is already false, so the fallback catch releases the socket for reuse. throw new H1TerminateLiteFallbackException( "H1 terminate MITM lite does not handle interim 1xx responses."); } diff --git a/src/Titanium.Web.Proxy/Handlers/Http11ToHttp2BridgeHandler.cs b/src/Titanium.Web.Proxy/Handlers/Http11ToHttp2BridgeHandler.cs index db8d22a63..cdd3b9612 100644 --- a/src/Titanium.Web.Proxy/Handlers/Http11ToHttp2BridgeHandler.cs +++ b/src/Titanium.Web.Proxy/Handlers/Http11ToHttp2BridgeHandler.cs @@ -748,7 +748,6 @@ await clientStream.WriteBodyAsync(fastBody, response.IsChunked, exchange.Response.StreamBodyWriter = null; // Drain and release origin resources; the replacement body is unrelated. await originStreamBody(Stream.Null, cancellationToken); - originStreamBody = null; } if (response.Locked) diff --git a/src/Titanium.Web.Proxy/Handlers/Http2ToHttp11BridgeHandler.cs b/src/Titanium.Web.Proxy/Handlers/Http2ToHttp11BridgeHandler.cs index 32e88f448..7d4e30922 100644 --- a/src/Titanium.Web.Proxy/Handlers/Http2ToHttp11BridgeHandler.cs +++ b/src/Titanium.Web.Proxy/Handlers/Http2ToHttp11BridgeHandler.cs @@ -1135,9 +1135,10 @@ private static async ValueTask RelayStreamToClientAsync( try { var readVt = source.ReadAsync(buf.AsMemory(), CancellationToken.None); - read = readVt.IsCompletedSuccessfully - ? readVt.Result - : await readVt.ConfigureAwait(false); + if (readVt.IsCompletedSuccessfully) + read = readVt.Result; + else + read = await readVt.ConfigureAwait(false); } catch (Exception readEx) { diff --git a/src/Titanium.Web.Proxy/Http/HttpHeaderHygiene.cs b/src/Titanium.Web.Proxy/Http/HttpHeaderHygiene.cs index dabeea3ec..1510a3587 100644 --- a/src/Titanium.Web.Proxy/Http/HttpHeaderHygiene.cs +++ b/src/Titanium.Web.Proxy/Http/HttpHeaderHygiene.cs @@ -22,13 +22,10 @@ public static bool ContainsForbiddenDelimiter(string? value) { if (string.IsNullOrEmpty(value)) return false; - foreach (var c in value) - { - if (c is '\0' or '\r' or '\n') - return true; - } - return false; + // IndexOfAny is vectorized and allocation-free. A LINQ Any() walk would box a delegate + // on a path that runs for every header written to an HTTP/1.1 wire. + return value.AsSpan().IndexOfAny('\0', '\r', '\n') >= 0; } public static void ThrowIfForbidden(ByteString name, ByteString value) diff --git a/src/Titanium.Web.Proxy/Http2/Http2DeferredOutboundData.cs b/src/Titanium.Web.Proxy/Http2/Http2DeferredOutboundData.cs index b239749e1..16ccbdbf0 100644 --- a/src/Titanium.Web.Proxy/Http2/Http2DeferredOutboundData.cs +++ b/src/Titanium.Web.Proxy/Http2/Http2DeferredOutboundData.cs @@ -137,7 +137,7 @@ public void CancelStream(int streamId) /// runs after the lock is released with the same stream id and /// may remove the stream (which cancels this queue). /// - public bool TryDrain(Http2FlowController flow, Http2FrameWriter writer, + public bool TryDrain(Http2FlowController flow, Http2FrameWriter writer, // NOSONAR S3776 -- Drain holds the queue lock and must keep DATA frame order; splitting it risks a short-window reorder. Action? onEndStreamSent = null, Action? onEndStreamQueued = null) { var wrote = false; diff --git a/src/Titanium.Web.Proxy/Http2/Http2Helper.Copy.Headers.cs b/src/Titanium.Web.Proxy/Http2/Http2Helper.Copy.Headers.cs index ea23fb4a5..deb6146fd 100644 --- a/src/Titanium.Web.Proxy/Http2/Http2Helper.Copy.Headers.cs +++ b/src/Titanium.Web.Proxy/Http2/Http2Helper.Copy.Headers.cs @@ -55,7 +55,7 @@ internal partial class Http2Helper int maxDecodedHeaderListBytes, ILogger logger, Action removeAndFinalizeStream, - Func, ValueTask> lockedOutputWrite, + Func, ValueTask> lockedOutputWrite, // NOSONAR S1172 -- Kept so this dispatch stays argument-aligned with CopyHttp2FrameAsync. bool forceStaticHpackTable, Http2Settings localSettings, // NOSONAR S1172 -- retained for CopyHttp2FrameAsync call-site IL match. HeaderCollection headerDecodeScratch, diff --git a/src/Titanium.Web.Proxy/Http2/Http2Helper.Copy.Relay.cs b/src/Titanium.Web.Proxy/Http2/Http2Helper.Copy.Relay.cs index b9d20a312..2f101649e 100644 --- a/src/Titanium.Web.Proxy/Http2/Http2Helper.Copy.Relay.cs +++ b/src/Titanium.Web.Proxy/Http2/Http2Helper.Copy.Relay.cs @@ -141,7 +141,7 @@ private static Task RelayCompressedWithOriginPoolAsync( // NOSONAR S107 -- Origi hbStreamId, blockToRelay, endStreamFlag, appendSuffix); } - private static async Task AssignOriginAndEnqueueAsync( + private static async Task AssignOriginAndEnqueueAsync( // NOSONAR S107 -- Cold origin assign keeps the relay signature explicit (no per-open context object). Http2ConnectionState connectionState, bool isClient, Http2Settings remoteSettings, @@ -151,11 +151,12 @@ private static async Task AssignOriginAndEnqueueAsync( { var assignTask = connectionState.OriginRelayPool! .AssignStreamAsync(hbStreamId, cancellationToken); + // AsTask consumes the ValueTask. Await that Task; a second await of assignTask is undefined. var tracked = assignTask.AsTask(); connectionState.PendingOriginAssignments[hbStreamId] = tracked; try { - var assignment = await assignTask.ConfigureAwait(false); + var assignment = await tracked.ConfigureAwait(false); EnqueueRelayedHeaderBlock(connectionState, isClient, remoteSettings, assignment.OriginStreamId, blockToRelay, endStreamFlag, appendSuffix, assignment.Leg.Writer, assignment.Leg.WriteLock, assignment.Leg.Stream); diff --git a/src/Titanium.Web.Proxy/Http2/Http2Helper.Copy.cs b/src/Titanium.Web.Proxy/Http2/Http2Helper.Copy.cs index efa743269..b0627f52e 100644 --- a/src/Titanium.Web.Proxy/Http2/Http2Helper.Copy.cs +++ b/src/Titanium.Web.Proxy/Http2/Http2Helper.Copy.cs @@ -684,16 +684,13 @@ async Task TrySendGracefulGoAwayAsync() { if (length > 0) ArrayPool.Shared.Return(payloadRented); - if (dataEndStream) + // already RequestClosed; may now finalize if response half is done + if (dataEndStream && compressedDataState.IsClosed) { - // already RequestClosed; may now finalize if response half is done - if (compressedDataState.IsClosed) - { - connectionState.OriginRelayPool?.ReleaseStream(dataStreamId); - connectionState.RemoveStream(dataStreamId); - ScheduleFinalize(compressedDataState, onAfterResponse, logger, - connectionState); - } + connectionState.OriginRelayPool?.ReleaseStream(dataStreamId); + connectionState.RemoveStream(dataStreamId); + ScheduleFinalize(compressedDataState, onAfterResponse, logger, + connectionState); } continue; @@ -798,7 +795,7 @@ void EnqueueWireFrame(ReadOnlySpan payload, bool endStream) if (ownedPayload) payloadRented = Array.Empty(); // already returned ReportException(logger, new ProxyHttpException( - "HTTP/2 deferred DATA queue exceeded its per-stream cap.", null, null)); + DeferredDataQueueCapExceededMessage, null, null)); await lockedOwnLegWrite(() => SendRstStreamAsync(new Http2FrameHeader(), new byte[9], peerStreamId != 0 ? peerStreamId : dataStreamId, Http2ErrorCode.EnhanceYourCalm, input)); @@ -842,7 +839,7 @@ await lockedOwnLegWrite(() => SendRstStreamAsync(new Http2FrameHeader(), { ArrayPool.Shared.Return(payloadRented); ReportException(logger, new ProxyHttpException( - "HTTP/2 deferred DATA queue exceeded its per-stream cap.", null, null)); + DeferredDataQueueCapExceededMessage, null, null)); await lockedOwnLegWrite(() => SendRstStreamAsync(new Http2FrameHeader(), new byte[9], peerStreamId != 0 ? peerStreamId : dataStreamId, Http2ErrorCode.EnhanceYourCalm, input)); @@ -863,7 +860,7 @@ await lockedOwnLegWrite(() => SendRstStreamAsync(new Http2FrameHeader(), { ArrayPool.Shared.Return(payloadRented); ReportException(logger, new ProxyHttpException( - "HTTP/2 deferred DATA queue exceeded its per-stream cap.", null, null)); + DeferredDataQueueCapExceededMessage, null, null)); await lockedOwnLegWrite(() => SendRstStreamAsync(new Http2FrameHeader(), new byte[9], peerStreamId != 0 ? peerStreamId : dataStreamId, Http2ErrorCode.EnhanceYourCalm, input)); @@ -1685,7 +1682,7 @@ await lockedOwnLegWrite(() => SendRstStreamAsync( else if (queued == QueueSendDataResult.CapExceeded) { ReportException(logger, new ProxyHttpException( - "HTTP/2 deferred DATA queue exceeded its per-stream cap.", null, args)); + DeferredDataQueueCapExceededMessage, null, args)); await lockedOwnLegWrite(() => SendRstStreamAsync(new Http2FrameHeader(), new byte[9], streamId, Http2ErrorCode.EnhanceYourCalm, input)); RemoveAndFinalizeStream(streamId); @@ -2310,7 +2307,7 @@ await lockedOwnLegWrite(async () => || absorbState.ResponseClosed; if (responseStillInFlight) { - sendPacket = false; + // continue below skips the sendPacket check for this frame. connectionState.ServerOutboundDeferred.CancelStream(streamId); absorbState.RequestClosed = true; if (absorbState.IsClosed) @@ -2718,7 +2715,7 @@ void EnqueueNonRelayWire(ReadOnlySpan payload, bool endStreamFlag) async ValueTask RejectDeferredCapAsync() { ReportException(logger, new ProxyHttpException( - "HTTP/2 deferred DATA queue exceeded its per-stream cap.", null, args)); + DeferredDataQueueCapExceededMessage, null, args)); await lockedOwnLegWrite(() => SendRstStreamAsync(new Http2FrameHeader(), new byte[9], streamId, Http2ErrorCode.EnhanceYourCalm, input)); // Cap overflow must drop the stream. Leaving it tracked stalls the peer diff --git a/src/Titanium.Web.Proxy/Http2/Http2Helper.Send.cs b/src/Titanium.Web.Proxy/Http2/Http2Helper.Send.cs index 29f36e403..ba61fcfb0 100644 --- a/src/Titanium.Web.Proxy/Http2/Http2Helper.Send.cs +++ b/src/Titanium.Web.Proxy/Http2/Http2Helper.Send.cs @@ -115,7 +115,7 @@ private static void QueueDataFrame(Http2ConnectionState connectionState, Stream /// Off-loop variant used by : may await flow-control credit because it /// does not run on the shared frame reader. /// - private static async ValueTask QueueSendDataAsync(Http2ConnectionState connectionState, bool towardServer, + private static async ValueTask QueueSendDataAsync(Http2ConnectionState connectionState, bool towardServer, // NOSONAR S107 -- DATA send keeps frame fields explicit; a context object would allocate on this path. SemaphoreSlim writeLock, int streamId, ReadOnlyMemory data, bool endStream, int maxFrameSize, Http2FlowController flow, Stream output, CancellationToken cancellationToken) { @@ -154,7 +154,7 @@ private static async ValueTask QueueSendDataAsync(Http2ConnectionState connectio /// must reset the stream (silently dropping them would hang the peer). /// /// - private static QueueSendDataResult QueueSendData(Http2ConnectionState connectionState, bool towardServer, + private static QueueSendDataResult QueueSendData(Http2ConnectionState connectionState, bool towardServer, // NOSONAR S107, S3776 -- Per-frame DATA queue shares flow-control state; splitting or boxing the args would change the send path. SemaphoreSlim writeLock, int streamId, ReadOnlyMemory data, bool endStream, int maxFrameSize, Http2FlowController flow, Http2DeferredOutboundData deferred, Stream output) { @@ -271,7 +271,7 @@ internal static void QueueRstStreamFrame(Http2ConnectionState connectionState, S /// would let this block overtake earlier-encoded HEADERS still queued on the dedicated writer, /// and the peer's HPACK decoder would desynchronize (COMPRESSION_ERROR). /// - internal static void QueueSendTrailer(Http2ConnectionState connectionState, bool towardServer, + internal static void QueueSendTrailer(Http2ConnectionState connectionState, bool towardServer, // NOSONAR S107 -- Trailer encode keeps the same explicit frame fields as QueueSendHeader. SemaphoreSlim writeLock, Http2Settings settings, Http2FrameHeader frameHeader, byte[] frameHeaderBuffer, int streamId, HeaderCollection trailingHeaders, bool endStream, Stream output) { @@ -433,7 +433,7 @@ private static ValueTask WriteHeaderBlockAsync(Http2FrameHeader frameHeader, byt hasPriority, data, maxFrameSize, output); } - private static async ValueTask WriteHeaderBlockMultiAsync(Http2FrameHeader frameHeader, + private static async ValueTask WriteHeaderBlockMultiAsync(Http2FrameHeader frameHeader, // NOSONAR S107 -- CONTINUATION writer keeps the same explicit frame fields as WriteHeaderBlockAsync. byte[] frameHeaderBuffer, int streamId, Http2FrameType type, bool endStream, bool hasPriority, ReadOnlyMemory data, int maxFrameSize, Stream output) { @@ -878,7 +878,7 @@ internal static void EnqueueControlFrame(Http2FrameWriter writer, Http2FrameType writer.EnqueueRented(rented, total); } - private static ValueTask AsValueTask(Task task) => new(task); + private static ValueTask AsValueTask(Task task) => new(task); // NOSONAR S1144 -- Reflection test seam for the completed-task ValueTask wrapper. /// /// Writes two buffers back-to-back without an async state machine when both complete synchronously. @@ -954,7 +954,7 @@ internal static async Task EmitInterimResponseAsync(SessionEventArgs args, int s /// HTTP/2 frames the body with DATA/END_STREAM (Transfer-Encoding is never used over h2), so the /// chunked header is always stripped regardless of which shape applies. /// - internal static async Task EmitSyntheticResponseAsync(SessionEventArgs args, int streamId, + internal static async Task EmitSyntheticResponseAsync(SessionEventArgs args, int streamId, // NOSONAR S107 -- Synthetic response keeps the caller’s frame buffers explicit. Http2ConnectionState connectionState, Stream clientStream, CancellationToken cancellationToken, Func? onAfterResponse = null, ILogger? logger = null, ReadOnlyMemory? wireBody = null) @@ -1075,7 +1075,7 @@ await EmitBufferedSyntheticResponseAsync(response, streamId, connectionState, fr /// Emits HEADERS + DATA from an already-buffered wire body without requiring /// (coalesce path for unread MITM lite / reverse). /// - private static async Task EmitBufferedSyntheticResponseAsync(Response response, int streamId, + private static async Task EmitBufferedSyntheticResponseAsync(Response response, int streamId, // NOSONAR S107 -- Buffered synthetic emit keeps the caller’s frame buffers explicit. Http2ConnectionState connectionState, Http2FrameHeader frameHeader, byte[] frameHeaderBuffer, Stream clientStream, ReadOnlyMemory wireBody, CancellationToken cancellationToken) { diff --git a/src/Titanium.Web.Proxy/Http2/Http2Helper.cs b/src/Titanium.Web.Proxy/Http2/Http2Helper.cs index 1d10b9c07..84308a69f 100644 --- a/src/Titanium.Web.Proxy/Http2/Http2Helper.cs +++ b/src/Titanium.Web.Proxy/Http2/Http2Helper.cs @@ -46,6 +46,8 @@ internal partial class Http2Helper private static readonly byte[] ConnectMethodBytes = "CONNECT"u8.ToArray(); private const string SyntheticResponseFailedMessage = "HTTP/2 synthetic response failed"; + private const string DeferredDataQueueCapExceededMessage = + "HTTP/2 deferred DATA queue exceeded its per-stream cap."; /// /// Connection-level WINDOW_UPDATE increment matching Chrome/Edge (0xEF0001). Grows the peer's diff --git a/src/Titanium.Web.Proxy/Http2/Http2OriginConnection.cs b/src/Titanium.Web.Proxy/Http2/Http2OriginConnection.cs index f7a4f0842..5b5c8f792 100644 --- a/src/Titanium.Web.Proxy/Http2/Http2OriginConnection.cs +++ b/src/Titanium.Web.Proxy/Http2/Http2OriginConnection.cs @@ -1317,44 +1317,69 @@ private void ApplySettings(ReadOnlySpan payload) { var identifier = (payload[i] << 8) | payload[i + 1]; var value = (int)BinaryPrimitives.ReadUInt32BigEndian(payload.Slice(i + 2, 4)); + if (!TryApplySetting(identifier, value)) + return; + } + } - if (identifier == (int)Http2SettingsId.HeaderTableSize) - originSettings.UpdateHeaderTableSize(value); - else if (identifier == (int)Http2SettingsId.MaxFrameSize) - { - // RFC 7540 §6.5.2: values outside [16384, 16777215] are a connection-level PROTOCOL_ERROR. - if (value < 16384 || value > 16777215) - { - Fail(new IOException( - $"HTTP/2 protocol error: SETTINGS_MAX_FRAME_SIZE value {value} is out of range [16384, 16777215].")); - return; - } + /// Returns false when the setting is a connection error and the rest of the frame must be ignored. + private bool TryApplySetting(int identifier, int value) + { + if (identifier == (int)Http2SettingsId.HeaderTableSize) + originSettings.UpdateHeaderTableSize(value); + else if (identifier == (int)Http2SettingsId.MaxFrameSize) + { + if (!TryApplyMaxFrameSize(value)) + return false; + } + else if (identifier == (int)Http2SettingsId.InitialWindowSize) + { + if (!TryApplyInitialWindowSize(value)) + return false; + } + else if (identifier == (int)Http2SettingsId.MaxConcurrentStreams) + originSettings.MaxConcurrentStreams = value; + else if (identifier == (int)Http2SettingsId.EnableConnectProtocol) + ApplyEnableConnectProtocolSetting(value); - originSettings.MaxFrameSize = value; - } - else if (identifier == (int)Http2SettingsId.InitialWindowSize) - { - // RFC 7540 §6.5.2: values above 2^31-1 are a connection-level FLOW_CONTROL_ERROR. - // A wire value > 2^31-1 wraps to a negative int when cast; checking < 0 catches that. - if (value < 0) - { - Fail(new IOException( - $"HTTP/2 protocol error: SETTINGS_INITIAL_WINDOW_SIZE value exceeds the maximum of 2,147,483,647.")); - return; - } + return true; + } - if (sendFlow.OnInitialWindowSizeChanged(value)) - { - Fail(new IOException( - "HTTP/2 protocol error: SETTINGS_INITIAL_WINDOW_SIZE drove a stream window above 2^31-1.")); - return; - } - } - else if (identifier == (int)Http2SettingsId.MaxConcurrentStreams) - originSettings.MaxConcurrentStreams = value; - else if (identifier == (int)Http2SettingsId.EnableConnectProtocol) - ApplyEnableConnectProtocolSetting(value); + /// RFC 7540 §6.5.2: values outside [16384, 16777215] are a connection-level PROTOCOL_ERROR. + private bool TryApplyMaxFrameSize(int value) + { + if (value < 16384 || value > 16777215) + { + Fail(new IOException( + $"HTTP/2 protocol error: SETTINGS_MAX_FRAME_SIZE value {value} is out of range [16384, 16777215].")); + return false; } + + originSettings.MaxFrameSize = value; + return true; + } + + /// + /// RFC 7540 §6.5.2: values above 2^31-1 are a connection-level FLOW_CONTROL_ERROR. + /// A wire value > 2^31-1 wraps to a negative int when cast; checking < 0 catches that. + /// + private bool TryApplyInitialWindowSize(int value) + { + if (value < 0) + { + Fail(new IOException( + "HTTP/2 protocol error: SETTINGS_INITIAL_WINDOW_SIZE value exceeds the maximum of 2,147,483,647.")); + return false; + } + + if (sendFlow.OnInitialWindowSizeChanged(value)) + { + Fail(new IOException( + "HTTP/2 protocol error: SETTINGS_INITIAL_WINDOW_SIZE drove a stream window above 2^31-1.")); + return false; + } + + return true; } /// RFC 8441 §3: value MUST be 0 or 1; a sender MUST NOT send 0 after previously sending 1. diff --git a/src/Titanium.Web.Proxy/Http3/H3H1QpackResponseReader.cs b/src/Titanium.Web.Proxy/Http3/H3H1QpackResponseReader.cs index 6610fa039..5a21b66c3 100644 --- a/src/Titanium.Web.Proxy/Http3/H3H1QpackResponseReader.cs +++ b/src/Titanium.Web.Proxy/Http3/H3H1QpackResponseReader.cs @@ -42,23 +42,15 @@ internal readonly struct Result var isChunked = false; var connectionClose = false; - while (reader.TryConsumeHeaderLineFromBuffer(out var emptyLine, out var lineBytes)) - { - if (emptyLine) - return Finish(builder, contentLength, isChunked, connectionClose); - ConsumeLine(lineBytes, builder, ref contentLength, ref isChunked, ref connectionClose, - alsoPopulate); - } + if (TryFinishFromBuffer(reader, builder, ref contentLength, ref isChunked, ref connectionClose, + alsoPopulate, out var buffered)) + return buffered; while (true) { - while (reader.TryConsumeHeaderLineFromBuffer(out var emptyLine, out var lineBytes)) - { - if (emptyLine) - return Finish(builder, contentLength, isChunked, connectionClose); - ConsumeLine(lineBytes, builder, ref contentLength, ref isChunked, ref connectionClose, - alsoPopulate); - } + if (TryFinishFromBuffer(reader, builder, ref contentLength, ref isChunked, ref connectionClose, + alsoPopulate, out buffered)) + return buffered; // Prefer byte-buffer fill + TryConsume on the hot path. When the buffer is full without // an LF (or EOF leaves a partial line), fall back to ReadLine spanning (cold path). @@ -81,6 +73,31 @@ internal readonly struct Result } } + private static bool TryFinishFromBuffer( + HttpStream reader, + QpackEncoder.ResponseBlockBuilder builder, + ref long contentLength, + ref bool isChunked, + ref bool connectionClose, + HeaderCollection? alsoPopulate, + out Result? result) + { + while (reader.TryConsumeHeaderLineFromBuffer(out var emptyLine, out var lineBytes)) + { + if (emptyLine) + { + result = Finish(builder, contentLength, isChunked, connectionClose); + return true; + } + + ConsumeLine(lineBytes, builder, ref contentLength, ref isChunked, ref connectionClose, + alsoPopulate); + } + + result = null; + return false; + } + private static Result Finish( QpackEncoder.ResponseBlockBuilder builder, long contentLength, diff --git a/src/Titanium.Web.Proxy/Http3/Http3Frame.cs b/src/Titanium.Web.Proxy/Http3/Http3Frame.cs index cd3979af3..abe86afe5 100644 --- a/src/Titanium.Web.Proxy/Http3/Http3Frame.cs +++ b/src/Titanium.Web.Proxy/Http3/Http3Frame.cs @@ -132,6 +132,15 @@ public static ValueTask WriteAsync( return WriteLargeAsync(stream, frameType, payload, completeWrites, cancellationToken); } + /// + /// Writes a zero-payload frame (used for GOAWAY and some SETTINGS without parameters). + /// + public static ValueTask WriteAsync( + Stream stream, + ulong frameType, + CancellationToken cancellationToken) + => WriteAsync(stream, frameType, ReadOnlyMemory.Empty, cancellationToken); + /// /// Copies the frame into and does not release it until the /// write has been consumed. Sync completion calls GetResult @@ -240,15 +249,6 @@ private static async ValueTask AwaitWriteAndReturnAsync(ValueTask writeVt, byte[ } } - /// - /// Writes a zero-payload frame (used for GOAWAY and some SETTINGS without parameters). - /// - public static ValueTask WriteAsync( - Stream stream, - ulong frameType, - CancellationToken cancellationToken) - => WriteAsync(stream, frameType, ReadOnlyMemory.Empty, cancellationToken); - /// /// Writes HEADERS then DATA as a single stream write when the combined frames fit a modest /// buffer. Used for already-buffered medium/large bodies (lossy / bodies ≥ 16 KiB) so the diff --git a/src/Titanium.Web.Proxy/Http3/Http3OriginBridge.Http2.cs b/src/Titanium.Web.Proxy/Http3/Http3OriginBridge.Http2.cs index 038327ac6..c8a5082d2 100644 --- a/src/Titanium.Web.Proxy/Http3/Http3OriginBridge.Http2.cs +++ b/src/Titanium.Web.Proxy/Http3/Http3OriginBridge.Http2.cs @@ -57,6 +57,29 @@ internal static async Task ForwardOverHttp2FastAsync( if (request.Authority.Length == 0 && !string.IsNullOrEmpty(request.Host)) request.Authority = request.Host.GetByteString(); + var origin = ResolveFastHttp2Origin(fwd, request, server); + + try + { + var target = new Http2OriginTarget( + origin.Host, origin.Port, origin.ConnectHost, origin.ConnectPort, origin.PoolKey); + var exchange = await SendHttp2OriginFastWithGoAwayRetryAsync( + server, logger, fwd, target, coldOpenSessionFactory, cancellationToken); + + StampForwardedHttp2Response(fwd, request, exchange); + } + finally + { + request.HttpVersion = clientHttpVersion; + } + } + + private readonly record struct FastHttp2Origin( + string Host, int Port, string PoolKey, string? ConnectHost, int? ConnectPort); + + private static FastHttp2Origin ResolveFastHttp2Origin( + H3H2FastForward fwd, Request request, ProxyServer server) + { string? connectHost = null; int? connectPort = null; if (fwd.ProxyEndPoint is TransparentBaseProxyEndPoint transparent @@ -66,75 +89,65 @@ internal static async Task ForwardOverHttp2FastAsync( connectPort = transparent.ForwardPort; } - string host; - int port; - string poolKey; if (fwd.ProxyEndPoint is TransparentBaseProxyEndPoint fastEp && fastEp.CachedH2OriginPoolKey != null && request.Authority.Equals(fastEp.CachedH2OriginAuthority)) { - host = fastEp.CachedH2OriginHost!; - port = fastEp.CachedH2OriginPort; - poolKey = fastEp.CachedH2OriginPoolKey; - } - else - { - (host, port) = ResolveH2OriginAuthority(request); - poolKey = Http2OriginConnectionPool.BuildPoolKey( - server, fwd.ProxyEndPoint, fwd.CustomUpStreamProxy, fwd.UpStreamEndPoint, - host, port, connectHost, connectPort); - if (fwd.ProxyEndPoint is TransparentBaseProxyEndPoint cacheEp) - { - cacheEp.CachedH2OriginAuthority = request.Authority; - cacheEp.CachedH2OriginHost = host; - cacheEp.CachedH2OriginPort = port; - cacheEp.CachedH2OriginPoolKey = poolKey; - } + return new FastHttp2Origin( + fastEp.CachedH2OriginHost!, fastEp.CachedH2OriginPort, fastEp.CachedH2OriginPoolKey, + connectHost, connectPort); } - try + var (host, port) = ResolveH2OriginAuthority(request); + var poolKey = Http2OriginConnectionPool.BuildPoolKey( + server, fwd.ProxyEndPoint, fwd.CustomUpStreamProxy, fwd.UpStreamEndPoint, + host, port, connectHost, connectPort); + if (fwd.ProxyEndPoint is TransparentBaseProxyEndPoint cacheEp) { - var target = new Http2OriginTarget(host, port, connectHost, connectPort, poolKey); - var exchange = await SendHttp2OriginFastWithGoAwayRetryAsync( - server, logger, fwd, target, coldOpenSessionFactory, cancellationToken); + cacheEp.CachedH2OriginAuthority = request.Authority; + cacheEp.CachedH2OriginHost = host; + cacheEp.CachedH2OriginPort = port; + cacheEp.CachedH2OriginPoolKey = poolKey; + } - var response = exchange.Response; - response.HttpVersion = HttpHeader.Version30; - response.RequestMethod = request.Method; - if (response.StreamBodyWriter == null) - { - response.IsBodyRead = true; - response.ContentLength = exchange.Body.Length; - // Interception-off: skip response.Body — emit via Preencoded HEADERS+DATA. - if (fwd.PreencodedQpackHeaders == null) - fwd.PreencodedQpackHeaders = QpackEncoder.EncodeResponse(response, context: null); - fwd.PreencodedBody = exchange.Body; - fwd.PreencodedBodyLength = exchange.Body.Length; - // MITM unchanged-lite seeds Response before the call. Stamp Body onto the - // exchange Response (which replaces that seed) so BeforeResponse / - // FinishMitm fallback can read it — match Tcp ForwardOverTcpFastAsync. - // Reverse (Response null) keeps skip-Body coalesce. - if (fwd.Response != null) - { - response.Body = exchange.Body; - response.BodyIsWireEncoded = true; - response.IsBodyReceived = true; - } - } - // else: leave StreamBodyWriter + IsBodyRead alone so SendResponseAsync streams DATA. + return new FastHttp2Origin(host, port, poolKey, connectHost, connectPort); + } - if (exchange.TrailingHeaders != null && !response.HasTrailingHeaders) + private static void StampForwardedHttp2Response( + H3H2FastForward fwd, Request request, Http2OriginExchange exchange) + { + var response = exchange.Response; + response.HttpVersion = HttpHeader.Version30; + response.RequestMethod = request.Method; + if (response.StreamBodyWriter == null) + { + response.IsBodyRead = true; + response.ContentLength = exchange.Body.Length; + // Interception-off: skip response.Body — emit via Preencoded HEADERS+DATA. + if (fwd.PreencodedQpackHeaders == null) + fwd.PreencodedQpackHeaders = QpackEncoder.EncodeResponse(response, context: null); + fwd.PreencodedBody = exchange.Body; + fwd.PreencodedBodyLength = exchange.Body.Length; + // MITM unchanged-lite seeds Response before the call. Stamp Body onto the + // exchange Response (which replaces that seed) so BeforeResponse / + // FinishMitm fallback can read it — match Tcp ForwardOverTcpFastAsync. + // Reverse (Response null) keeps skip-Body coalesce. + if (fwd.Response != null) { - foreach (var header in exchange.TrailingHeaders) - response.TrailingHeaders.AddHeader(header); + response.Body = exchange.Body; + response.BodyIsWireEncoded = true; + response.IsBodyReceived = true; } - - fwd.Response = response; } - finally + // else: leave StreamBodyWriter + IsBodyRead alone so SendResponseAsync streams DATA. + + if (exchange.TrailingHeaders != null && !response.HasTrailingHeaders) { - request.HttpVersion = clientHttpVersion; + foreach (var header in exchange.TrailingHeaders) + response.TrailingHeaders.AddHeader(header); } + + fwd.Response = response; } private static async Task SendHttp2OriginFastWithGoAwayRetryAsync( diff --git a/src/Titanium.Web.Proxy/Http3/Http3OriginBridge.Tcp.cs b/src/Titanium.Web.Proxy/Http3/Http3OriginBridge.Tcp.cs index 1927fb6d6..9a3b2813c 100644 --- a/src/Titanium.Web.Proxy/Http3/Http3OriginBridge.Tcp.cs +++ b/src/Titanium.Web.Proxy/Http3/Http3OriginBridge.Tcp.cs @@ -689,11 +689,12 @@ async ValueTask CopyResponseAsync() } // Keep request upload live while copying the response (true duplex). - var copyVt = CopyResponseAsync(); + // Each branch consumes its own ValueTask once. The no-upload path awaits + // directly so it does not allocate the Task that WhenAll needs. if (pendingUpload != null) - await Task.WhenAll(pendingUpload, copyVt.AsTask()); + await Task.WhenAll(pendingUpload, CopyResponseAsync().AsTask()); else - await copyVt; + await CopyResponseAsync(); }; } else if (uploadTask != null) diff --git a/src/Titanium.Web.Proxy/Http3/Http3RequestStream.cs b/src/Titanium.Web.Proxy/Http3/Http3RequestStream.cs index cce2cf555..00eb1814c 100644 --- a/src/Titanium.Web.Proxy/Http3/Http3RequestStream.cs +++ b/src/Titanium.Web.Proxy/Http3/Http3RequestStream.cs @@ -366,11 +366,12 @@ SessionEventArgs ColdOpenSessionFactory() var stub = new SessionEventArgs(server, endPoint, nullStream, null, stubCts); stub.IsFastPath = true; stub.CustomUpStreamProxy = fwd.CustomUpStreamProxy; - stub.UpstreamHttpProtocol = mitmUnchangedH3H3 - ? UpstreamHttpProtocol.Http3 - : mitmUnchangedH3H2 - ? UpstreamHttpProtocol.Http2 - : UpstreamHttpProtocol.Http11; + if (mitmUnchangedH3H3) + stub.UpstreamHttpProtocol = UpstreamHttpProtocol.Http3; + else if (mitmUnchangedH3H2) + stub.UpstreamHttpProtocol = UpstreamHttpProtocol.Http2; + else + stub.UpstreamHttpProtocol = UpstreamHttpProtocol.Http11; return stub; } @@ -716,7 +717,7 @@ await SendPreencodedResponseAsync(stream, fwd.PreencodedQpackHeaders, try { if (!cts.IsCancellationRequested) - cts.Cancel(); + await cts.CancelAsync().ConfigureAwait(false); } catch (ObjectDisposedException) { diff --git a/tests/Titanium.Inspector.Tests/DeferredInspectorUiCoverageTests.cs b/tests/Titanium.Inspector.Tests/DeferredInspectorUiCoverageTests.cs index 3717212bc..ce847db69 100644 --- a/tests/Titanium.Inspector.Tests/DeferredInspectorUiCoverageTests.cs +++ b/tests/Titanium.Inspector.Tests/DeferredInspectorUiCoverageTests.cs @@ -11,6 +11,7 @@ namespace Titanium.Inspector.Tests; [TestClass] public class DeferredInspectorUiCoverageTests { + private static readonly int[] DeferredUiOrder = [1, 2]; [TestMethod] public async Task DeferredUi_PresentsOnTheNextLiveTurn_AndPostsWhileTheAppIsUp() { @@ -39,7 +40,7 @@ await session.Dispatch(() => { // StatusText presents a stashed queue once the dispatcher is live. _ = vm.StatusText; - CollectionAssert.AreEqual(new[] { 1, 2 }, order); + CollectionAssert.AreEqual(DeferredUiOrder, order); var inline = 0; vm.QueueLiveUi(() => inline++); diff --git a/tests/Titanium.Inspector.Tests/LiveSessionUpdateUxTests.cs b/tests/Titanium.Inspector.Tests/LiveSessionUpdateUxTests.cs index 7322ad7e8..5e24e3608 100644 --- a/tests/Titanium.Inspector.Tests/LiveSessionUpdateUxTests.cs +++ b/tests/Titanium.Inspector.Tests/LiveSessionUpdateUxTests.cs @@ -117,7 +117,7 @@ public void OpaqueReason_RaisesInpc() } [TestMethod] - public void SessionUpdated_BeforeBatchedCapture_DoesNotDuplicateGridRow() + public async Task SessionUpdated_BeforeBatchedCapture_DoesNotDuplicateGridRow() { var path = Path.Combine(Path.GetTempPath(), "twp-dup-race-" + Guid.NewGuid().ToString("N") + ".json"); try @@ -155,7 +155,7 @@ public void SessionUpdated_BeforeBatchedCapture_DoesNotDuplicateGridRow() var deadline = DateTime.UtcNow.AddSeconds(2); while (vm.Sessions.Count == 0 && DateTime.UtcNow < deadline) { - Thread.Sleep(20); + await Task.Delay(20); } Assert.AreEqual(1, vm.Sessions.Count); diff --git a/tests/Titanium.Web.Proxy.IntegrationTests/Http2ContinuationAndAbuseBudgetTests.cs b/tests/Titanium.Web.Proxy.IntegrationTests/Http2ContinuationAndAbuseBudgetTests.cs index fe6b92ddd..9a1812f95 100644 --- a/tests/Titanium.Web.Proxy.IntegrationTests/Http2ContinuationAndAbuseBudgetTests.cs +++ b/tests/Titanium.Web.Proxy.IntegrationTests/Http2ContinuationAndAbuseBudgetTests.cs @@ -24,6 +24,8 @@ namespace Titanium.Web.Proxy.IntegrationTests; [TestClass] public class Http2ContinuationAndAbuseBudgetTests { + private static readonly int[] StreamsWithinAckedCap = [1, 3]; + private static readonly int[] StreamPastAckedCap = [5]; private static TestServer sharedServer = null!; [ClassInitialize] @@ -455,7 +457,7 @@ async Task OpenGetAsync(int streamId) return null; } - var earlyRefusal = await WaitForRstAsync(new[] { 1, 3 }, TimeSpan.FromMilliseconds(400)); + var earlyRefusal = await WaitForRstAsync(StreamsWithinAckedCap, TimeSpan.FromMilliseconds(400)); Assert.IsNull(earlyRefusal, "Streams within the ACKed cap must not be refused because a later SETTINGS lowered the cap."); @@ -463,7 +465,7 @@ await rawClient.Connection.WriteFrameAsync(Http2FrameType.Settings, 0, Http2Fram Array.Empty()); await OpenGetAsync(5); - var refused = await WaitForRstAsync(new[] { 5 }, TimeSpan.FromSeconds(3)); + var refused = await WaitForRstAsync(StreamPastAckedCap, TimeSpan.FromSeconds(3)); Assert.AreEqual(5, refused, "After the second SETTINGS is ACKed, a new stream past that cap must be REFUSED_STREAM."); diff --git a/tests/Titanium.Web.Proxy.UnitTests/HandlerAndProtocolHelperCoverageTests.cs b/tests/Titanium.Web.Proxy.UnitTests/HandlerAndProtocolHelperCoverageTests.cs index 0f895ad8d..c8c27a872 100644 --- a/tests/Titanium.Web.Proxy.UnitTests/HandlerAndProtocolHelperCoverageTests.cs +++ b/tests/Titanium.Web.Proxy.UnitTests/HandlerAndProtocolHelperCoverageTests.cs @@ -691,7 +691,7 @@ public async Task WriteTerminateLiteMiddlewareResponse_ConnectionCloseWhenClient var req = new Request { Method = "GET", HttpVersion = HttpHeader.Version11 }; req.Headers.AddHeader("Connection", "close"); - var writeMw = typeof(ProxyServer).GetMethod("WriteTerminateLiteMiddlewareResponseAsync", PrivateStatic)!; + var writeMw = typeof(ProxyServer).GetMethod("WriteTerminateLiteMiddlewareResponseAsync", PrivateInstance)!; var ctx = new ProxyMiddlewareContext { Session = new object(), @@ -704,7 +704,7 @@ public async Task WriteTerminateLiteMiddlewareResponse_ConnectionCloseWhenClient var buf = new byte[2048]; try { accepted.Receive(buf); } catch { /* ignore */ } }); - await (Task)writeMw.Invoke(null, [clientStream, req, ctx, CancellationToken.None])!; + await (Task)writeMw.Invoke(proxy, [clientStream, req, ctx, CancellationToken.None])!; await drain; Assert.AreEqual(403, ctx.HandledStatusCode); Assert.IsTrue(ctx.IsHandled); diff --git a/tests/Titanium.Web.Proxy.UnitTests/Http2BufferedSendOrderingTests.cs b/tests/Titanium.Web.Proxy.UnitTests/Http2BufferedSendOrderingTests.cs index b9b1418c3..9d564c604 100644 --- a/tests/Titanium.Web.Proxy.UnitTests/Http2BufferedSendOrderingTests.cs +++ b/tests/Titanium.Web.Proxy.UnitTests/Http2BufferedSendOrderingTests.cs @@ -25,6 +25,7 @@ namespace Titanium.Web.Proxy.UnitTests; [TestClass] public class Http2BufferedSendOrderingTests { + private static readonly int[] ExpectedOriginHeaderStreamIds = [1, 3]; private sealed record Frame(byte Type, byte Flags, int StreamId, byte[] Payload); private sealed class Capture : Http2.Hpack.IHeaderListener @@ -175,7 +176,7 @@ public async Task SendBufferedRequestAfterAdmission_LaterStreamStartedFirst_Reac .Where(f => f.Type == (byte)Http2FrameType.Headers) .Select(f => f.StreamId) .ToArray(); - CollectionAssert.AreEqual(new[] { 1, 3 }, streamOrder, + CollectionAssert.AreEqual(ExpectedOriginHeaderStreamIds, streamOrder, "Origin must see HEADERS in increasing stream-id order."); var frames = ParseFrames(origin.ToArray()); diff --git a/tests/Titanium.Web.Proxy.UnitTests/SonarGateCoverageBumpTests.cs b/tests/Titanium.Web.Proxy.UnitTests/SonarGateCoverageBumpTests.cs index 3c5e43417..353155275 100644 --- a/tests/Titanium.Web.Proxy.UnitTests/SonarGateCoverageBumpTests.cs +++ b/tests/Titanium.Web.Proxy.UnitTests/SonarGateCoverageBumpTests.cs @@ -712,7 +712,10 @@ public async Task ReadRequestLine_PrefetchHttp10BlankEofAndCancel() { var blank = await s.ReadRequestLineWithResultAsync(); Assert.IsFalse(blank.Cancelled); + // Method is annotated non-nullable but a blank request line leaves it null. +#pragma warning disable MSTEST0025 Assert.IsNull(blank.Status.Method); +#pragma warning restore MSTEST0025 var next = await s.ReadRequestLine(); Assert.AreEqual("GET", next.Method); } @@ -720,7 +723,10 @@ public async Task ReadRequestLine_PrefetchHttp10BlankEofAndCancel() await using (var s = MakeClientStream([])) { var eof = await s.ReadRequestLine(); + // Method is annotated non-nullable but EOF leaves it null. +#pragma warning disable MSTEST0025 Assert.IsNull(eof.Method); +#pragma warning restore MSTEST0025 } var gate = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); diff --git a/tests/Titanium.Web.Proxy.UnitTests/SonarNewCodeCoverageTests.cs b/tests/Titanium.Web.Proxy.UnitTests/SonarNewCodeCoverageTests.cs index ca5e8d47d..249e9a119 100644 --- a/tests/Titanium.Web.Proxy.UnitTests/SonarNewCodeCoverageTests.cs +++ b/tests/Titanium.Web.Proxy.UnitTests/SonarNewCodeCoverageTests.cs @@ -1741,7 +1741,7 @@ public async Task H1TerminateLite_MiddlewareAndLiveOriginForward() var clientStream = new HttpClientStream(proxy, clientConn, new NetworkStream(clientSock, ownsSocket: false), proxy.BufferPool, CancellationToken.None); - var writeMw = typeof(ProxyServer).GetMethod("WriteTerminateLiteMiddlewareResponseAsync", PrivateStatic)!; + var writeMw = typeof(ProxyServer).GetMethod("WriteTerminateLiteMiddlewareResponseAsync", PrivateInstance)!; var ctx = new ProxyMiddlewareContext { Session = new object(), @@ -1755,16 +1755,16 @@ public async Task H1TerminateLite_MiddlewareAndLiveOriginForward() var buf = new byte[1024]; try { accepted.Receive(buf); } catch { /* ignore */ } }); - await (Task)writeMw.Invoke(null, [clientStream, req, ctx, CancellationToken.None])!; + await (Task)writeMw.Invoke(proxy, [clientStream, req, ctx, CancellationToken.None])!; await drain; - var tryMw = typeof(ProxyServer).GetMethod("TryRunTerminateLiteMiddlewareAsync", PrivateStatic)!; + var tryMw = typeof(ProxyServer).GetMethod("TryRunTerminateLiteMiddlewareAsync", PrivateInstance)!; var handled = new HandleAllMiddleware(); - var keep = await (Task)tryMw.Invoke(null, + var keep = await (Task)tryMw.Invoke(proxy, [ep, clientStream, req, new IProxyMiddleware[] { handled }, CancellationToken.None])!; Assert.IsNotNull(keep); - var passthrough = await (Task)tryMw.Invoke(null, + var passthrough = await (Task)tryMw.Invoke(proxy, [ep, clientStream, req, Array.Empty(), CancellationToken.None])!; Assert.IsNull(passthrough); From 68147e5127b073ce33635aebf94bf96fed425354 Mon Sep 17 00:00:00 2001 From: justcoding121 Date: Sun, 4 Oct 2026 16:49:56 +0530 Subject: [PATCH 2/2] Keep the EnsureBodiesLoadedAsync overloads adjacent. --- .../Services/SessionStore.cs | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/src/Titanium.Inspector/Services/SessionStore.cs b/src/Titanium.Inspector/Services/SessionStore.cs index 380c88605..ffa5f7f6c 100644 --- a/src/Titanium.Inspector/Services/SessionStore.cs +++ b/src/Titanium.Inspector/Services/SessionStore.cs @@ -373,6 +373,15 @@ public async Task EnsureBodiesLoadedAsync(SessionSnapshot snapshot, Cancellation } } + public async Task EnsureBodiesLoadedAsync(IEnumerable snapshots, CancellationToken ct = default) + { + foreach (var snap in snapshots) + { + ct.ThrowIfCancellationRequested(); + await EnsureBodiesLoadedAsync(snap, ct).ConfigureAwait(false); + } + } + /// /// Under heavy capture, the spill writer may still be draining thousands of HARs. /// Keep waiting while work is queued; only mark missing once the channel is idle @@ -415,15 +424,6 @@ private async Task WaitForSpilledBodyAsync( return false; } - public async Task EnsureBodiesLoadedAsync(IEnumerable snapshots, CancellationToken ct = default) - { - foreach (var snap in snapshots) - { - ct.ThrowIfCancellationRequested(); - await EnsureBodiesLoadedAsync(snap, ct).ConfigureAwait(false); - } - } - /// /// For export: hydrate spilled bodies, invoke , then unload unless selected. ///