From dd898f9f44b8e5dfc47039542307df6f9c76b415 Mon Sep 17 00:00:00 2001 From: Aleksandr Platonenkov Date: Thu, 13 Aug 2026 15:06:43 -0300 Subject: [PATCH 1/8] =?UTF-8?q?perf(client):=20=D0=BB=D0=B8=D0=BD=D0=B5?= =?UTF-8?q?=D0=B9=D0=BD=D0=B0=D1=8F=20=D1=81=D0=B1=D0=BE=D1=80=D0=BA=D0=B0?= =?UTF-8?q?=20WebSocket-=D1=81=D0=BE=D0=BE=D0=B1=D1=89=D0=B5=D0=BD=D0=B8?= =?UTF-8?q?=D0=B9=20=D0=B2=D0=BC=D0=B5=D1=81=D1=82=D0=BE=20=D0=BA=D0=B2?= =?UTF-8?q?=D0=B0=D0=B4=D1=80=D0=B0=D1=82=D0=B8=D1=87=D0=BD=D0=BE=D0=B9?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ReceiveLoopAsync собирал многочанковое сообщение через byteResult.Concat(buffer.Take(result.Count)).ToArray() — на каждый чанк аллоцировался новый массив во всю накопленную длину, плюс поэлементное LINQ-перечисление вместо блочного копирования. Все промежуточные массивы крупнее 85 КБ уходили в LOH. Теперь чанки копируются Buffer.BlockCopy в scratch-буфер, который растёт до максимального сообщения соединения и дальше переиспользуется; одночанковые сообщения копируются напрямую, минуя scratch. Буфер приёма берётся из ArrayPool, а не аллоцируется на каждое соединение. Замер (300 сообщений по 2 МиБ, аллокации на сообщение): 1 чанк 3.50x payload -> 3.01x 8 чанков 6.50x -> 3.01x 32 чанка 18.52x -> 3.01x Полный стек, 3000 страниц ledger_data по 2 МиБ в 32 чанка с растущим живым heap потребителя: время 100.7 с -> 55.8 с (29.8 -> 53.8 страниц/с) аллокации 158.4 ГиБ -> 67.3 ГиБ (54.1 -> 23.0 МиБ на страницу) сборок gen2 891 -> 398 тренд последней децили к первой 1.39x -> 1.13x Заодно убрана мёртвая переменная timedOut: она объявлялась и проверялась, но никогда не присваивалась с момента появления в 3c7e38e. --- .../Client/BenchmarkLedgerDataCrawl.cs | 177 +++++++++ .../Client/BenchmarkWebSocketAssembly.cs | 161 ++++++++ Tests/Xrpl.Tests/Client/BulkMessageServer.cs | 261 +++++++++++++ .../Xrpl.Tests/Client/PagedResponseServer.cs | 355 ++++++++++++++++++ .../Client/TestUWebSocketMessageAssembly.cs | 135 +++++++ Xrpl/Client/WebSocketClient.cs | 77 +++- 6 files changed, 1153 insertions(+), 13 deletions(-) create mode 100644 Tests/Xrpl.Tests/Client/BenchmarkLedgerDataCrawl.cs create mode 100644 Tests/Xrpl.Tests/Client/BenchmarkWebSocketAssembly.cs create mode 100644 Tests/Xrpl.Tests/Client/BulkMessageServer.cs create mode 100644 Tests/Xrpl.Tests/Client/PagedResponseServer.cs create mode 100644 Tests/Xrpl.Tests/Client/TestUWebSocketMessageAssembly.cs diff --git a/Tests/Xrpl.Tests/Client/BenchmarkLedgerDataCrawl.cs b/Tests/Xrpl.Tests/Client/BenchmarkLedgerDataCrawl.cs new file mode 100644 index 00000000..dc86527a --- /dev/null +++ b/Tests/Xrpl.Tests/Client/BenchmarkLedgerDataCrawl.cs @@ -0,0 +1,177 @@ +using Microsoft.VisualStudio.TestTools.UnitTesting; + +using System; +using System.Collections.Generic; +using System.Diagnostics; +using System.Threading.Tasks; + +using Xrpl.Client; + +namespace Xrpl.Tests.ClientLib +{ + /// + /// Manual benchmark of a long paged crawl through the full client stack (socket receive loop, + /// message routing, RequestManager, JSON round-trip). Deliberately named outside the + /// TestU/TestI filters so it never runs in CI; invoke it explicitly: + /// dotnet test --filter "FullyQualifiedName~BenchmarkLedgerDataCrawl". + /// Knobs: CRAWL_PAGES, CRAWL_PAYLOAD_BYTES, CRAWL_FRAGMENTS, CRAWL_RETAIN_PER_PAGE. + /// CRAWL_RETAIN_PER_PAGE models a consumer that keeps every crawled object alive (a full + /// ledger-state snapshot), which is what makes each forced gen2 collection progressively + /// more expensive as the crawl advances. + /// + [TestClass] + public class BenchmarkLedgerDataCrawl + { + private static int EnvInt(string name, int fallback) + { + string? raw = Environment.GetEnvironmentVariable(name); + return int.TryParse(raw, out int value) && value >= 0 ? value : fallback; + } + + /// Stand-in for one crawled ledger object the consumer keeps in its snapshot. + private sealed class RetainedEntry + { + public RetainedEntry(string index) + { + Index = index; + } + + public string Index { get; } + } + + [TestMethod] + public async Task BenchmarkSequentialPaging() + { + int pages = EnvInt("CRAWL_PAGES", 2000); + int payloadBytes = EnvInt("CRAWL_PAYLOAD_BYTES", 2 * 1024 * 1024); + int fragments = EnvInt("CRAWL_FRAGMENTS", 32); + int retainPerPage = EnvInt("CRAWL_RETAIN_PER_PAGE", 0); + + List snapshot = new List(pages * retainPerPage); + + using PagedResponseServer server = new PagedResponseServer(payloadBytes, fragments); + using XrplClient client = new XrplClient(server.Url); + + await client.Connect().ConfigureAwait(false); + + // One warm-up page so JIT and pooled buffers are not charged to the measured window. + await RequestPageAsync(client).ConfigureAwait(false); + + GC.Collect(2, GCCollectionMode.Forced, blocking: true); + GC.WaitForPendingFinalizers(); + GC.Collect(2, GCCollectionMode.Forced, blocking: true); + + long allocatedBefore = GC.GetTotalAllocatedBytes(precise: true); + int gen0Before = GC.CollectionCount(0); + int gen1Before = GC.CollectionCount(1); + int gen2Before = GC.CollectionCount(2); + long heapBefore = GC.GetTotalMemory(false); + long lohBefore = LohBytes(); + + double[] pageMs = new double[pages]; + long startTicks = Stopwatch.GetTimestamp(); + + for (int i = 0; i < pages; i++) + { + long before = Stopwatch.GetTimestamp(); + await RequestPageAsync(client).ConfigureAwait(false); + + for (int entry = 0; entry < retainPerPage; entry++) + { + snapshot.Add(new RetainedEntry(((long)i * retainPerPage + entry).ToString("X16"))); + } + + pageMs[i] = (Stopwatch.GetTimestamp() - before) * 1000.0 / Stopwatch.Frequency; + } + + double totalSeconds = (Stopwatch.GetTimestamp() - startTicks) / (double)Stopwatch.Frequency; + long allocated = GC.GetTotalAllocatedBytes(precise: true) - allocatedBefore; + long heapAfter = GC.GetTotalMemory(false); + long lohAfter = LohBytes(); + + await client.Disconnect().ConfigureAwait(false); + + Console.WriteLine("=== ledger_data crawl benchmark (full client stack) ==="); + Console.WriteLine($"pages : {pages}"); + Console.WriteLine($"payload : {payloadBytes / 1024.0 / 1024.0:F2} MiB"); + Console.WriteLine($"fragments per page: {fragments}"); + Console.WriteLine($"retained objects : {snapshot.Count:N0} ({retainPerPage}/page)"); + Console.WriteLine($"total time : {totalSeconds:F2} s ({pages / totalSeconds:F2} pages/s)"); + Console.WriteLine($"first decile avg : {DecileAverage(pageMs, 0):F1} ms/page"); + Console.WriteLine($"last decile avg : {DecileAverage(pageMs, 9):F1} ms/page"); + Console.WriteLine($"trend (last/first): {DecileAverage(pageMs, 9) / DecileAverage(pageMs, 0):F2}x"); + Console.WriteLine($"p50 / p99 / max : {Percentile(pageMs, 50):F1} / {Percentile(pageMs, 99):F1} / " + + $"{Percentile(pageMs, 100):F1} ms"); + Console.WriteLine($"allocated total : {allocated / 1024.0 / 1024.0 / 1024.0:F2} GiB"); + Console.WriteLine($"allocated per page: {allocated / (double)pages / 1024.0 / 1024.0:F1} MiB " + + $"({allocated / (double)pages / payloadBytes:F1}x payload)"); + Console.WriteLine($"gen0/gen1/gen2 : {GC.CollectionCount(0) - gen0Before} / " + + $"{GC.CollectionCount(1) - gen1Before} / {GC.CollectionCount(2) - gen2Before}"); + Console.WriteLine($"managed heap : {heapBefore / 1024.0 / 1024.0:F1} -> {heapAfter / 1024.0 / 1024.0:F1} MiB"); + Console.WriteLine($"LOH size : {lohBefore / 1024.0 / 1024.0:F1} -> {lohAfter / 1024.0 / 1024.0:F1} MiB"); + Console.WriteLine("decile profile (ms/page): " + string.Join(" | ", DecileProfile(pageMs))); + } + + private static async Task RequestPageAsync(XrplClient client) + { + Dictionary request = new Dictionary + { + ["command"] = "ledger_data", + ["ledger_index"] = 96000000, + ["binary"] = true, + ["limit"] = 2048 + }; + + Dictionary response = await client.Request(request).ConfigureAwait(false); + if (response == null) + { + throw new InvalidOperationException("empty ledger_data response"); + } + } + + private static long LohBytes() + { + GCMemoryInfo info = GC.GetGCMemoryInfo(); + ReadOnlySpan generations = info.GenerationInfo; + return generations.Length > 3 ? generations[3].SizeAfterBytes : 0; + } + + private static double DecileAverage(double[] values, int decile) + { + int size = Math.Max(1, values.Length / 10); + int from = decile * size; + int to = Math.Min(values.Length, from + size); + if (from >= to) + { + return 0; + } + + double sum = 0; + for (int i = from; i < to; i++) + { + sum += values[i]; + } + + return sum / (to - from); + } + + private static string[] DecileProfile(double[] values) + { + string[] profile = new string[10]; + for (int i = 0; i < 10; i++) + { + profile[i] = DecileAverage(values, i).ToString("F1"); + } + + return profile; + } + + private static double Percentile(double[] values, int percentile) + { + double[] sorted = (double[])values.Clone(); + Array.Sort(sorted); + int index = (int)Math.Round((percentile / 100.0) * (sorted.Length - 1)); + return sorted[Math.Clamp(index, 0, sorted.Length - 1)]; + } + } +} diff --git a/Tests/Xrpl.Tests/Client/BenchmarkWebSocketAssembly.cs b/Tests/Xrpl.Tests/Client/BenchmarkWebSocketAssembly.cs new file mode 100644 index 00000000..75f6e4c7 --- /dev/null +++ b/Tests/Xrpl.Tests/Client/BenchmarkWebSocketAssembly.cs @@ -0,0 +1,161 @@ +using Microsoft.VisualStudio.TestTools.UnitTesting; + +using System; +using System.Diagnostics; +using System.Threading; +using System.Threading.Tasks; + +using Xrpl.Client; + +namespace Xrpl.Tests.ClientLib +{ + /// + /// Manual benchmark of the WebSocket multi-chunk message assembly path. Deliberately named + /// outside the TestU/TestI filters so it never runs in CI; invoke it explicitly: + /// dotnet test --filter "FullyQualifiedName~BenchmarkWebSocketAssembly". + /// Knobs: BENCH_MESSAGES, BENCH_PAYLOAD_BYTES, BENCH_FRAGMENTS. + /// + [TestClass] + public class BenchmarkWebSocketAssembly + { + private static int EnvInt(string name, int fallback) + { + string? raw = Environment.GetEnvironmentVariable(name); + return int.TryParse(raw, out int value) && value > 0 ? value : fallback; + } + + [TestMethod] + public async Task BenchmarkMultiChunkAssembly() + { + int messageCount = EnvInt("BENCH_MESSAGES", 600); + int payloadBytes = EnvInt("BENCH_PAYLOAD_BYTES", 2 * 1024 * 1024); + int fragments = EnvInt("BENCH_FRAGMENTS", 32); + + using BulkMessageServer server = new BulkMessageServer(messageCount, payloadBytes, fragments); + + long[] timestamps = new long[messageCount]; + int received = 0; + int corrupted = 0; + TaskCompletionSource allReceived = + new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + + WebSocketClient client = WebSocketClient.Create(server.Url); + client.OnMessageReceived((message, _) => + { + int index = received; + if (message.Length != payloadBytes) + { + Interlocked.Increment(ref corrupted); + } + + if (index < timestamps.Length) + { + timestamps[index] = Stopwatch.GetTimestamp(); + } + + received = index + 1; + if (received >= messageCount) + { + allReceived.TrySetResult(true); + } + + return Task.CompletedTask; + }); + + // Warm up the JIT and the socket path before the measured window. + await client.Connect().ConfigureAwait(false); + + GC.Collect(2, GCCollectionMode.Forced, blocking: true); + GC.WaitForPendingFinalizers(); + GC.Collect(2, GCCollectionMode.Forced, blocking: true); + + long allocatedBefore = GC.GetTotalAllocatedBytes(precise: true); + int gen0Before = GC.CollectionCount(0); + int gen1Before = GC.CollectionCount(1); + int gen2Before = GC.CollectionCount(2); + long lohBefore = LohBytes(); + long startTicks = Stopwatch.GetTimestamp(); + + // Releases the server; nothing has been sent before this point. + client.SendMessage("go"); + + Task completed = await Task.WhenAny( + allReceived.Task, + Task.Delay(TimeSpan.FromMinutes(30))).ConfigureAwait(false); + + long stopTicks = Stopwatch.GetTimestamp(); + long allocatedAfter = GC.GetTotalAllocatedBytes(precise: true); + long lohAfter = LohBytes(); + + client.CancelIntentionally(); + client.Dispose(); + + Assert.AreSame(allReceived.Task, completed, "benchmark did not finish within 30 minutes"); + Assert.AreEqual(0, corrupted, "at least one message was assembled with the wrong length"); + + double totalSeconds = (stopTicks - startTicks) / (double)Stopwatch.Frequency; + long allocated = allocatedAfter - allocatedBefore; + + double firstDecile = DecileAverageMs(timestamps, startTicks, 0); + double lastDecile = DecileAverageMs(timestamps, startTicks, 9); + double slowest = SlowestMs(timestamps, startTicks); + + Console.WriteLine("=== WebSocket multi-chunk assembly benchmark ==="); + Console.WriteLine($"messages : {messageCount}"); + Console.WriteLine($"payload : {payloadBytes / 1024.0 / 1024.0:F2} MiB"); + Console.WriteLine($"fragments per msg : {fragments} ({payloadBytes / fragments / 1024} KiB each)"); + Console.WriteLine($"total time : {totalSeconds:F2} s ({messageCount / totalSeconds:F2} msg/s)"); + Console.WriteLine($"first decile avg : {firstDecile:F2} ms/msg"); + Console.WriteLine($"last decile avg : {lastDecile:F2} ms/msg"); + Console.WriteLine($"slowest message : {slowest:F2} ms"); + Console.WriteLine($"trend (last/first): {lastDecile / firstDecile:F2}x"); + Console.WriteLine($"allocated total : {allocated / 1024.0 / 1024.0:F1} MiB"); + Console.WriteLine($"allocated per msg : {allocated / (double)messageCount / 1024.0 / 1024.0:F2} MiB " + + $"({allocated / (double)messageCount / payloadBytes:F2}x payload)"); + Console.WriteLine($"gen0/gen1/gen2 : {GC.CollectionCount(0) - gen0Before} / " + + $"{GC.CollectionCount(1) - gen1Before} / {GC.CollectionCount(2) - gen2Before}"); + Console.WriteLine($"LOH size : {lohBefore / 1024.0 / 1024.0:F1} MiB -> {lohAfter / 1024.0 / 1024.0:F1} MiB"); + } + + private static long LohBytes() + { + GCMemoryInfo info = GC.GetGCMemoryInfo(); + ReadOnlySpan generations = info.GenerationInfo; + return generations.Length > 3 ? generations[3].SizeAfterBytes : 0; + } + + private static double DecileAverageMs(long[] timestamps, long startTicks, int decile) + { + int size = Math.Max(1, timestamps.Length / 10); + int from = decile * size; + int to = Math.Min(timestamps.Length, from + size); + if (from >= to) + { + return 0; + } + + long previous = from == 0 ? startTicks : timestamps[from - 1]; + double totalMs = (timestamps[to - 1] - previous) * 1000.0 / Stopwatch.Frequency; + return totalMs / (to - from); + } + + private static double SlowestMs(long[] timestamps, long startTicks) + { + double slowest = 0; + long previous = startTicks; + + foreach (long timestamp in timestamps) + { + double ms = (timestamp - previous) * 1000.0 / Stopwatch.Frequency; + if (ms > slowest) + { + slowest = ms; + } + + previous = timestamp; + } + + return slowest; + } + } +} diff --git a/Tests/Xrpl.Tests/Client/BulkMessageServer.cs b/Tests/Xrpl.Tests/Client/BulkMessageServer.cs new file mode 100644 index 00000000..62445f30 --- /dev/null +++ b/Tests/Xrpl.Tests/Client/BulkMessageServer.cs @@ -0,0 +1,261 @@ +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. + /// + internal sealed class BulkMessageServer : IDisposable + { + private readonly TcpListener _listener; + private readonly CancellationTokenSource _cts = new(); + private readonly int _messageCount; + private readonly int _fragments; + private readonly int[] _lengthCycle; + private readonly byte[] _payload; + private readonly TaskCompletionSource _finished = + new(TaskCreationOptions.RunContinuationsAsynchronously); + + /// How many messages to push once the client connects. + /// Size of the longest message payload, in bytes. + /// + /// Number of WebSocket frames each message is split into; the client sees exactly this + /// many receive chunks per message. + /// + /// + /// Optional cycle of message lengths, each a prefix of the full payload. Lets a test mix + /// long and short messages on one connection, which is what exposes stale bytes left in a + /// reused assembly buffer. Defaults to every message being the full payload. + /// + public BulkMessageServer(int messageCount, int payloadBytes, int fragments, int[]? lengthCycle = null) + { + _messageCount = messageCount; + _fragments = Math.Max(1, fragments); + _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; } + + public string Url => "ws://127.0.0.1:" + Port + "/"; + + /// Payload every message carries, as the client should see it. + public string PayloadText => Encoding.UTF8.GetString(_payload); + + public int PayloadBytes => _payload.Length; + + /// Completes with the number of messages written once the server is done sending. + public Task SendCompleted => _finished.Task; + + /// + /// Builds an ASCII payload shaped like a paged rippled response, so the byte content is + /// non-uniform and any accidental truncation during assembly is visible. + /// + private static byte[] BuildPayload(int payloadBytes) + { + StringBuilder builder = new StringBuilder(payloadBytes + 64); + builder.Append("{\"id\":1,\"status\":\"success\",\"type\":\"response\",\"result\":{\"state\":["); + + int index = 0; + while (builder.Length < payloadBytes - 80) + { + if (index > 0) + { + builder.Append(','); + } + + builder.Append("{\"i\":").Append(index).Append(",\"d\":\""); + builder.Append((char)('A' + (index % 26)), 48); + builder.Append("\"}"); + index++; + } + + builder.Append("]}}"); + + while (builder.Length < payloadBytes) + { + builder.Append(' '); + } + + return Encoding.UTF8.GetBytes(builder.ToString(0, payloadBytes)); + } + + private async Task AcceptAsync() + { + try + { + 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(_messageCount); + + await Task.Delay(Timeout.InfiniteTimeSpan, _cts.Token).ConfigureAwait(false); + } + catch (OperationCanceledException) + { + _finished.TrySetCanceled(); + } + catch (Exception ex) + { + _finished.TrySetException(ex); + } + } + + private async Task DrainAsync(NetworkStream stream) + { + byte[] sink = new byte[4096]; + + try + { + while (await stream.ReadAsync(sink, _cts.Token).ConfigureAwait(false) > 0) + { + } + } + catch + { + // 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 new file mode 100644 index 00000000..1d8b0f2c --- /dev/null +++ b/Tests/Xrpl.Tests/Client/PagedResponseServer.cs @@ -0,0 +1,355 @@ +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. + /// + internal sealed class PagedResponseServer : IDisposable + { + private readonly TcpListener _listener; + private readonly CancellationTokenSource _cts = new(); + private readonly int _fragments; + private readonly string _resultBody; + private int _served; + + /// Target size of each response, in bytes. + /// Number of WebSocket frames each response is split into. + 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(); + } + + 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); + + /// + /// Body of the result object: a binary-form ledger_data page, i.e. a list of + /// {data, index} pairs, which is what a bulk ledger crawl actually receives. + /// + private static string BuildResultBody(int approximatePayloadBytes) + { + StringBuilder builder = new StringBuilder(approximatePayloadBytes + 256); + builder.Append("{\"ledger_hash\":\"842B57C1CC0613299A686D3E9F310EC0422C84D3911E5056389AA7E5808A93C8\","); + builder.Append("\"ledger_index\":\"96000000\",\"validated\":true,\"state\":["); + + int index = 0; + while (builder.Length < approximatePayloadBytes) + { + if (index > 0) + { + builder.Append(','); + } + + builder.Append("{\"data\":\""); + AppendHex(builder, index, 900); + builder.Append("\",\"index\":\""); + AppendHex(builder, index + 7, 64); + builder.Append("\"}"); + index++; + } + + builder.Append("],\"marker\":\""); + AppendHex(builder, index, 64); + builder.Append("\"}"); + + return builder.ToString(); + } + + private static void AppendHex(StringBuilder builder, int seed, int length) + { + const string Digits = "0123456789ABCDEF"; + int state = (seed * 1103515245) ^ 0x5F3A; + + for (int i = 0; i < length; i++) + { + state = (state * 1103515245) + 12345; + builder.Append(Digits[(state >> 16) & 0xF]); + } + } + + private async Task AcceptAsync() + { + try + { + 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) + { + return; + } + + string id = ExtractId(message); + await WriteResponseAsync(stream, id).ConfigureAwait(false); + Interlocked.Increment(ref _served); + } + } + 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. + } + } + + /// Pulls the JSON string value of the request's "id" property. + private static string ExtractId(string message) + { + int keyIndex = message.IndexOf("\"id\"", StringComparison.Ordinal); + if (keyIndex < 0) + { + return "\"0\""; + } + + int colon = message.IndexOf(':', keyIndex); + if (colon < 0) + { + return "\"0\""; + } + + int start = colon + 1; + while (start < message.Length && char.IsWhiteSpace(message[start])) + { + start++; + } + + if (start < message.Length && message[start] == '"') + { + int end = message.IndexOf('"', start + 1); + return end < 0 ? "\"0\"" : message.Substring(start, end - start + 1); + } + + int stop = start; + while (stop < message.Length && message[stop] != ',' && message[stop] != '}') + { + stop++; + } + + 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++) + { + int offset = fragment * fragmentBytes; + 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. + /// + private async Task ReadTextFrameAsync(NetworkStream stream) + { + while (true) + { + byte[] head = new byte[2]; + if (!await ReadExactAsync(stream, head, 2).ConfigureAwait(false)) + { + return null; + } + + int opcode = head[0] & 0x0F; + bool masked = (head[1] & 0x80) != 0; + long length = head[1] & 0x7F; + + if (length == 126) + { + byte[] extended = new byte[2]; + if (!await ReadExactAsync(stream, extended, 2).ConfigureAwait(false)) + { + return null; + } + + length = BinaryPrimitives.ReadUInt16BigEndian(extended); + } + else if (length == 127) + { + byte[] extended = new byte[8]; + if (!await ReadExactAsync(stream, extended, 8).ConfigureAwait(false)) + { + return null; + } + + length = (long)BinaryPrimitives.ReadUInt64BigEndian(extended); + } + + byte[] mask = new byte[4]; + if (masked && !await ReadExactAsync(stream, mask, 4).ConfigureAwait(false)) + { + return null; + } + + byte[] payload = new byte[length]; + if (length > 0 && !await ReadExactAsync(stream, payload, (int)length).ConfigureAwait(false)) + { + return null; + } + + if (masked) + { + for (int i = 0; i < payload.Length; i++) + { + payload[i] ^= mask[i % 4]; + } + } + + if (opcode == 0x8) + { + return null; + } + + if (opcode == 0x1 || opcode == 0x2) + { + return Encoding.UTF8.GetString(payload); + } + } + } + + 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/TestUWebSocketMessageAssembly.cs b/Tests/Xrpl.Tests/Client/TestUWebSocketMessageAssembly.cs new file mode 100644 index 00000000..ecea00c2 --- /dev/null +++ b/Tests/Xrpl.Tests/Client/TestUWebSocketMessageAssembly.cs @@ -0,0 +1,135 @@ +using Microsoft.VisualStudio.TestTools.UnitTesting; + +using System; +using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; + +using Xrpl.Client; + +namespace Xrpl.Tests.ClientLib +{ + /// + /// Regression coverage for assembling a WebSocket message that arrives as several receive + /// chunks. The assembly buffer is reused for the life of the connection, so the tests check + /// both that every message comes out byte-exact and that the cost per message does not grow + /// with the number of chunks it was split into. + /// + [TestClass] + public class TestUWebSocketMessageAssembly + { + private const int WaitSeconds = 120; + + [TestMethod] + public async Task TestUMultiChunkMessageArrivesIntact() + { + // Deliberately larger than the client's 1 MiB receive buffer and split far more finely + // than the buffer would split it on its own. + const int PayloadBytes = 3 * 1024 * 1024; + + using BulkMessageServer server = new BulkMessageServer(4, PayloadBytes, fragments: 96); + IReadOnlyList messages = await ReceiveAsync(server, 4).ConfigureAwait(false); + + Assert.AreEqual(4, messages.Count); + foreach (string message in messages) + { + Assert.AreEqual(server.PayloadText, message); + } + } + + [TestMethod] + public async Task TestUShortMessageAfterLongOneIsNotPaddedWithStaleBytes() + { + const int PayloadBytes = 2 * 1024 * 1024; + int[] lengthCycle = { PayloadBytes, PayloadBytes / 8, PayloadBytes / 2, 1024 }; + + using BulkMessageServer server = new BulkMessageServer( + messageCount: 12, + payloadBytes: PayloadBytes, + fragments: 16, + lengthCycle: lengthCycle); + + IReadOnlyList messages = await ReceiveAsync(server, 12).ConfigureAwait(false); + + Assert.AreEqual(12, messages.Count); + for (int i = 0; i < messages.Count; i++) + { + int expectedLength = lengthCycle[i % lengthCycle.Length]; + Assert.AreEqual(expectedLength, messages[i].Length, $"message {i} has the wrong length"); + Assert.AreEqual(server.PayloadText.Substring(0, expectedLength), messages[i], + $"message {i} does not match the expected prefix"); + } + } + + /// + /// Guards the shape of the fix: assembly used to be quadratic in the number of chunks, so + /// allocation per message grew with the split. The bound is deliberately loose — the point + /// is that a 64-way split must not cost an order of magnitude more than the payload. + /// + [TestMethod] + public async Task TestUAssemblyAllocationDoesNotGrowWithChunkCount() + { + const int MessageCount = 200; + const int PayloadBytes = 1024 * 1024; + + // Floor per message is the exact-sized byte[] plus the UTF-16 string handed to the + // callback, i.e. about 3x the payload. Quadratic assembly at 64 chunks cost ~34x. + const double AllowedTimesPayload = 12.0; + + using BulkMessageServer server = new BulkMessageServer(MessageCount, PayloadBytes, fragments: 64); + + GC.Collect(2, GCCollectionMode.Forced, blocking: true); + GC.WaitForPendingFinalizers(); + GC.Collect(2, GCCollectionMode.Forced, blocking: true); + + long allocatedBefore = GC.GetTotalAllocatedBytes(precise: true); + IReadOnlyList messages = await ReceiveAsync(server, MessageCount).ConfigureAwait(false); + long allocated = GC.GetTotalAllocatedBytes(precise: true) - allocatedBefore; + + Assert.AreEqual(MessageCount, messages.Count); + + double perMessage = allocated / (double)MessageCount / PayloadBytes; + Assert.IsTrue( + perMessage < AllowedTimesPayload, + $"allocated {perMessage:F1}x payload per message, expected below {AllowedTimesPayload:F1}x"); + } + + private static async Task> ReceiveAsync(BulkMessageServer server, int expected) + { + List messages = new List(expected); + TaskCompletionSource done = + new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + + WebSocketClient client = WebSocketClient.Create(server.Url); + client.OnMessageReceived((message, _) => + { + messages.Add(message); + if (messages.Count >= expected) + { + done.TrySetResult(true); + } + + return Task.CompletedTask; + }); + + try + { + await client.Connect().ConfigureAwait(false); + + // The server holds off until the client speaks, so nothing is missed. + client.SendMessage("go"); + + 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"); + } + finally + { + client.CancelIntentionally(); + client.Dispose(); + } + + return messages; + } + } +} diff --git a/Xrpl/Client/WebSocketClient.cs b/Xrpl/Client/WebSocketClient.cs index dd8636e8..90bc6faf 100644 --- a/Xrpl/Client/WebSocketClient.cs +++ b/Xrpl/Client/WebSocketClient.cs @@ -1,10 +1,10 @@  using System; +using System.Buffers; using System.Collections.Generic; using System.Diagnostics; using System.IO; -using System.Linq; using System.Net.WebSockets; using System.Text; using System.Threading; @@ -476,43 +476,64 @@ await Task.WhenAny( private async Task ReceiveLoopAsync() { - byte[] buffer = new byte[ReceiveChunkSize]; + // One receive buffer per connection, rented rather than allocated: at ReceiveChunkSize + // it is a large-object-heap array, and a reconnect loop would otherwise leak one per session. + byte[] buffer = ArrayPool.Shared.Rent(ReceiveChunkSize); + + // 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. + byte[]? assemblyBuffer = null; try { // Continue receiving while Open OR CloseSent (waiting for server's Close frame after CloseOutputAsync) - while (_ws != null && - (_ws.State == WebSocketState.Open || _ws.State == WebSocketState.CloseSent) && + while (_ws != null && + (_ws.State == WebSocketState.Open || _ws.State == WebSocketState.CloseSent) && !_cancellationToken.IsCancellationRequested) { - byte[] byteResult = Array.Empty(); + byte[]? completeMessage = null; + int assembledLength = 0; WebSocketReceiveResult result; - bool timedOut = false; do { - result = await _ws.ReceiveAsync(new ArraySegment(buffer), _cancellationToken).ConfigureAwait(false); + result = await _ws.ReceiveAsync(new ArraySegment(buffer, 0, ReceiveChunkSize), _cancellationToken).ConfigureAwait(false); if (result.MessageType == WebSocketMessageType.Close) { WebSocketCloseStatus? closeStatus = result.CloseStatus; string? closeDescription = result.CloseStatusDescription; - + _onClosed?.Invoke(this); await CallOnDisconnectedAsync(closeStatus, closeDescription).ConfigureAwait(false); return; } - else + else if (result.EndOfMessage && assembledLength == 0) + { + // Whole message arrived in a single chunk - copy it out directly. + completeMessage = new byte[result.Count]; + Buffer.BlockCopy(buffer, 0, completeMessage, 0, result.Count); + } + else if (result.Count > 0) { - byteResult = byteResult.Concat(buffer.Take(result.Count)).ToArray(); + EnsureAssemblyCapacity(ref assemblyBuffer, assembledLength + result.Count); + Buffer.BlockCopy(buffer, 0, assemblyBuffer, assembledLength, result.Count); + assembledLength += result.Count; } } while (!result.EndOfMessage); - if (timedOut) - continue; + if (completeMessage == null) + { + completeMessage = new byte[assembledLength]; + if (assembledLength > 0) + { + Buffer.BlockCopy(assemblyBuffer!, 0, completeMessage, 0, assembledLength); + } + } - CallOnMessage(byteResult); + CallOnMessage(completeMessage); } } catch (OperationCanceledException) when (_isIntentionalDisconnect || _cancellationToken.IsCancellationRequested || IsDisposed) @@ -611,6 +632,36 @@ private async Task ReceiveLoopAsync() _onConnectionError?.Invoke(ex, this); await CallOnDisconnectedAsync(WebSocketCloseStatus.EndpointUnavailable, "Unknown error: " + ex.Message).ConfigureAwait(false); } + finally + { + ArrayPool.Shared.Return(buffer); + } + } + + /// + /// Grows so it can hold + /// bytes, doubling the current capacity and preserving what has already been assembled. + /// + private static void EnsureAssemblyCapacity(ref byte[]? assemblyBuffer, int requiredLength) + { + if (assemblyBuffer != null && assemblyBuffer.Length >= requiredLength) + { + return; + } + + int capacity = assemblyBuffer?.Length ?? ReceiveChunkSize; + while (capacity < requiredLength) + { + capacity = capacity <= Array.MaxLength / 2 ? capacity * 2 : requiredLength; + } + + byte[] grown = new byte[capacity]; + if (assemblyBuffer != null) + { + Buffer.BlockCopy(assemblyBuffer, 0, grown, 0, assemblyBuffer.Length); + } + + assemblyBuffer = grown; } private void CallOnMessage(byte[] result) From ff4f14651d414cbf7effef17fe8bf3d157f6f5dd Mon Sep 17 00:00:00 2001 From: Aleksandr Platonenkov Date: Thu, 13 Aug 2026 15:11:46 -0300 Subject: [PATCH 2/8] =?UTF-8?q?fix(client):=20=D0=BE=D1=81=D0=B2=D0=BE?= =?UTF-8?q?=D0=B1=D0=BE=D0=B6=D0=B4=D0=B0=D1=82=D1=8C=20=D1=82=D0=B0=D0=B9?= =?UTF-8?q?=D0=BC=D0=B5=D1=80=20=D1=82=D0=B0=D0=B9=D0=BC=D0=B0=D1=83=D1=82?= =?UTF-8?q?=D0=B0=20=D0=B7=D0=B0=D0=BF=D1=80=D0=BE=D1=81=D0=B0,=20=D0=B0?= =?UTF-8?q?=20=D0=BD=D0=B5=20=D1=82=D0=BE=D0=BB=D1=8C=D0=BA=D0=BE=20=D0=BE?= =?UTF-8?q?=D1=81=D1=82=D0=B0=D0=BD=D0=B0=D0=B2=D0=BB=D0=B8=D0=B2=D0=B0?= =?UTF-8?q?=D1=82=D1=8C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Resolve и Reject вызывали timer.Stop(), но не Dispose(). System.Timers.Timer наследует Component с финализатором, поэтому каждый завершённый запрос оставлял финализируемый объект; за длинный постраничный обход это тысячи недоосвобождённых таймеров, каждый из которых через замыкание Elapsed ещё и удерживает сериализованный текст своего запроса. Dispose останавливает таймер и снимает его с очереди финализации. --- Xrpl/Client/RequestManager.cs | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) diff --git a/Xrpl/Client/RequestManager.cs b/Xrpl/Client/RequestManager.cs index be85d284..cec67976 100644 --- a/Xrpl/Client/RequestManager.cs +++ b/Xrpl/Client/RequestManager.cs @@ -68,8 +68,10 @@ public void Resolve(Guid id, BaseResponse response) return; } - if (timeoutsAwaitingResponse.TryRemove(id, out var timer)) - timer.Stop(); + // Dispose, not Stop: Timer is a finalizable Component, and a stopped but undisposed + // one per request piles up on the finalization queue over a long paged run. + if (timeoutsAwaitingResponse.TryRemove(id, out Timer timer)) + timer.Dispose(); try { @@ -99,8 +101,10 @@ public void Reject(Guid id, T error) where T : Exception Debug.WriteLine($"Reject called for non-existent promise {id} (likely already resolved)"); return; } - if (timeoutsAwaitingResponse.TryRemove(id, out var timer)) - timer.Stop(); + // Dispose, not Stop: Timer is a finalizable Component, and a stopped but undisposed + // one per request piles up on the finalization queue over a long paged run. + if (timeoutsAwaitingResponse.TryRemove(id, out Timer timer)) + timer.Dispose(); var setException = taskInfo.TaskCompletionResult.GetType().GetMethod("TrySetException", new Type[] { typeof(Exception) }, null); setException.Invoke(taskInfo.TaskCompletionResult, new[] { error }); From e8dca3728225458ac2cecf513ec3f8a8e6f98b63 Mon Sep 17 00:00:00 2001 From: Aleksandr Platonenkov Date: Thu, 13 Aug 2026 15:12:06 -0300 Subject: [PATCH 3/8] =?UTF-8?q?perf(client):=20=D1=83=D0=B1=D1=80=D0=B0?= =?UTF-8?q?=D1=82=D1=8C=20=D1=80=D0=B5=D1=84=D0=BB=D0=B5=D0=BA=D1=81=D0=B8?= =?UTF-8?q?=D1=8E=20=D1=81=20=D0=BF=D1=83=D1=82=D0=B8=20=D1=80=D0=B0=D0=B7?= =?UTF-8?q?=D0=B1=D0=BE=D1=80=D0=B0=20=D0=BA=D0=B0=D0=B6=D0=B4=D0=BE=D0=B3?= =?UTF-8?q?=D0=BE=20=D0=BE=D1=82=D0=B2=D0=B5=D1=82=D0=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Resolve, Reject и ObserveTaskException доставали TrySetResult / TrySetException / Task через GetType().GetMethod(...) + Invoke на каждый ответ. TaskInfo теперь несёт типизированные делегаты SetResult и SetException и саму CompletionTask, проставляемые при создании запроса. Свойства добавлены, а не заменены: TaskInfo — публичный тип, поэтому для экземпляров, собранных вне RequestManager, оставлен прежний путь через рефлексию. --- Xrpl/Client/RequestManager.cs | 69 ++++++++++++++++++++++++++++------- Xrpl/Client/TaskInfo.cs | 20 ++++++++++ 2 files changed, 76 insertions(+), 13 deletions(-) diff --git a/Xrpl/Client/RequestManager.cs b/Xrpl/Client/RequestManager.cs index cec67976..6dc9806d 100644 --- a/Xrpl/Client/RequestManager.cs +++ b/Xrpl/Client/RequestManager.cs @@ -4,6 +4,7 @@ using System.Collections.Concurrent; using System.Collections.Generic; using System.Diagnostics; +using System.Reflection; using System.Text.Json; using System.Text.Json.Nodes; using System.Threading; @@ -75,9 +76,8 @@ public void Resolve(Guid id, BaseResponse response) try { - var deserialized = JsonSerializer.Deserialize(response.Result?.ToString() ?? "{}", taskInfo.Type, serializerOptions); - var setResult = taskInfo.TaskCompletionResult.GetType().GetMethod("TrySetResult"); - setResult.Invoke(taskInfo.TaskCompletionResult, new[] { deserialized }); + object deserialized = JsonSerializer.Deserialize(response.Result?.ToString() ?? "{}", taskInfo.Type, serializerOptions); + CompleteWithResult(taskInfo, deserialized); this.DeletePromise(id, taskInfo); } catch (Exception ex) @@ -105,30 +105,67 @@ public void Reject(Guid id, T error) where T : Exception // one per request piles up on the finalization queue over a long paged run. if (timeoutsAwaitingResponse.TryRemove(id, out Timer timer)) timer.Dispose(); - var setException = taskInfo.TaskCompletionResult.GetType().GetMethod("TrySetException", new Type[] { typeof(Exception) }, null); - setException.Invoke(taskInfo.TaskCompletionResult, new[] { error }); - + CompleteWithException(taskInfo, error); + // Observe the exception to prevent UnobservedTaskException in consuming apps // This is critical for MAUI/mobile apps that have global exception handlers - ObserveTaskException(taskInfo.TaskCompletionResult); + ObserveTaskException(taskInfo); this.DeletePromise(id, taskInfo); } + /// + /// Completes the pending request with a deserialized result, using the typed delegate + /// captured when the request was created and falling back to reflection for + /// instances built outside this manager. + /// + private static void CompleteWithResult(TaskInfo taskInfo, object result) + { + if (taskInfo.SetResult is not null) + { + taskInfo.SetResult(result); + return; + } + + MethodInfo setResult = taskInfo.TaskCompletionResult.GetType().GetMethod("TrySetResult"); + setResult.Invoke(taskInfo.TaskCompletionResult, new[] { result }); + } + + /// + /// Faults the pending request. See for the fallback rules. + /// + private static void CompleteWithException(TaskInfo taskInfo, Exception error) + { + if (taskInfo.SetException is not null) + { + taskInfo.SetException(error); + return; + } + + MethodInfo setException = taskInfo.TaskCompletionResult.GetType() + .GetMethod("TrySetException", new Type[] { typeof(Exception) }, null); + setException.Invoke(taskInfo.TaskCompletionResult, new object[] { error }); + } + /// /// Observes the exception on a TaskCompletionSource's Task to prevent UnobservedTaskException. /// When a Task faults but is never awaited, .NET raises UnobservedTaskException event. /// By adding a ContinueWith that reads the exception, we mark it as "observed". /// - private void ObserveTaskException(object taskCompletionSource) + private void ObserveTaskException(TaskInfo taskInfo) { try { - // Get the Task property from TaskCompletionSource - var taskProperty = taskCompletionSource.GetType().GetProperty("Task"); - if (taskProperty == null) return; - - var task = taskProperty.GetValue(taskCompletionSource) as Task; + Task task = taskInfo.CompletionTask; + if (task == null) + { + // Externally built TaskInfo: fall back to reading the Task property reflectively. + PropertyInfo taskProperty = taskInfo.TaskCompletionResult.GetType().GetProperty("Task"); + if (taskProperty == null) return; + + task = taskProperty.GetValue(taskInfo.TaskCompletionResult) as Task; + } + if (task == null) return; // Add a continuation that observes the exception (reads it to mark as handled) @@ -218,6 +255,9 @@ public XrplGRequest CreateGRequest( TaskInfo taskInfo = new TaskInfo(); taskInfo.TaskId = newId; taskInfo.TaskCompletionResult = task; + taskInfo.SetResult = result => task.TrySetResult(result); + taskInfo.SetException = error => task.TrySetException(error); + taskInfo.CompletionTask = task.Task; taskInfo.RemoveUponCompletion = true; taskInfo.Type = typeof(T); @@ -295,6 +335,9 @@ public XrplRequest CreateRequest( TaskInfo taskInfo = new TaskInfo(); taskInfo.TaskId = newId; taskInfo.TaskCompletionResult = task; + taskInfo.SetResult = result => task.TrySetResult((Dictionary)result); + taskInfo.SetException = error => task.TrySetException(error); + taskInfo.CompletionTask = task.Task; taskInfo.RemoveUponCompletion = true; taskInfo.Type = typeof(Dictionary); diff --git a/Xrpl/Client/TaskInfo.cs b/Xrpl/Client/TaskInfo.cs index ba478f14..6b4db338 100644 --- a/Xrpl/Client/TaskInfo.cs +++ b/Xrpl/Client/TaskInfo.cs @@ -1,5 +1,6 @@ using System; using System.Threading; +using System.Threading.Tasks; namespace Xrpl.Client { @@ -11,6 +12,25 @@ public class TaskInfo public object TaskCompletionResult { get; set; } + /// + /// Completes with a deserialized result. Set when the + /// request is created, so the response path does not have to reach for the strongly typed + /// TrySetResult through reflection. Null for externally built instances. + /// + public Func SetResult { get; set; } + + /// + /// Faults . Set when the request is created; null for + /// externally built instances. + /// + public Func SetException { get; set; } + + /// + /// The task behind , used to observe faults without + /// reflection. Null for externally built instances. + /// + public Task CompletionTask { get; set; } + public bool RemoveUponCompletion { get; set; } public CancellationTokenRegistration? CancellationRegistration { get; set; } From 89a3673ccdbf2c0a4c1a0f194f13635de7ff3ebb Mon Sep 17 00:00:00 2001 From: Aleksandr Platonenkov Date: Thu, 13 Aug 2026 15:15:02 -0300 Subject: [PATCH 4/8] =?UTF-8?q?refactor(client):=20=D1=82=D0=BE=D1=87?= =?UTF-8?q?=D0=BD=D0=B5=D0=B5=20=D0=BA=D0=BE=D0=BC=D0=BC=D0=B5=D0=BD=D1=82?= =?UTF-8?q?=D0=B0=D1=80=D0=B8=D0=B9=20=D0=BF=D1=80=D0=BE=20=D0=BF=D1=83?= =?UTF-8?q?=D0=BB=20=D0=B8=20=D0=BA=D0=BE=D0=BF=D0=B8=D1=80=D0=BE=D0=B2?= =?UTF-8?q?=D0=B0=D0=BD=D0=B8=D0=B5=20=D1=82=D0=BE=D0=BB=D1=8C=D0=BA=D0=BE?= =?UTF-8?q?=20=D1=81=D0=BE=D0=B1=D1=80=D0=B0=D0=BD=D0=BD=D0=BE=D0=B3=D0=BE?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit По итогам селф-ревью: комментарий у аренды буфера говорил про «утечку», хотя речь про свежий мегабайт в LOH на каждое переподключение. При росте scratch-буфера копируется ровно собранная часть, а не весь старый массив вместе с мусором за концом. --- Xrpl/Client/WebSocketClient.cs | 15 +++++++++------ 1 file changed, 9 insertions(+), 6 deletions(-) diff --git a/Xrpl/Client/WebSocketClient.cs b/Xrpl/Client/WebSocketClient.cs index 90bc6faf..8a0c885d 100644 --- a/Xrpl/Client/WebSocketClient.cs +++ b/Xrpl/Client/WebSocketClient.cs @@ -477,7 +477,9 @@ await Task.WhenAny( private async Task ReceiveLoopAsync() { // One receive buffer per connection, rented rather than allocated: at ReceiveChunkSize - // it is a large-object-heap array, and a reconnect loop would otherwise leak one per session. + // it lands on the large object heap, so a reconnect loop would otherwise burn a fresh + // uncompactable megabyte per session. Returned in the finally below, which only runs + // once the loop has stopped awaiting the socket. byte[] buffer = ArrayPool.Shared.Rent(ReceiveChunkSize); // Scratch buffer for messages that arrive in more than one chunk. It grows to the @@ -517,7 +519,7 @@ private async Task ReceiveLoopAsync() } else if (result.Count > 0) { - EnsureAssemblyCapacity(ref assemblyBuffer, assembledLength + result.Count); + EnsureAssemblyCapacity(ref assemblyBuffer, assembledLength + result.Count, assembledLength); Buffer.BlockCopy(buffer, 0, assemblyBuffer, assembledLength, result.Count); assembledLength += result.Count; } @@ -640,9 +642,10 @@ private async Task ReceiveLoopAsync() /// /// Grows so it can hold - /// bytes, doubling the current capacity and preserving what has already been assembled. + /// bytes, doubling the current capacity and carrying over the first + /// bytes already assembled. /// - private static void EnsureAssemblyCapacity(ref byte[]? assemblyBuffer, int requiredLength) + private static void EnsureAssemblyCapacity(ref byte[]? assemblyBuffer, int requiredLength, int preserveLength) { if (assemblyBuffer != null && assemblyBuffer.Length >= requiredLength) { @@ -656,9 +659,9 @@ private static void EnsureAssemblyCapacity(ref byte[]? assemblyBuffer, int requi } byte[] grown = new byte[capacity]; - if (assemblyBuffer != null) + if (preserveLength > 0) { - Buffer.BlockCopy(assemblyBuffer, 0, grown, 0, assemblyBuffer.Length); + Buffer.BlockCopy(assemblyBuffer!, 0, grown, 0, preserveLength); } assemblyBuffer = grown; From 11a2976115753030ca4fe608d1375ec321d0d478 Mon Sep 17 00:00:00 2001 From: Aleksandr Platonenkov Date: Thu, 13 Aug 2026 15:15:38 -0300 Subject: [PATCH 5/8] =?UTF-8?q?docs(changes):=20=D0=B7=D0=B0=D0=BF=D0=B8?= =?UTF-8?q?=D1=81=D1=8C=20=D0=BE=20=D0=BB=D0=B8=D0=BD=D0=B5=D0=B9=D0=BD?= =?UTF-8?q?=D0=BE=D0=B9=20=D1=81=D0=B1=D0=BE=D1=80=D0=BA=D0=B5=20=D1=81?= =?UTF-8?q?=D0=BE=D0=BE=D0=B1=D1=89=D0=B5=D0=BD=D0=B8=D0=B9=20=D0=B8=20?= =?UTF-8?q?=D0=BF=D1=80=D0=B0=D0=B2=D0=BA=D0=B0=D1=85=20RequestManager?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- CHANGES.md | 3 +++ 1 file changed, 3 insertions(+) diff --git a/CHANGES.md b/CHANGES.md index 43133aa8..cf2b4bae 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -24,6 +24,9 @@ * `.ci-config/bump-nightly-pin.sh` does the move: newest `xrpld` build from the nightly apt channel, `ARG XRPLD_VERSION` rewritten, `rippled.batchv11.cfg` regenerated from the develop commit **encoded in that version string** — config and binary cannot drift apart, which is the failure mode the old manual two-step invited. `--check` reports the pin, the newest build and the pin's age without touching anything. Both timestamp formats are compared by their common `YYYYMMDDHHMM` prefix * the workflow bumps only once the pin is older than `MAX_PIN_AGE_DAYS` (21) — nightly publishes several builds a day, and a weekly PR would be noise rather than signal; `workflow_dispatch` takes a `force` input for the exceptions. It then builds and starts the stand on the new pin and requires the AMM sentinel amendment to come up enabled at genesis, which is what proves the regenerated config was accepted rather than silently ignored, and attaches the definitions diff against the new build to the PR body — a `node-only` field there is the SDK being behind develop, reported instead of hidden * credentials, idempotency and the tracking-issue fallback follow release-watch exactly, including the one-notification-per-failure-streak rule +* **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 +* **Request timeout timers were stopped but never disposed** — `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 +* **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 * **Cancellation no longer disappears into the autofill fee fallbacks** — `FetchCounterpartySignerCount` and `FetchLoan` wrap their client call in a broad `catch`, which is right for the case they exist for (the counterparty account or the Loan object is not there yet, and preclaim will report it) but also swallowed an `OperationCanceledException` raised from the caller's own token. Autofill then carried on and wrote a fee derived from the fallback — one signer, no loan — for a request the caller had already abandoned. Both catches now carry `when (!cancellationToken.IsCancellationRequested)`, which lets a caller's cancellation through while a client-side timeout, which does not cancel that token, still falls back as before. Covered in both directions: a cancelled token must throw and leave no `Fee` behind, an unreadable Loan object must still fall back * **`MPTokenIssuanceSet` validation reports a malformed `Flags` as `ValidationException`** — it went through `Convert.ToUInt32`, which throws `FormatException` or `InvalidCastException` on a non-numeric value, while the `ImmutableFlags` check two lines below reports `ValidationException` like the rest of the validators. Callers catching `ValidationException` did not catch the other two * **The conformance fixtures are re-pinned to the 3.3.0 tag** — `transactions.macro` and `LedgerFormats.h` now come from the release commit (`00a178fb`) instead of a July `develop` sha and `3.3.0-rc1`; `ledger_entries.macro` stays on `develop` (`9859e5ce`) for the reason its `.ref` already gives — `sfLEVersion` exists only there. Both macro files are byte-identical to upstream and re-verifiable with the `curl … | diff` line in each `.ref`. This is what makes the guards test against the version CI actually runs: From dd2a3f6cf4175ebc87349ba702efd61fe0995a15 Mon Sep 17 00:00:00 2001 From: Aleksandr Platonenkov Date: Thu, 13 Aug 2026 15:21:51 -0300 Subject: [PATCH 6/8] =?UTF-8?q?refactor(client):=20=D1=83=D0=B1=D1=80?= =?UTF-8?q?=D0=B0=D1=82=D1=8C=20=D0=BC=D1=91=D1=80=D1=82=D0=B2=D0=BE=D0=B5?= =?UTF-8?q?=20=D0=BF=D0=BE=D0=BB=D0=B5=20tasks=20=D0=B8=D0=B7=20XrplClient?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit private readonly ConcurrentDictionary tasks нигде не присваивалось и нигде не читалось — readonly-поле без инициализации, всегда null. Остаток от реализации, где клиент сам вёл учёт запросов; сейчас этим занимается RequestManager. XrplClient не partial, обращений по имени через рефлексию нет. System.Collections.Concurrent держался только этой строкой и убран вместе с ней. --- Xrpl/Client/IXrplClient.cs | 3 --- 1 file changed, 3 deletions(-) diff --git a/Xrpl/Client/IXrplClient.cs b/Xrpl/Client/IXrplClient.cs index 0a87af1b..90633b2f 100644 --- a/Xrpl/Client/IXrplClient.cs +++ b/Xrpl/Client/IXrplClient.cs @@ -1,5 +1,4 @@ using System; -using System.Collections.Concurrent; using System.Collections.Generic; using System.Text.Json; using System.Threading; @@ -487,8 +486,6 @@ public class ClientOptions : ConnectionOptions ///// Current web socket client state //public WebSocketState SocketState => client.State; - private readonly ConcurrentDictionary tasks; - public XrplClient(string server, ClientOptions? options = null) { From 3d3d704f815c3a6fae704a7eee6b9385f3c78217 Mon Sep 17 00:00:00 2001 From: Aleksandr Platonenkov Date: Thu, 13 Aug 2026 15:28:29 -0300 Subject: [PATCH 7/8] =?UTF-8?q?fix(client):=20=D0=BD=D0=B5=20=D0=BE=D1=81?= =?UTF-8?q?=D1=82=D0=B0=D0=B2=D0=BB=D1=8F=D1=82=D1=8C=20=D1=82=D0=B0=D0=B9?= =?UTF-8?q?=D0=BC=D0=B5=D1=80=20=D0=B8=20=D1=80=D0=B5=D0=B3=D0=B8=D1=81?= =?UTF-8?q?=D1=82=D1=80=D0=B0=D1=86=D0=B8=D1=8E=20=D0=BE=D1=82=D0=BC=D0=B5?= =?UTF-8?q?=D0=BD=D1=8B=20=D0=BE=D1=82=20=D1=83=D0=B6=D0=B5=20=D0=BE=D1=82?= =?UTF-8?q?=D0=BC=D0=B5=D0=BD=D1=91=D0=BD=D0=BD=D0=BE=D0=B3=D0=BE=20=D1=82?= =?UTF-8?q?=D0=BE=D0=BA=D0=B5=D0=BD=D0=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Замечания CodeRabbit к PR #89. - RequestManager: уже отменённый токен выполняет колбэк Register синхронно, поэтому Reject завершал запрос раньше, чем фабрика успевала зарегистрировать таймер таймаута — и снимать было нечего. Таймер добавлялся уже после удаления промиса и оставался в timeoutsAwaitingResponse навсегда: когда он срабатывал, Reject уходил в ранний return, не дойдя до снятия. По той же причине оставалась неосвобождённой CancellationTokenRegistration — присваивание в TaskInfo происходило после того, как DeletePromise уже отработал. Обе фабрики теперь проверяют, жив ли ещё промис, и подчищают за собой; снятие таймера вынесено в DisposeTimeout и вызывается в том числе на ранних возвратах Resolve и Reject, что закрывает и узкую гонку с параллельной отменой - PagedResponseServer: offset не ограничивался длиной payload, хотя в BulkMessageServer ограничение уже стояло. При маленьком payload и большом числе фрагментов деление с округлением вверх уводило offset за конец, length уходил в минус и AsMemory бросал ArgumentOutOfRangeException мимо catch-ей AcceptAsync - TestUWebSocketMessageAssembly помечен [DoNotParallelize]: GC.GetTotalAllocatedBytes считает аллокации всего процесса, а прогон параллелит на уровне классов Тесты: TestURequestManagerCancellation — уже отменённый токен не оставляет ни таймера, ни промиса (обе фабрики), живой запрос по-прежнему взводит таймаут и снимает его при завершении. Первые два падают на коде до правки. --- .../Xrpl.Tests/Client/PagedResponseServer.cs | 4 +- .../Client/TestURequestManagerCancellation.cs | 102 ++++++++++++++++++ .../Client/TestUWebSocketMessageAssembly.cs | 3 + Xrpl/Client/RequestManager.cs | 63 +++++++++-- 4 files changed, 163 insertions(+), 9 deletions(-) create mode 100644 Tests/Xrpl.Tests/Client/TestURequestManagerCancellation.cs diff --git a/Tests/Xrpl.Tests/Client/PagedResponseServer.cs b/Tests/Xrpl.Tests/Client/PagedResponseServer.cs index 1d8b0f2c..8162114c 100644 --- a/Tests/Xrpl.Tests/Client/PagedResponseServer.cs +++ b/Tests/Xrpl.Tests/Client/PagedResponseServer.cs @@ -173,7 +173,9 @@ private async Task WriteResponseAsync(NetworkStream stream, string id) for (int fragment = 0; fragment < _fragments; fragment++) { - int offset = fragment * fragmentBytes; + // 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; diff --git a/Tests/Xrpl.Tests/Client/TestURequestManagerCancellation.cs b/Tests/Xrpl.Tests/Client/TestURequestManagerCancellation.cs new file mode 100644 index 00000000..fb4fc8c7 --- /dev/null +++ b/Tests/Xrpl.Tests/Client/TestURequestManagerCancellation.cs @@ -0,0 +1,102 @@ +using Microsoft.VisualStudio.TestTools.UnitTesting; + +using System; +using System.Collections.Concurrent; +using System.Collections.Generic; +using System.Reflection; +using System.Threading; +using System.Threading.Tasks; + +using Xrpl.Client; +using Xrpl.Models.Methods; + +using Timer = System.Timers.Timer; + +namespace Xrpl.Tests.ClientLib +{ + /// + /// A token that is already cancelled runs its registration callback inline, so the request is + /// rejected in the middle of being built — before its timeout timer exists. These tests pin + /// that nothing is left behind on that path: no timer in the timeout map, no pending promise. + /// + [TestClass] + public class TestURequestManagerCancellation + { + private static ConcurrentDictionary Timeouts(RequestManager manager) + { + FieldInfo field = typeof(RequestManager).GetField( + "timeoutsAwaitingResponse", + BindingFlags.Instance | BindingFlags.NonPublic); + + Assert.IsNotNull(field, "timeoutsAwaitingResponse is gone - update this test"); + return (ConcurrentDictionary)field.GetValue(manager); + } + + private static ConcurrentDictionary Promises(RequestManager manager) + { + FieldInfo field = typeof(RequestManager).GetField( + "promisesAwaitingResponse", + BindingFlags.Instance | BindingFlags.NonPublic); + + Assert.IsNotNull(field, "promisesAwaitingResponse is gone - update this test"); + return (ConcurrentDictionary)field.GetValue(manager); + } + + [TestMethod] + public async Task TestUCancelledTokenLeavesNoTimerBehindOnCreateRequest() + { + RequestManager manager = new RequestManager(); + using CancellationTokenSource cts = new CancellationTokenSource(); + cts.Cancel(); + + Dictionary request = new Dictionary { ["command"] = "ping" }; + RequestManager.XrplRequest created = manager.CreateRequest( + request, + TimeSpan.FromSeconds(30), + cts.Token); + + await Assert.ThrowsExactlyAsync(() => created.Promise); + + Assert.AreEqual(0, Timeouts(manager).Count, "timeout timer outlived the cancelled request"); + Assert.AreEqual(0, Promises(manager).Count, "promise outlived the cancelled request"); + } + + [TestMethod] + public async Task TestUCancelledTokenLeavesNoTimerBehindOnCreateGRequest() + { + RequestManager manager = new RequestManager(); + using CancellationTokenSource cts = new CancellationTokenSource(); + cts.Cancel(); + + RequestManager.XrplGRequest created = manager.CreateGRequest( + new PingRequest(), + TimeSpan.FromSeconds(30), + cts.Token); + + await Assert.ThrowsExactlyAsync(() => created.Promise); + + Assert.AreEqual(0, Timeouts(manager).Count, "timeout timer outlived the cancelled request"); + Assert.AreEqual(0, Promises(manager).Count, "promise outlived the cancelled request"); + } + + /// + /// The plain path still registers a timer and still cleans it up once the request finishes, + /// so the guard above cannot pass by never registering one in the first place. + /// + [TestMethod] + public void TestULiveRequestRegistersAndThenReleasesItsTimer() + { + RequestManager manager = new RequestManager(); + + Dictionary request = new Dictionary { ["command"] = "ping" }; + RequestManager.XrplRequest created = manager.CreateRequest(request, TimeSpan.FromSeconds(30)); + + Assert.AreEqual(1, Timeouts(manager).Count, "a live request must arm its timeout"); + + manager.Reject(created.Id, new OperationCanceledException("done")); + + Assert.AreEqual(0, Timeouts(manager).Count, "completing a request must release its timeout"); + Assert.AreEqual(0, Promises(manager).Count); + } + } +} diff --git a/Tests/Xrpl.Tests/Client/TestUWebSocketMessageAssembly.cs b/Tests/Xrpl.Tests/Client/TestUWebSocketMessageAssembly.cs index ecea00c2..4a4eee72 100644 --- a/Tests/Xrpl.Tests/Client/TestUWebSocketMessageAssembly.cs +++ b/Tests/Xrpl.Tests/Client/TestUWebSocketMessageAssembly.cs @@ -15,7 +15,10 @@ namespace Xrpl.Tests.ClientLib /// both that every message comes out byte-exact and that the cost per message does not grow /// with the number of chunks it was split into. /// + // GC.GetTotalAllocatedBytes is process-wide, so the allocation assertion below would pick up + // whatever other test classes allocate alongside it. [TestClass] + [DoNotParallelize] public class TestUWebSocketMessageAssembly { private const int WaitSeconds = 120; diff --git a/Xrpl/Client/RequestManager.cs b/Xrpl/Client/RequestManager.cs index 6dc9806d..2a18e1be 100644 --- a/Xrpl/Client/RequestManager.cs +++ b/Xrpl/Client/RequestManager.cs @@ -66,13 +66,11 @@ public void Resolve(Guid id, BaseResponse response) if (!promisesAwaitingResponse.TryGetValue(id, out var taskInfo) || taskInfo == null) { Debug.WriteLine($"Resolve called for non-existent promise {id} (likely already cancelled/timed out)"); + DisposeTimeout(id); return; } - // Dispose, not Stop: Timer is a finalizable Component, and a stopped but undisposed - // one per request piles up on the finalization queue over a long paged run. - if (timeoutsAwaitingResponse.TryRemove(id, out Timer timer)) - timer.Dispose(); + DisposeTimeout(id); try { @@ -99,12 +97,14 @@ public void Reject(Guid id, T error) where T : Exception if (!promisesAwaitingResponse.TryGetValue(id, out var taskInfo) || taskInfo == null) { Debug.WriteLine($"Reject called for non-existent promise {id} (likely already resolved)"); + + // A timer registered after its request had already finished has no other chance of + // being cleaned up: this is the Elapsed callback of exactly such a timer. + DisposeTimeout(id); return; } - // Dispose, not Stop: Timer is a finalizable Component, and a stopped but undisposed - // one per request piles up on the finalization queue over a long paged run. - if (timeoutsAwaitingResponse.TryRemove(id, out Timer timer)) - timer.Dispose(); + + DisposeTimeout(id); CompleteWithException(taskInfo, error); // Observe the exception to prevent UnobservedTaskException in consuming apps @@ -114,6 +114,19 @@ public void Reject(Guid id, T error) where T : Exception this.DeletePromise(id, taskInfo); } + /// + /// Removes the timeout timer of and disposes it. + /// Dispose, not Stop: is a finalizable Component, and a stopped but + /// undisposed one per request piles up on the finalization queue over a long paged run. + /// + private void DisposeTimeout(Guid id) + { + if (timeoutsAwaitingResponse.TryRemove(id, out Timer timer)) + { + timer.Dispose(); + } + } + /// /// Completes the pending request with a deserialized result, using the typed delegate /// captured when the request was created and falling back to reflection for @@ -280,6 +293,14 @@ public XrplGRequest CreateGRequest( } }); taskInfo.CancellationRegistration = registration; + + // An already cancelled token runs its callback inline, so Reject completed and + // removed the promise before the registration was stored — nothing would ever + // dispose it. + if (!promisesAwaitingResponse.ContainsKey(newId)) + { + _ = registration.DisposeAsync(); + } } if (timeout != System.Threading.Timeout.InfiniteTimeSpan) @@ -299,6 +320,15 @@ public XrplGRequest CreateGRequest( }; timer.Start(); timeoutsAwaitingResponse.TryAdd(newId, timer); + + // Same inline-cancellation case: the request may already be finished, and the + // Reject that finished it ran before this timer existed, so it had nothing to + // remove. Whatever is registered under this id now belongs to a request that is + // already gone. + if (!promisesAwaitingResponse.ContainsKey(newId)) + { + DisposeTimeout(newId); + } } return new XrplGRequest() @@ -360,6 +390,14 @@ public XrplRequest CreateRequest( } }); taskInfo.CancellationRegistration = registration; + + // An already cancelled token runs its callback inline, so Reject completed and + // removed the promise before the registration was stored — nothing would ever + // dispose it. + if (!promisesAwaitingResponse.ContainsKey(newId)) + { + _ = registration.DisposeAsync(); + } } if (timeout != System.Threading.Timeout.InfiniteTimeSpan) @@ -379,6 +417,15 @@ public XrplRequest CreateRequest( }; timer.Start(); timeoutsAwaitingResponse.TryAdd(newId, timer); + + // Same inline-cancellation case: the request may already be finished, and the + // Reject that finished it ran before this timer existed, so it had nothing to + // remove. Whatever is registered under this id now belongs to a request that is + // already gone. + if (!promisesAwaitingResponse.ContainsKey(newId)) + { + DisposeTimeout(newId); + } } return new XrplRequest() From 52361682f90c6cd44cf1c4333bc4bed1f857298c Mon Sep 17 00:00:00 2001 From: Aleksandr Platonenkov Date: Thu, 13 Aug 2026 15:29:36 -0300 Subject: [PATCH 8/8] =?UTF-8?q?docs(changes):=20=D0=BF=D0=B5=D1=80=D0=B5?= =?UTF-8?q?=D0=BD=D0=B5=D1=81=D1=82=D0=B8=20=D0=B7=D0=B0=D0=BF=D0=B8=D1=81?= =?UTF-8?q?=D0=B8=20=D0=B2=20=D1=80=D0=B0=D0=B7=D0=B4=D0=B5=D0=BB=2010.11.?= =?UTF-8?q?1.0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Раздел 10.11.0.0 уже выпущен, 10.11.1.0 открыт в dev и висит релизным PR #88 — записи изначально ушли не туда. Заодно запись про таймеры расширена случаем уже отменённого токена и добавлена запись про удалённое мёртвое поле tasks. --- CHANGES.md | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/CHANGES.md b/CHANGES.md index 0485439d..e1db87e1 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -7,6 +7,10 @@ * `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 * `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 +* **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 ## 10.11.0.0 08/04/2026 @@ -32,9 +36,6 @@ * `.ci-config/bump-nightly-pin.sh` does the move: newest `xrpld` build from the nightly apt channel, `ARG XRPLD_VERSION` rewritten, `rippled.batchv11.cfg` regenerated from the develop commit **encoded in that version string** — config and binary cannot drift apart, which is the failure mode the old manual two-step invited. `--check` reports the pin, the newest build and the pin's age without touching anything. Both timestamp formats are compared by their common `YYYYMMDDHHMM` prefix * the workflow bumps only once the pin is older than `MAX_PIN_AGE_DAYS` (21) — nightly publishes several builds a day, and a weekly PR would be noise rather than signal; `workflow_dispatch` takes a `force` input for the exceptions. It then builds and starts the stand on the new pin and requires the AMM sentinel amendment to come up enabled at genesis, which is what proves the regenerated config was accepted rather than silently ignored, and attaches the definitions diff against the new build to the PR body — a `node-only` field there is the SDK being behind develop, reported instead of hidden * credentials, idempotency and the tracking-issue fallback follow release-watch exactly, including the one-notification-per-failure-streak rule -* **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 -* **Request timeout timers were stopped but never disposed** — `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 -* **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 * **Cancellation no longer disappears into the autofill fee fallbacks** — `FetchCounterpartySignerCount` and `FetchLoan` wrap their client call in a broad `catch`, which is right for the case they exist for (the counterparty account or the Loan object is not there yet, and preclaim will report it) but also swallowed an `OperationCanceledException` raised from the caller's own token. Autofill then carried on and wrote a fee derived from the fallback — one signer, no loan — for a request the caller had already abandoned. Both catches now carry `when (!cancellationToken.IsCancellationRequested)`, which lets a caller's cancellation through while a client-side timeout, which does not cancel that token, still falls back as before. Covered in both directions: a cancelled token must throw and leave no `Fee` behind, an unreadable Loan object must still fall back * **`MPTokenIssuanceSet` validation reports a malformed `Flags` as `ValidationException`** — it went through `Convert.ToUInt32`, which throws `FormatException` or `InvalidCastException` on a non-numeric value, while the `ImmutableFlags` check two lines below reports `ValidationException` like the rest of the validators. Callers catching `ValidationException` did not catch the other two * **The conformance fixtures are re-pinned to the 3.3.0 tag** — `transactions.macro` and `LedgerFormats.h` now come from the release commit (`00a178fb`) instead of a July `develop` sha and `3.3.0-rc1`; `ledger_entries.macro` stays on `develop` (`9859e5ce`) for the reason its `.ref` already gives — `sfLEVersion` exists only there. Both macro files are byte-identical to upstream and re-verifiable with the `curl … | diff` line in each `.ref`. This is what makes the guards test against the version CI actually runs: