From 1a25475774c95ece835f10e3eb6a2ce41b5f36ea Mon Sep 17 00:00:00 2001 From: Aleksandr Platonenkov Date: Thu, 13 Aug 2026 16:21:02 -0300 Subject: [PATCH 1/3] =?UTF-8?q?perf(client):=20scratch-=D0=B1=D1=83=D1=84?= =?UTF-8?q?=D0=B5=D1=80=20=D1=81=D0=B1=D0=BE=D1=80=D0=BA=D0=B8=20=D1=82?= =?UTF-8?q?=D0=BE=D0=B6=D0=B5=20=D0=B8=D0=B7=20=D0=BF=D1=83=D0=BB=D0=B0=20?= =?UTF-8?q?+=20=D0=B7=D0=B0=D0=BC=D0=B5=D1=87=D0=B0=D0=BD=D0=B8=D1=8F=20Co?= =?UTF-8?q?deRabbit=20=D0=BA=20PR=20#88?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - WebSocketClient: assemblyBuffer брался обычной аллокацией, хотя буфер приёма уже шёл из ArrayPool. Замерено на .NET 10: общий пул возвращает тот же массив после Return и на 2, и на 4 МиБ, так что это снимает последнюю LOH-аллокацию на соединение, а не добавляет косвенности. Старый массив возвращается в пул при росте, текущий — в том же finally, что и буфер приёма - тестовые WS-серверы: framing, хендшейк, приём заголовков и Dispose вынесены в WebSocketTestServerBase. Дублирование здесь уже один раз стрельнуло — кламп offset существовал в BulkMessageServer и отсутствовал в PagedResponseServer, что и поймал предыдущий проход ревью - BulkMessageServer проверяет lengthCycle в конструкторе: цикл отправки идёт отдельной задачей, поэтому длина больше payload вылезала не ошибкой аргумента, а таймаутом приёма через две минуты - CHANGES: «the other six converters» перечисляло девять. Само число верное — WithoutConverter зовут ровно шесть типов; три node-конвертера его не зовут и в эту шестёрку не входят, они проверялись по другому признаку. Формулировка разведена, счёт не менялся --- CHANGES.md | 4 +- Tests/Xrpl.Tests/Client/BulkMessageServer.cs | 202 ++++------------- .../Xrpl.Tests/Client/PagedResponseServer.cs | 188 ++-------------- .../Client/WebSocketTestServerBase.cs | 205 ++++++++++++++++++ Xrpl/Client/WebSocketClient.cs | 18 +- 5 files changed, 284 insertions(+), 333 deletions(-) create mode 100644 Tests/Xrpl.Tests/Client/WebSocketTestServerBase.cs diff --git a/CHANGES.md b/CHANGES.md index e1db87e1..ad9e1f5c 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -5,9 +5,9 @@ * **Fix infinite recursion in `LONFTokenConverter.Write` — the metadata of an NFT transaction could not be serialized at all** — regression introduced in 10.3.0.0 with the `Newtonsoft.Json` → `System.Text.Json` migration; affects every release from 10.3.0.0 on. `JsonSerializer.Serialize(tx.Meta)` threw `JsonException: A possible object cycle was detected` for any transaction whose `AffectedNodes` contain an `NFTokenPage`, which is every `NFTokenMint`, `NFTokenBurn`, `NFTokenAcceptOffer` and `NFTokenModify` that touched a page. Verified against mainnet on all six NFT transaction types — the four above failed, `NFTokenCreateOffer` and `NFTokenCancelOffer` (no page in their metadata) went through: * The converter broke its own recursion the way the other polymorphic converters do — strip itself from `options.Converters` via `JsonSerializerOptionsCache.WithoutConverter` and re-enter the serializer. That works only for a converter that is *registered in the list*. `LONFTokenConverter` is declared as a `[JsonConverter]` **attribute on the `NFToken` type itself** (`LONFTokenPage.cs`), and a converter attached to a type outranks the options list, so System.Text.Json handed the value straight back to `Write` no matter what the list looked like. The frame repeated until the writer hit `MaxDepth`. Raising `MaxDepth` is not a workaround: at 64 and 128 it is a catchable `JsonException`, at 256 the stack overflows and the process dies * `NFToken` has two fields, so `Write` now emits them directly instead of delegating. The wire shape is unchanged — `{"NFToken":{"NFTokenID":"…","URI":"…"}}`, the envelope `Read` already looks for — and the documented null behaviour is preserved by honouring `options.DefaultIgnoreCondition` rather than hard-coding one: `XrplJsonOptions.Default` (`WhenWritingNull`) omits a null `URI`, plain options keep it as `null` - * The other six converters that call `WithoutConverter` were audited against the same two conditions — declared as a type-level attribute **and** re-serializing that same declared type. None hit both. `LOConverter` is registered in the options list (its one attribute use is property-level) and writes the concrete runtime type; `GenericStringConverter`, `MetaBinaryConverter`, `LedgerBinaryConverter` and `TransactionRequestConverter` are only ever attached to properties; the three node converters are type-level but serialize a *different* class (`value.NewFields.GetType()`). `TransactionResponseConverter` is the one other type-level case, and the same trap was already defused there by the `TransactionResponseUnknown` sentinel, so that no value ever carries the annotated type at runtime. Nothing else was changed + * The six other converter types that call `WithoutConverter` — `LOConverter`, `GenericStringConverter`, `MetaBinaryConverter`, `LedgerBinaryConverter`, `TransactionRequestConverter` and `TransactionResponseConverter` — were audited against the same two conditions — declared as a type-level attribute **and** re-serializing that same declared type. None hit both. `LOConverter` is registered in the options list (its one attribute use is property-level) and writes the concrete runtime type; `GenericStringConverter`, `MetaBinaryConverter`, `LedgerBinaryConverter` and `TransactionRequestConverter` are only ever attached to properties. The three node converters (`CreatedNodeConverter`, `ModifiedNodeConverter`, `DeletedNodeConverter`) do not call `WithoutConverter` at all and so are not among those six, but they are type-level and were checked for the same trap anyway: they serialize a *different* class (`value.NewFields.GetType()`). `TransactionResponseConverter` is the one other type-level case, and the same trap was already defused there by the `TransactionResponseUnknown` sentinel, so that no value ever carries the annotated type at runtime. Nothing else was changed * `TestULONFTokenConverter` had `Read` coverage only, which is how the bug survived. It now pins the written shape, both round trips (`URI` set and null), null handling under `XrplJsonOptions.Default` and under plain options, a multi-token `NFTokenPage`, and — the regression test proper — serializing a `Meta` carrying an `NFTokenPage` in `CreatedNode.NewFields`, `ModifiedNode.FinalFields`, `ModifiedNode.PreviousFields` and `DeletedNode.FinalFields`, since the page can arrive in any of them. All offline, on prepared JSON -* **WebSocket message assembly was quadratic in the number of receive chunks** — `ReceiveLoopAsync` grew a multi-chunk message with `byteResult = byteResult.Concat(buffer.Take(result.Count)).ToArray()`. Every chunk allocated a fresh array the size of everything received so far and refilled it one byte at a time through a LINQ enumerator, so a message split into *k* chunks copied roughly `k/2` times its own length; every intermediate array was well past the 85 KB threshold and therefore landed on the uncompacted large object heap. `ledger_data` at `limit=2048` is a few megabytes and arrives in dozens of chunks over a real link, which is exactly where the cost concentrates. Chunks are now `Buffer.BlockCopy`-ed into a scratch buffer that grows to the largest message on the connection and is reused from then on; a message that arrives whole in one chunk skips the scratch entirely, and the receive buffer itself is rented from `ArrayPool` rather than allocated per connection. Measured on a local fragmenting WebSocket server, 300 messages of 2 MiB, allocation per message: 3.50x payload at one chunk, 6.50x at eight, 18.52x at thirty-two — now a flat 3.01x, which is the floor (the exact-sized `byte[]` plus the UTF-16 string handed to the callback). Through the full client stack, a 3000-page `ledger_data` crawl with each page arriving in 32 chunks and a consumer retaining every object: 100.7 s → 55.8 s, 158.4 GiB → 67.3 GiB allocated, 891 → 398 gen2 collections, and the last-decile-to-first-decile page time drops from 1.39x to 1.13x. `ReceiveChunkSize` was measured at 1 MiB and 64 KiB and left at 1 MiB — now that the buffer is pooled, shrinking it changed nothing outside run-to-run noise. `TestUWebSocketMessageAssembly` pins the byte-exactness of a 96-chunk message, that a short message after a long one picks up no stale bytes from the reused buffer, and that per-message allocation stays under 12x payload at 64 chunks (34.5x before the fix). A dead `timedOut` local, declared and tested but never assigned since it appeared, is gone +* **WebSocket message assembly was quadratic in the number of receive chunks** — `ReceiveLoopAsync` grew a multi-chunk message with `byteResult = byteResult.Concat(buffer.Take(result.Count)).ToArray()`. Every chunk allocated a fresh array the size of everything received so far and refilled it one byte at a time through a LINQ enumerator, so a message split into *k* chunks copied roughly `k/2` times its own length; every intermediate array was well past the 85 KB threshold and therefore landed on the uncompacted large object heap. `ledger_data` at `limit=2048` is a few megabytes and arrives in dozens of chunks over a real link, which is exactly where the cost concentrates. Chunks are now `Buffer.BlockCopy`-ed into a scratch buffer that grows to the largest message on the connection and is reused from then on; a message that arrives whole in one chunk skips the scratch entirely, and both the receive buffer and that scratch buffer are rented from `ArrayPool` rather than allocated per connection (measured on .NET 10: the shared pool does hand back the same multi-megabyte array after a return, so this is a real saving and not just indirection). Measured on a local fragmenting WebSocket server, 300 messages of 2 MiB, allocation per message: 3.50x payload at one chunk, 6.50x at eight, 18.52x at thirty-two — now a flat 3.01x, which is the floor (the exact-sized `byte[]` plus the UTF-16 string handed to the callback). Through the full client stack, a 3000-page `ledger_data` crawl with each page arriving in 32 chunks and a consumer retaining every object: 100.7 s → 55.8 s, 158.4 GiB → 67.3 GiB allocated, 891 → 398 gen2 collections, and the last-decile-to-first-decile page time drops from 1.39x to 1.13x. `ReceiveChunkSize` was measured at 1 MiB and 64 KiB and left at 1 MiB — now that the buffer is pooled, shrinking it changed nothing outside run-to-run noise. `TestUWebSocketMessageAssembly` pins the byte-exactness of a 96-chunk message, that a short message after a long one picks up no stale bytes from the reused buffer, and that per-message allocation stays under 12x payload at 64 chunks (34.5x before the fix). A dead `timedOut` local, declared and tested but never assigned since it appeared, is gone * **Request timeout timers outlived their requests** — `RequestManager.Resolve`/`Reject` called `timer.Stop()`. `System.Timers.Timer` derives from `Component` and carries a finalizer, so every completed request left a finalizable object behind, each of them holding its request's serialized text alive through the `Elapsed` closure; over a long paged crawl that is thousands of them. `Dispose()` stops the timer and takes it off the finalization queue. A second, worse case sat next to it: a token that is **already cancelled** runs its `Register` callback inline, so `Reject` completed the request in the middle of the factory method — before the timeout timer existed and therefore with nothing to remove. The factory then registered the timer for a promise that was already gone, and when it fired, `Reject` took its missing-promise early return without removing it, so the entry stayed in `timeoutsAwaitingResponse` for the life of the process. The `CancellationTokenRegistration` leaked on the same path, its assignment to `TaskInfo` happening after `DeletePromise` had already run. Both factories now check whether the promise survived and clean up after themselves; timer removal moved into `DisposeTimeout`, which is also called on the early returns of `Resolve` and `Reject` and so closes the narrow race with a concurrent cancellation as well. `TestURequestManagerCancellation` pins that an already-cancelled token leaves neither timer nor promise behind in either factory, and that a live request still arms its timeout and releases it on completion * **Reflection on the per-response path is gone** — `Resolve`, `Reject` and `ObserveTaskException` reached for `TrySetResult`, `TrySetException` and `Task` through `GetType().GetMethod(...)` + `Invoke` on every single response. `TaskInfo` now carries typed `SetResult`/`SetException` delegates and the `CompletionTask` itself, wired when the request is created. The properties were added rather than substituted: `TaskInfo` is public, so instances built outside `RequestManager` keep the old reflective path * **Dead `tasks` field removed from `XrplClient`** — `private readonly ConcurrentDictionary tasks` was never assigned and never read, so it was permanently null; a leftover from when the client tracked pending requests itself, which `RequestManager` has done for a long time diff --git a/Tests/Xrpl.Tests/Client/BulkMessageServer.cs b/Tests/Xrpl.Tests/Client/BulkMessageServer.cs index 62445f30..2e52f7ef 100644 --- a/Tests/Xrpl.Tests/Client/BulkMessageServer.cs +++ b/Tests/Xrpl.Tests/Client/BulkMessageServer.cs @@ -1,25 +1,19 @@ using System; -using System.Net; using System.Net.Sockets; using System.Text; -using System.Threading; using System.Threading.Tasks; -using Xrpl.Tests.MockRippled; - namespace Xrpl.Tests { /// - /// Minimal WebSocket server that completes a handshake and then pushes a fixed number of - /// large text messages, each split into a controlled number of WebSocket continuation - /// frames. Fragmenting at the protocol level (rather than relying on how the socket happens - /// to slice the stream) makes the number of client-side receive chunks per message exact and - /// reproducible, which is what the assembly path is sensitive to. + /// WebSocket server that pushes a fixed number of large text messages once the client says go, + /// each split into a controlled number of WebSocket continuation frames. Fragmenting at the + /// protocol level (rather than relying on how the socket happens to slice the stream) makes the + /// number of client-side receive chunks per message exact and reproducible, which is what the + /// assembly path is sensitive to. /// - internal sealed class BulkMessageServer : IDisposable + internal sealed class BulkMessageServer : WebSocketTestServerBase { - private readonly TcpListener _listener; - private readonly CancellationTokenSource _cts = new(); private readonly int _messageCount; private readonly int _fragments; private readonly int[] _lengthCycle; @@ -45,17 +39,22 @@ public BulkMessageServer(int messageCount, int payloadBytes, int fragments, int[ _payload = BuildPayload(payloadBytes); _lengthCycle = lengthCycle is { Length: > 0 } ? lengthCycle : new[] { _payload.Length }; - _listener = new TcpListener(IPAddress.Loopback, 0); - _listener.Start(); - Port = ((IPEndPoint)_listener.LocalEndpoint).Port; - _ = AcceptAsync(); - } - - public int Port { get; } + // Checked here rather than left to fail mid-send: the send loop runs detached, so a bad + // length would only surface as the test's receive timeout minutes later. + foreach (int length in _lengthCycle) + { + if (length < 0 || length > _payload.Length) + { + throw new ArgumentOutOfRangeException( + nameof(lengthCycle), + $"length {length} must be between 0 and the payload length {_payload.Length}"); + } + } - public string Url => "ws://127.0.0.1:" + Port + "/"; + StartAccepting(); + } - /// Payload every message carries, as the client should see it. + /// Payload every message is a prefix of, as the client should see it. public string PayloadText => Encoding.UTF8.GetString(_payload); public int PayloadBytes => _payload.Length; @@ -96,76 +95,42 @@ private static byte[] BuildPayload(int payloadBytes) return Encoding.UTF8.GetBytes(builder.ToString(0, payloadBytes)); } - private async Task AcceptAsync() + protected override async Task ServeAsync(NetworkStream stream) { - try + // Wait for the client's go-ahead before pushing anything, so no message can land + // before the caller has opened its measurement window. + byte[] goAhead = new byte[256]; + if (await stream.ReadAsync(goAhead, Token).ConfigureAwait(false) == 0) { - using TcpClient client = await _listener.AcceptTcpClientAsync(_cts.Token).ConfigureAwait(false); - client.NoDelay = true; - NetworkStream stream = client.GetStream(); - - string request = await ReadUntilHeadersEndAsync(stream).ConfigureAwait(false); - string key = Helpers.GetHandshakeRequestKey(request); - byte[] response = Encoding.ASCII.GetBytes(Helpers.GetHandshakeResponse(Helpers.HashKey(key))); - await stream.WriteAsync(response, _cts.Token).ConfigureAwait(false); - await stream.FlushAsync(_cts.Token).ConfigureAwait(false); - - // Wait for the client's go-ahead before pushing anything, so no message can land - // before the caller has opened its measurement window. - byte[] goAhead = new byte[256]; - if (await stream.ReadAsync(goAhead, _cts.Token).ConfigureAwait(false) == 0) - { - _finished.TrySetResult(0); - return; - } - - // Drain and discard whatever the client sends afterwards (keep-alive pings, close - // frames) so its socket never blocks on a full send window. - _ = DrainAsync(stream); - - for (int i = 0; i < _messageCount; i++) - { - int messageLength = _lengthCycle[i % _lengthCycle.Length]; - int fragmentBytes = (messageLength + _fragments - 1) / _fragments; - - for (int fragment = 0; fragment < _fragments; fragment++) - { - // Ceil division can leave trailing frames past the end; they are sent empty - // so the frame count stays exactly as requested. - int offset = Math.Min(fragment * fragmentBytes, messageLength); - int length = Math.Min(fragmentBytes, messageLength - offset); - bool isFirst = fragment == 0; - bool isLast = fragment == _fragments - 1; - - await stream.WriteAsync(BuildFrameHeader(length, isFirst, isLast), _cts.Token) - .ConfigureAwait(false); - await stream.WriteAsync(_payload.AsMemory(offset, length), _cts.Token) - .ConfigureAwait(false); - await stream.FlushAsync(_cts.Token).ConfigureAwait(false); - } - } + _finished.TrySetResult(0); + return; + } - _finished.TrySetResult(_messageCount); + // Drain and discard whatever the client sends afterwards (keep-alive pings, close + // frames) so its socket never blocks on a full send window. + _ = DrainAsync(stream); - await Task.Delay(Timeout.InfiniteTimeSpan, _cts.Token).ConfigureAwait(false); - } - catch (OperationCanceledException) - { - _finished.TrySetCanceled(); - } - catch (Exception ex) + for (int i = 0; i < _messageCount; i++) { - _finished.TrySetException(ex); + int messageLength = _lengthCycle[i % _lengthCycle.Length]; + await WriteFragmentedMessageAsync(stream, _payload.AsMemory(0, messageLength), _fragments) + .ConfigureAwait(false); } + + _finished.TrySetResult(_messageCount); } + protected override void OnCancelled() => _finished.TrySetCanceled(); + + protected override void OnFaulted(Exception error) => _finished.TrySetException(error); + private async Task DrainAsync(NetworkStream stream) { byte[] sink = new byte[4096]; try { - while (await stream.ReadAsync(sink, _cts.Token).ConfigureAwait(false) > 0) + while (await stream.ReadAsync(sink, Token).ConfigureAwait(false) > 0) { } } @@ -174,88 +139,5 @@ private async Task DrainAsync(NetworkStream stream) // The connection going away is the normal end of this loop. } } - - /// - /// Unmasked server-to-client frame header. The first frame of a message carries the text - /// opcode (0x1), every following frame carries continuation (0x0); FIN is set on the last. - /// - private static byte[] BuildFrameHeader(int payloadLength, bool isFirst, bool isLast) - { - byte first = (byte)((isLast ? 0x80 : 0x00) | (isFirst ? 0x01 : 0x00)); - - if (payloadLength <= 125) - { - return new byte[] { first, (byte)payloadLength }; - } - - if (payloadLength <= ushort.MaxValue) - { - return new byte[] - { - first, - 126, - (byte)(payloadLength >> 8), - (byte)payloadLength - }; - } - - return new byte[] - { - first, - 127, - 0, 0, 0, 0, - (byte)(payloadLength >> 24), - (byte)(payloadLength >> 16), - (byte)(payloadLength >> 8), - (byte)payloadLength - }; - } - - private async Task ReadUntilHeadersEndAsync(NetworkStream stream) - { - byte[] buffer = new byte[4096]; - StringBuilder request = new StringBuilder(); - - while (true) - { - int read = await stream.ReadAsync(buffer, _cts.Token).ConfigureAwait(false); - if (read == 0) - { - break; - } - - request.Append(Encoding.ASCII.GetString(buffer, 0, read)); - - if (request.ToString().Contains("\r\n\r\n", StringComparison.Ordinal)) - { - break; - } - } - - return request.ToString(); - } - - public void Dispose() - { - try - { - _cts.Cancel(); - } - catch - { - // best effort - } - - try - { - _listener.Stop(); - } - catch - { - // best effort - } - - _cts.Dispose(); - } } } diff --git a/Tests/Xrpl.Tests/Client/PagedResponseServer.cs b/Tests/Xrpl.Tests/Client/PagedResponseServer.cs index 8162114c..ce23bd38 100644 --- a/Tests/Xrpl.Tests/Client/PagedResponseServer.cs +++ b/Tests/Xrpl.Tests/Client/PagedResponseServer.cs @@ -1,27 +1,20 @@ using System; using System.Buffers.Binary; -using System.IO; -using System.Net; using System.Net.Sockets; using System.Text; using System.Threading; using System.Threading.Tasks; -using Xrpl.Tests.MockRippled; - namespace Xrpl.Tests { /// - /// Minimal WebSocket server that answers every client request with a large, paged - /// rippled-shaped response echoing the request id. Each response is split into a controlled - /// number of WebSocket continuation frames, so the client assembles it from exactly that many - /// receive chunks — which is what a multi-megabyte ledger_data page looks like over a - /// real link. + /// WebSocket server that answers every client request with a large, paged rippled-shaped + /// response echoing the request id. Each response is split into a controlled number of + /// WebSocket continuation frames, so the client assembles it from exactly that many receive + /// chunks — which is what a multi-megabyte ledger_data page looks like over a real link. /// - internal sealed class PagedResponseServer : IDisposable + internal sealed class PagedResponseServer : WebSocketTestServerBase { - private readonly TcpListener _listener; - private readonly CancellationTokenSource _cts = new(); private readonly int _fragments; private readonly string _resultBody; private int _served; @@ -33,16 +26,9 @@ public PagedResponseServer(int approximatePayloadBytes, int fragments) _fragments = Math.Max(1, fragments); _resultBody = BuildResultBody(approximatePayloadBytes); - _listener = new TcpListener(IPAddress.Loopback, 0); - _listener.Start(); - Port = ((IPEndPoint)_listener.LocalEndpoint).Port; - _ = AcceptAsync(); + StartAccepting(); } - public int Port { get; } - - public string Url => "ws://127.0.0.1:" + Port + "/"; - /// Number of requests answered so far. public int Served => Volatile.Read(ref _served); @@ -91,40 +77,23 @@ private static void AppendHex(StringBuilder builder, int seed, int length) } } - private async Task AcceptAsync() + protected override async Task ServeAsync(NetworkStream stream) { - try + while (!Token.IsCancellationRequested) { - using TcpClient client = await _listener.AcceptTcpClientAsync(_cts.Token).ConfigureAwait(false); - client.NoDelay = true; - NetworkStream stream = client.GetStream(); - - string request = await ReadUntilHeadersEndAsync(stream).ConfigureAwait(false); - string key = Helpers.GetHandshakeRequestKey(request); - byte[] response = Encoding.ASCII.GetBytes(Helpers.GetHandshakeResponse(Helpers.HashKey(key))); - await stream.WriteAsync(response, _cts.Token).ConfigureAwait(false); - await stream.FlushAsync(_cts.Token).ConfigureAwait(false); - - while (!_cts.IsCancellationRequested) + string? message = await ReadTextFrameAsync(stream).ConfigureAwait(false); + if (message == null) { - string? message = await ReadTextFrameAsync(stream).ConfigureAwait(false); - if (message == null) - { - return; - } - - string id = ExtractId(message); - await WriteResponseAsync(stream, id).ConfigureAwait(false); - Interlocked.Increment(ref _served); + return; } - } - catch (OperationCanceledException) - { - // Normal shutdown. - } - catch (Exception ex) when (ex is IOException or SocketException or ObjectDisposedException) - { - // The client going away is the normal end of this loop. + + string id = ExtractId(message); + string envelope = "{\"id\":" + id + ",\"status\":\"success\",\"type\":\"response\",\"result\":" + + _resultBody + "}"; + + await WriteFragmentedMessageAsync(stream, Encoding.UTF8.GetBytes(envelope), _fragments) + .ConfigureAwait(false); + Interlocked.Increment(ref _served); } } @@ -164,59 +133,6 @@ private static string ExtractId(string message) return message.Substring(start, stop - start).Trim(); } - private async Task WriteResponseAsync(NetworkStream stream, string id) - { - string envelope = "{\"id\":" + id + ",\"status\":\"success\",\"type\":\"response\",\"result\":" + - _resultBody + "}"; - byte[] payload = Encoding.UTF8.GetBytes(envelope); - int fragmentBytes = (payload.Length + _fragments - 1) / _fragments; - - for (int fragment = 0; fragment < _fragments; fragment++) - { - // Ceil division can leave trailing frames past the end; they are sent empty - // so the frame count stays exactly as requested. - int offset = Math.Min(fragment * fragmentBytes, payload.Length); - int length = Math.Min(fragmentBytes, payload.Length - offset); - bool isFirst = fragment == 0; - bool isLast = fragment == _fragments - 1; - - await stream.WriteAsync(BuildFrameHeader(length, isFirst, isLast), _cts.Token) - .ConfigureAwait(false); - await stream.WriteAsync(payload.AsMemory(offset, length), _cts.Token).ConfigureAwait(false); - await stream.FlushAsync(_cts.Token).ConfigureAwait(false); - } - } - - /// - /// Unmasked server-to-client frame header. The first frame of a message carries the text - /// opcode (0x1), every following frame carries continuation (0x0); FIN is set on the last. - /// - private static byte[] BuildFrameHeader(int payloadLength, bool isFirst, bool isLast) - { - byte first = (byte)((isLast ? 0x80 : 0x00) | (isFirst ? 0x01 : 0x00)); - - if (payloadLength <= 125) - { - return new byte[] { first, (byte)payloadLength }; - } - - if (payloadLength <= ushort.MaxValue) - { - return new byte[] { first, 126, (byte)(payloadLength >> 8), (byte)payloadLength }; - } - - return new byte[] - { - first, - 127, - 0, 0, 0, 0, - (byte)(payloadLength >> 24), - (byte)(payloadLength >> 16), - (byte)(payloadLength >> 8), - (byte)payloadLength - }; - } - /// /// Reads one client frame. Returns the decoded text of the first text frame seen, or null /// once the peer closes. Control frames other than Close are skipped. @@ -287,71 +203,5 @@ private static byte[] BuildFrameHeader(int payloadLength, bool isFirst, bool isL } } } - - private async Task ReadExactAsync(NetworkStream stream, byte[] buffer, int count) - { - int read = 0; - - while (read < count) - { - int chunk = await stream.ReadAsync(buffer.AsMemory(read, count - read), _cts.Token) - .ConfigureAwait(false); - if (chunk == 0) - { - return false; - } - - read += chunk; - } - - return true; - } - - private async Task ReadUntilHeadersEndAsync(NetworkStream stream) - { - byte[] buffer = new byte[4096]; - StringBuilder request = new StringBuilder(); - - while (true) - { - int read = await stream.ReadAsync(buffer, _cts.Token).ConfigureAwait(false); - if (read == 0) - { - break; - } - - request.Append(Encoding.ASCII.GetString(buffer, 0, read)); - - if (request.ToString().Contains("\r\n\r\n", StringComparison.Ordinal)) - { - break; - } - } - - return request.ToString(); - } - - public void Dispose() - { - try - { - _cts.Cancel(); - } - catch - { - // best effort - } - - try - { - _listener.Stop(); - } - catch - { - // best effort - } - - _cts.Dispose(); - } } } diff --git a/Tests/Xrpl.Tests/Client/WebSocketTestServerBase.cs b/Tests/Xrpl.Tests/Client/WebSocketTestServerBase.cs new file mode 100644 index 00000000..2ba1f0e4 --- /dev/null +++ b/Tests/Xrpl.Tests/Client/WebSocketTestServerBase.cs @@ -0,0 +1,205 @@ +using System; +using System.Net; +using System.Net.Sockets; +using System.Text; +using System.Threading; +using System.Threading.Tasks; + +using Xrpl.Tests.MockRippled; + +namespace Xrpl.Tests +{ + /// + /// Shared plumbing for the raw WebSocket servers the client tests drive: a loopback listener, + /// the HTTP upgrade handshake, frame framing and disposal. The framing in particular lives here + /// on purpose — it used to be copied per server, and the copies drifted. + /// + internal abstract class WebSocketTestServerBase : IDisposable + { + private readonly TcpListener _listener; + + protected WebSocketTestServerBase() + { + _listener = new TcpListener(IPAddress.Loopback, 0); + _listener.Start(); + Port = ((IPEndPoint)_listener.LocalEndpoint).Port; + } + + /// Token cancelled when the server is disposed. + protected CancellationToken Token => Cts.Token; + + protected CancellationTokenSource Cts { get; } = new CancellationTokenSource(); + + public int Port { get; } + + public string Url => "ws://127.0.0.1:" + Port + "/"; + + /// + /// Starts accepting. Call at the end of the derived constructor, once its own state is set: + /// the accept loop runs on the thread pool and may reach at any time. + /// + protected void StartAccepting() + { + _ = AcceptAsync(); + } + + /// Serves one connected client; the handshake has already completed. + protected abstract Task ServeAsync(NetworkStream stream); + + /// Called when the accept loop ends with an exception the server did not expect. + protected virtual void OnFaulted(Exception error) + { + } + + /// Called when the accept loop ends because the server was disposed. + protected virtual void OnCancelled() + { + } + + private async Task AcceptAsync() + { + try + { + using TcpClient client = await _listener.AcceptTcpClientAsync(Token).ConfigureAwait(false); + client.NoDelay = true; + NetworkStream stream = client.GetStream(); + + string request = await ReadUntilHeadersEndAsync(stream).ConfigureAwait(false); + string key = Helpers.GetHandshakeRequestKey(request); + byte[] response = Encoding.ASCII.GetBytes(Helpers.GetHandshakeResponse(Helpers.HashKey(key))); + await stream.WriteAsync(response, Token).ConfigureAwait(false); + await stream.FlushAsync(Token).ConfigureAwait(false); + + await ServeAsync(stream).ConfigureAwait(false); + + // Hold the connection open until the test disposes the server. + await Task.Delay(Timeout.InfiniteTimeSpan, Token).ConfigureAwait(false); + } + catch (OperationCanceledException) + { + OnCancelled(); + } + catch (Exception ex) + { + OnFaulted(ex); + } + } + + /// + /// Unmasked server-to-client frame header. The first frame of a message carries the text + /// opcode (0x1), every following frame carries continuation (0x0); FIN is set on the last. + /// + protected static byte[] BuildFrameHeader(int payloadLength, bool isFirst, bool isLast) + { + byte first = (byte)((isLast ? 0x80 : 0x00) | (isFirst ? 0x01 : 0x00)); + + if (payloadLength <= 125) + { + return new byte[] { first, (byte)payloadLength }; + } + + if (payloadLength <= ushort.MaxValue) + { + return new byte[] { first, 126, (byte)(payloadLength >> 8), (byte)payloadLength }; + } + + return new byte[] + { + first, + 127, + 0, 0, 0, 0, + (byte)(payloadLength >> 24), + (byte)(payloadLength >> 16), + (byte)(payloadLength >> 8), + (byte)payloadLength + }; + } + + /// + /// Writes one message split into frames. Ceil division can + /// leave trailing frames past the end of the payload; they are sent empty so the frame + /// count is exactly what the caller asked for. + /// + protected async Task WriteFragmentedMessageAsync(NetworkStream stream, ReadOnlyMemory payload, int fragments) + { + int fragmentBytes = (payload.Length + fragments - 1) / fragments; + + for (int fragment = 0; fragment < fragments; fragment++) + { + int offset = Math.Min(fragment * fragmentBytes, payload.Length); + int length = Math.Min(fragmentBytes, payload.Length - offset); + bool isFirst = fragment == 0; + bool isLast = fragment == fragments - 1; + + await stream.WriteAsync(BuildFrameHeader(length, isFirst, isLast), Token).ConfigureAwait(false); + await stream.WriteAsync(payload.Slice(offset, length), Token).ConfigureAwait(false); + await stream.FlushAsync(Token).ConfigureAwait(false); + } + } + + protected async Task ReadExactAsync(NetworkStream stream, byte[] buffer, int count) + { + int read = 0; + + while (read < count) + { + int chunk = await stream.ReadAsync(buffer.AsMemory(read, count - read), Token).ConfigureAwait(false); + if (chunk == 0) + { + return false; + } + + read += chunk; + } + + return true; + } + + private async Task ReadUntilHeadersEndAsync(NetworkStream stream) + { + byte[] buffer = new byte[4096]; + StringBuilder request = new StringBuilder(); + + while (true) + { + int read = await stream.ReadAsync(buffer, Token).ConfigureAwait(false); + if (read == 0) + { + break; + } + + request.Append(Encoding.ASCII.GetString(buffer, 0, read)); + + if (request.ToString().Contains("\r\n\r\n", StringComparison.Ordinal)) + { + break; + } + } + + return request.ToString(); + } + + public void Dispose() + { + try + { + Cts.Cancel(); + } + catch + { + // best effort + } + + try + { + _listener.Stop(); + } + catch + { + // best effort + } + + Cts.Dispose(); + } + } +} diff --git a/Xrpl/Client/WebSocketClient.cs b/Xrpl/Client/WebSocketClient.cs index 8a0c885d..e4ae5693 100644 --- a/Xrpl/Client/WebSocketClient.cs +++ b/Xrpl/Client/WebSocketClient.cs @@ -484,7 +484,9 @@ private async Task ReceiveLoopAsync() // Scratch buffer for messages that arrive in more than one chunk. It grows to the // largest message seen on this connection and is then reused, so steady-state assembly - // costs nothing beyond the single exact-sized array handed to the callbacks. + // costs nothing beyond the single exact-sized array handed to the callbacks. Rented for + // the same reason as the receive buffer above - at these sizes it is a large-object-heap + // array that would otherwise be thrown away per connection. byte[]? assemblyBuffer = null; try @@ -637,6 +639,11 @@ private async Task ReceiveLoopAsync() finally { ArrayPool.Shared.Return(buffer); + + if (assemblyBuffer != null) + { + ArrayPool.Shared.Return(assemblyBuffer); + } } } @@ -658,12 +665,19 @@ private static void EnsureAssemblyCapacity(ref byte[]? assemblyBuffer, int requi capacity = capacity <= Array.MaxLength / 2 ? capacity * 2 : requiredLength; } - byte[] grown = new byte[capacity]; + // Rent may hand back a longer array than asked for; the growth above keys off the + // actual length, so the next doubling starts from what we really got. + byte[] grown = ArrayPool.Shared.Rent(capacity); if (preserveLength > 0) { Buffer.BlockCopy(assemblyBuffer!, 0, grown, 0, preserveLength); } + if (assemblyBuffer != null) + { + ArrayPool.Shared.Return(assemblyBuffer); + } + assemblyBuffer = grown; } From 14730a29f5c4aa4659e747958242b0564e4ff9a9 Mon Sep 17 00:00:00 2001 From: Aleksandr Platonenkov Date: Thu, 13 Aug 2026 16:26:57 -0300 Subject: [PATCH 2/3] =?UTF-8?q?fix(tests):=20=D0=BD=D0=B5=20=D1=82=D0=B5?= =?UTF-8?q?=D1=80=D1=8F=D1=82=D1=8C=20=D0=BE=D1=88=D0=B8=D0=B1=D0=BA=D1=83?= =?UTF-8?q?=20=D1=81=D0=B5=D1=80=D0=B2=D0=B5=D1=80=D0=B0=20=D0=B8=20=D0=BD?= =?UTF-8?q?=D0=B5=20=D0=BE=D1=81=D1=82=D0=B0=D0=B2=D0=BB=D1=8F=D1=82=D1=8C?= =?UTF-8?q?=20=D1=81=D0=BB=D1=83=D1=88=D0=B0=D1=82=D0=B5=D0=BB=D1=8F=20?= =?UTF-8?q?=D0=BF=D1=80=D0=B8=20=D0=BE=D1=82=D0=BA=D0=B0=D0=B7=D0=B5=20?= =?UTF-8?q?=D0=BA=D0=BE=D0=BD=D1=81=D1=82=D1=80=D1=83=D0=BA=D1=82=D0=BE?= =?UTF-8?q?=D1=80=D0=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit По итогам селф-ревью рефакторинга харнесса. - PagedResponseServer раньше глушил IOException/SocketException/ObjectDisposedException сам, а после выноса в базу под общий catch попадали уже любые исключения и молча пропадали. База теперь запоминает их в Fault, BulkMessageServer по-прежнему дополнительно роняет в SendCompleted, а тест сборки показывает Fault в сообщении об ошибке — иначе поломка сервера выглядела бы как таймаут приёма на стороне клиента - конструктор базы уже поднимает слушателя к моменту, когда отрабатывает проверка lengthCycle в наследнике, так что бросок оставлял открытый сокет и неосвобождённый CTS без единого владельца. Перед throw вызывается Dispose --- Tests/Xrpl.Tests/Client/BulkMessageServer.cs | 9 ++++++++- Tests/Xrpl.Tests/Client/TestUWebSocketMessageAssembly.cs | 5 ++++- Tests/Xrpl.Tests/Client/WebSocketTestServerBase.cs | 8 ++++++++ 3 files changed, 20 insertions(+), 2 deletions(-) diff --git a/Tests/Xrpl.Tests/Client/BulkMessageServer.cs b/Tests/Xrpl.Tests/Client/BulkMessageServer.cs index 2e52f7ef..c0653e89 100644 --- a/Tests/Xrpl.Tests/Client/BulkMessageServer.cs +++ b/Tests/Xrpl.Tests/Client/BulkMessageServer.cs @@ -45,6 +45,9 @@ public BulkMessageServer(int messageCount, int payloadBytes, int fragments, int[ { if (length < 0 || length > _payload.Length) { + // The base constructor has already opened the listener, and a constructor that + // throws leaves nobody to dispose it. + Dispose(); throw new ArgumentOutOfRangeException( nameof(lengthCycle), $"length {length} must be between 0 and the payload length {_payload.Length}"); @@ -122,7 +125,11 @@ await WriteFragmentedMessageAsync(stream, _payload.AsMemory(0, messageLength), _ protected override void OnCancelled() => _finished.TrySetCanceled(); - protected override void OnFaulted(Exception error) => _finished.TrySetException(error); + protected override void OnFaulted(Exception error) + { + base.OnFaulted(error); + _finished.TrySetException(error); + } private async Task DrainAsync(NetworkStream stream) { diff --git a/Tests/Xrpl.Tests/Client/TestUWebSocketMessageAssembly.cs b/Tests/Xrpl.Tests/Client/TestUWebSocketMessageAssembly.cs index 4a4eee72..973388bc 100644 --- a/Tests/Xrpl.Tests/Client/TestUWebSocketMessageAssembly.cs +++ b/Tests/Xrpl.Tests/Client/TestUWebSocketMessageAssembly.cs @@ -124,7 +124,10 @@ private static async Task> ReceiveAsync(BulkMessageServer Task completed = await Task.WhenAny(done.Task, Task.Delay(TimeSpan.FromSeconds(WaitSeconds))) .ConfigureAwait(false); - Assert.AreSame(done.Task, completed, $"only {messages.Count} of {expected} messages arrived"); + Assert.AreSame( + done.Task, + completed, + $"only {messages.Count} of {expected} messages arrived; server fault: {server.Fault?.ToString() ?? "none"}"); } finally { diff --git a/Tests/Xrpl.Tests/Client/WebSocketTestServerBase.cs b/Tests/Xrpl.Tests/Client/WebSocketTestServerBase.cs index 2ba1f0e4..b7b05085 100644 --- a/Tests/Xrpl.Tests/Client/WebSocketTestServerBase.cs +++ b/Tests/Xrpl.Tests/Client/WebSocketTestServerBase.cs @@ -17,6 +17,7 @@ namespace Xrpl.Tests internal abstract class WebSocketTestServerBase : IDisposable { private readonly TcpListener _listener; + private Exception? _fault; protected WebSocketTestServerBase() { @@ -46,9 +47,16 @@ protected void StartAccepting() /// Serves one connected client; the handshake has already completed. protected abstract Task ServeAsync(NetworkStream stream); + /// + /// The exception that ended the accept loop, if it ended badly. Recorded rather than + /// swallowed so a broken server shows up as itself instead of as the caller's timeout. + /// + public Exception? Fault => Volatile.Read(ref _fault); + /// Called when the accept loop ends with an exception the server did not expect. protected virtual void OnFaulted(Exception error) { + Volatile.Write(ref _fault, error); } /// Called when the accept loop ends because the server was disposed. From 319592e8e38872e6b7682a1c2fc2006697fde15a Mon Sep 17 00:00:00 2001 From: Aleksandr Platonenkov Date: Thu, 13 Aug 2026 16:47:32 -0300 Subject: [PATCH 3/3] =?UTF-8?q?fix(tests):=20=D0=BD=D0=B5=20=D1=87=D0=B8?= =?UTF-8?q?=D1=82=D0=B0=D1=82=D1=8C=20Token=20=D1=83=20=D1=83=D0=B6=D0=B5?= =?UTF-8?q?=20=D0=BE=D1=81=D0=B2=D0=BE=D0=B1=D0=BE=D0=B6=D0=B4=D1=91=D0=BD?= =?UTF-8?q?=D0=BD=D0=BE=D0=B3=D0=BE=20CTS=20=D0=B2=20=D1=82=D0=B5=D1=81?= =?UTF-8?q?=D1=82=D0=BE=D0=B2=D0=BE=D0=BC=20WS-=D1=81=D0=B5=D1=80=D0=B2?= =?UTF-8?q?=D0=B5=D1=80=D0=B5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Замечания CodeRabbit к PR #90. - Dispose звал Cancel и сразу Dispose, не дожидаясь цикла приёма, а тот и DrainAsync читают Token на каждой итерации. Чтение CancellationTokenSource.Token после Dispose бросает ObjectDisposedException, и она уходила в общий catch — то есть в Fault попадал артефакт разбора стенда вместо настоящей ошибки сервера, а у BulkMessageServer ещё и в SendCompleted. Токен теперь снимается один раз в конструкторе: Dispose всегда отменяет перед освобождением, а по уже отменённому токену ожидания завершаются, не обращаясь к источнику - ReadUntilHeadersEndAsync может прочитать за "\r\n\r\n" и остаток молча теряет. Живого клиента это не задевает — фрейм нельзя слать до 101, — но база теперь общая, поэтому ограничение описано в док-комментарии --- .../Client/WebSocketTestServerBase.cs | 24 +++++++++++++++---- 1 file changed, 19 insertions(+), 5 deletions(-) diff --git a/Tests/Xrpl.Tests/Client/WebSocketTestServerBase.cs b/Tests/Xrpl.Tests/Client/WebSocketTestServerBase.cs index b7b05085..9ac8d07a 100644 --- a/Tests/Xrpl.Tests/Client/WebSocketTestServerBase.cs +++ b/Tests/Xrpl.Tests/Client/WebSocketTestServerBase.cs @@ -17,19 +17,25 @@ namespace Xrpl.Tests internal abstract class WebSocketTestServerBase : IDisposable { private readonly TcpListener _listener; + private readonly CancellationTokenSource _cts = new CancellationTokenSource(); private Exception? _fault; protected WebSocketTestServerBase() { + // Captured once rather than read from the source on every use: Dispose does not wait for + // the accept loop, and reading CancellationTokenSource.Token after Dispose throws + // ObjectDisposedException, which would surface as a bogus Fault during teardown. Dispose + // always cancels before disposing, so a token captured here stays usable afterwards - + // an already cancelled token completes waits without touching its source. + Token = _cts.Token; + _listener = new TcpListener(IPAddress.Loopback, 0); _listener.Start(); Port = ((IPEndPoint)_listener.LocalEndpoint).Port; } /// Token cancelled when the server is disposed. - protected CancellationToken Token => Cts.Token; - - protected CancellationTokenSource Cts { get; } = new CancellationTokenSource(); + protected CancellationToken Token { get; } public int Port { get; } @@ -163,6 +169,14 @@ protected async Task ReadExactAsync(NetworkStream stream, byte[] buffer, i return true; } + /// + /// Reads the client's HTTP upgrade request up to the blank line that ends the headers. + /// Anything the same read happened to pull in after that terminator is returned as part of + /// the string and then dropped by the caller: no client here pipelines a WebSocket frame + /// onto the upgrade request, because it has to wait for the 101 before it may send one. + /// A server that ever needs to accept such a client has to hand the leftover bytes to + /// instead. + /// private async Task ReadUntilHeadersEndAsync(NetworkStream stream) { byte[] buffer = new byte[4096]; @@ -191,7 +205,7 @@ public void Dispose() { try { - Cts.Cancel(); + _cts.Cancel(); } catch { @@ -207,7 +221,7 @@ public void Dispose() // best effort } - Cts.Dispose(); + _cts.Dispose(); } } }