From 45d9c22b22f9144c947a051e32585d3905761103 Mon Sep 17 00:00:00 2001 From: Aleksandr Platonenkov Date: Sun, 16 Aug 2026 14:19:01 -0300 Subject: [PATCH 1/5] =?UTF-8?q?perf(client):=20=D0=BD=D0=B5=20=D1=80=D0=B0?= =?UTF-8?q?=D0=B7=D0=B1=D0=B8=D1=80=D0=B0=D1=82=D1=8C=20=D0=BA=D0=B0=D0=B6?= =?UTF-8?q?=D0=B4=D1=8B=D0=B9=20=D0=BE=D1=82=D0=B2=D0=B5=D1=82=20=D0=B4?= =?UTF-8?q?=D0=B2=D0=B0=D0=B6=D0=B4=D1=8B=20=D0=B8=20=D0=BD=D0=B5=20=D0=BA?= =?UTF-8?q?=D0=BE=D0=BF=D0=B8=D1=80=D0=BE=D0=B2=D0=B0=D1=82=D1=8C=20=D0=B5?= =?UTF-8?q?=D0=B3=D0=BE=20=D0=B2=20UTF-16?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit RequestManager.Resolve рендерил уже разобранный result обратно в строку (JsonElement.ToString) и парсил её повторно. На странице ledger_data при limit=2048 (~1 МБ) это 7.30 МБ аллокаций на ответ — 7.42x его размера, все четыре куска мимо порога 85 КБ, то есть в LOH. Теперь result десериализуется прямо из узла, а при запросе JsonElement/object узел отдаётся как есть. Приём ответа переведён на UTF-8: Connection слушает OnBinaryMessage, у IsLikelyResponse и HandleResponse появились перегрузки на ReadOnlySpan, а UTF-16 строка материализуется лениво и только там, где действительно нужен текст (стримы и колбэки предупреждений/ошибок). Замер на локальном WebSocket-сервере, 600 страниц по ~1 МБ: 8.32 -> 2.68 МБ на ответ, 11.49 -> 7.96 мс, пик кучи 42.4 -> 20.8 МБ, пик LOH 39.2 -> 17.3 МБ, RSS 203.4 -> 71.8 МБ. Под заниженным DOTNET_GCHeapHardLimit прежний путь воспроизводит продовый OOM тем же стеком, новый проходит. Заодно убран третий разбор того же сообщения в ветке status == "error". --- CHANGES.md | 10 + .../Xrpl.Tests/Client/TestUResponseParsing.cs | 220 ++++++++++++++++++ Xrpl/Client/RequestManager.cs | 90 +++++-- Xrpl/Client/connection.cs | 83 ++++++- Xrpl/Xrpl.csproj | 2 +- 5 files changed, 381 insertions(+), 24 deletions(-) create mode 100644 Tests/Xrpl.Tests/Client/TestUResponseParsing.cs diff --git a/CHANGES.md b/CHANGES.md index ad9e1f5c..26ffd8d0 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -1,5 +1,15 @@ # Changes +## 10.12.0.0 08/16/2026 + +* **Every response was parsed twice and copied to UTF-16 twice** — the cost of reading a response, measured rather than reasoned about. `RequestManager.Resolve` did `JsonSerializer.Deserialize(response.Result?.ToString() ?? "{}", taskInfo.Type, ...)`. `BaseResponse.Result` is typed `object`, which System.Text.Json fills with a `JsonElement` that already owns a private copy of the `result` bytes — so `.ToString()` rendered that element back into a UTF-16 string and the serializer parsed the string a second time. On a `ledger_data` page at `limit=2048` (~1 MB) the four stages measured, per response, at: 1.97 MB for the UTF-16 copy of the message, 1.68 MB for the document built over it, 1.97 MB for the UTF-16 copy of the `result`, 1.68 MB for the second document — **7.30 MB, 7.42x the response size**, all four allocations past the 85 KB large-object threshold. Both halves are now gone: + * `DeserializeResult` works off the parsed node: `element.Deserialize(type, options)` for a typed model, and the element itself when the request asked for `JsonElement` or `object`, which is what a consumer that needs the raw ledger objects asks for (the typed `LOLedgerData.State` drops unknown fields). A `BaseResponse` assembled by hand rather than parsed off the wire keeps the old string path. Behaviour is otherwise unchanged, including a missing or JSON-`null` `result`, which still yields what deserializing `"{}"` yielded + * the socket path carries the frame as it arrived. `Connection` binds `OnBinaryMessage` instead of `OnMessageReceived`, `IsLikelyResponse` and `RequestManager.HandleResponse` have `ReadOnlySpan` overloads, and the UTF-16 string is materialized — once, lazily — only for what genuinely needs text: stream messages and the `OnWarning`/`OnServerWarning`/`OnError` callbacks. The `string` overloads stay for `Connection.OnMessage(string)` and for external callers + * measured end to end against a local WebSocket server, 600 `ledger_data` pages of ~1 MB: **8.32 → 2.68 MB allocated per response** (8.46x → 2.72x the payload), 11.49 → 7.96 ms per response, 87 → 126 responses/s, peak managed heap 42.4 → 20.8 MB, peak LOH 39.2 → 17.3 MB, peak working set 203.4 → 71.8 MB. Under a lowered `DOTNET_GCHeapHardLimit` the pre-fix path reproduced the production failure exactly — `XrplException: Failed to deserialize response for request : Exception of type 'System.OutOfMemoryException' was thrown`, with `JsonElement.ToString()` at the top of the inner stack — at a ceiling the fixed path completes 15/15 pages under + * the win is not specific to `ledger_data` or to `JsonElement`: the second parse was on the path of every command. The repo's own `BenchmarkLedgerDataCrawl`, which goes through `Request` → `Dictionary`, drops from 22.9 to 14.8 MiB allocated per 2 MiB page (11.4x → 7.4x) with LOH ending at 38.2 instead of 115.4 MiB — it stays above the `JsonElement` figure because building a `Dictionary` boxes every value, which this change does not address + * `TestUResponseParsing` pins the behaviour that had to survive — the untyped node handed through is self-contained and readable after a forced gen2 collection, a typed model deserializes to the same values, the `string` and UTF-8 overloads agree, a missing `result` still completes, an `error` status still rejects with the parsed `ErrorResponse` attached — and holds the allocation budget at 4x the response size, measured per thread so the class-parallel test run cannot perturb it +* **`error` responses were deserialized a third time** — the `status == "error"` branch of `HandleResponse` re-parsed the whole message into an `ErrorResponse` inside a `try`/`catch` that swallowed everything, to build the exception's `Response`. The message had already been deserialized into an `ErrorResponse` at the top of the same method; the second parse only produced an equal copy, and on a large error payload it was a second large-object allocation on a path that is already failing + ## 10.11.1.0 08/13/2026 * **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: diff --git a/Tests/Xrpl.Tests/Client/TestUResponseParsing.cs b/Tests/Xrpl.Tests/Client/TestUResponseParsing.cs new file mode 100644 index 00000000..dc2fdac1 --- /dev/null +++ b/Tests/Xrpl.Tests/Client/TestUResponseParsing.cs @@ -0,0 +1,220 @@ +using Microsoft.VisualStudio.TestTools.UnitTesting; + +using System; +using System.Collections.Generic; +using System.Text; +using System.Text.Json; +using System.Threading; +using System.Threading.Tasks; + +using Xrpl.Client; +using Xrpl.Client.Exceptions; +using Xrpl.Models.Ledger; +using Xrpl.Models.Methods; + +namespace Xrpl.Tests.ClientLib +{ + /// + /// Covers how turns a response into the requested type. The + /// response arrives already parsed, so the result member is deserialized straight from + /// its : rendering it back to text and parsing it a second time used + /// to cost two extra copies of the whole response per request, both large-object-heap sized on + /// a paged crawl. These tests pin the behaviour that must survive that, and the allocation + /// budget that must not creep back up. + /// + [TestClass] + public class TestUResponseParsing + { + private static string BuildLedgerDataMessage(Guid id, int entries) + { + StringBuilder builder = new StringBuilder(entries * 128 + 256); + builder.Append("{\"id\":\"").Append(id.ToString("D")) + .Append("\",\"status\":\"success\",\"type\":\"response\",\"result\":{") + .Append("\"ledger_hash\":\"842B57C1CC0613299A686D3E9F310EC0422C84D3911E5056389AA7E5808A93C8\",") + .Append("\"ledger_index\":96000000,\"validated\":true,\"marker\":\"AABBCCDD\",\"state\":["); + + for (int i = 0; i < entries; i++) + { + if (i > 0) + { + builder.Append(','); + } + + builder.Append("{\"LedgerEntryType\":\"AccountRoot\",\"Account\":\"rN7n7otQDd6FczFgLdSqtcsAUxDkw6fzRH\",") + .Append("\"Balance\":\"").Append(1000000 + i) + .Append("\",\"Flags\":0,\"OwnerCount\":").Append(i % 17) + .Append(",\"Sequence\":").Append(i + 1) + .Append(",\"index\":\"").Append(i.ToString("X64")).Append("\"}"); + } + + builder.Append("]}}"); + return builder.ToString(); + } + + /// Rewrites the 36-character id of a prebuilt message in place. + private static void WriteId(byte[] message, int offset, Guid id) + { + string text = id.ToString("D"); + for (int i = 0; i < text.Length; i++) + { + message[offset + i] = (byte)text[i]; + } + } + + private static RequestManager.XrplGRequest Pending(RequestManager manager) + { + return manager.CreateGRequest( + new LedgerDataRequest { Limit = 4 }, + System.Threading.Timeout.InfiniteTimeSpan); + } + + [TestMethod] + public void TestUntypedRequestGetsTheParsedResultNode() + { + RequestManager manager = new RequestManager(); + RequestManager.XrplGRequest pending = Pending(manager); + + manager.HandleResponse(BuildLedgerDataMessage(pending.Id, 4)); + + JsonElement result = (JsonElement)pending.Promise.GetAwaiter().GetResult(); + Assert.AreEqual(JsonValueKind.Object, result.ValueKind); + Assert.AreEqual(4, result.GetProperty("state").GetArrayLength()); + Assert.AreEqual(96000000, result.GetProperty("ledger_index").GetInt32()); + Assert.AreEqual("AABBCCDD", result.GetProperty("marker").GetString()); + + // The element must outlive the parse: it is handed out rather than copied, so it has + // to own its data and stay readable after everything else is collected. + GC.Collect(2, GCCollectionMode.Forced, blocking: true); + Assert.AreEqual(4, result.GetProperty("state").GetArrayLength()); + } + + [TestMethod] + public void TestTypedRequestDeserializesFromTheParsedResultNode() + { + RequestManager manager = new RequestManager(); + RequestManager.XrplGRequest pending = Pending(manager); + + manager.HandleResponse(BuildLedgerDataMessage(pending.Id, 3)); + + LOLedgerData result = (LOLedgerData)pending.Promise.GetAwaiter().GetResult(); + Assert.IsNotNull(result); + Assert.AreEqual(96000000u, result.LedgerIndex); + Assert.AreEqual("842B57C1CC0613299A686D3E9F310EC0422C84D3911E5056389AA7E5808A93C8", result.LedgerHash); + Assert.IsNotNull(result.State); + Assert.AreEqual(3, result.State.Count); + } + + [TestMethod] + public void TestUtf8AndStringOverloadsProduceTheSameResult() + { + RequestManager manager = new RequestManager(); + + RequestManager.XrplGRequest viaString = Pending(manager); + string message = BuildLedgerDataMessage(viaString.Id, 5); + manager.HandleResponse(message); + + RequestManager.XrplGRequest viaBytes = Pending(manager); + manager.HandleResponse(Encoding.UTF8.GetBytes(BuildLedgerDataMessage(viaBytes.Id, 5))); + + LOLedgerData fromString = (LOLedgerData)viaString.Promise.GetAwaiter().GetResult(); + LOLedgerData fromBytes = (LOLedgerData)viaBytes.Promise.GetAwaiter().GetResult(); + + Assert.AreEqual(fromString.LedgerIndex, fromBytes.LedgerIndex); + Assert.AreEqual(fromString.LedgerHash, fromBytes.LedgerHash); + Assert.AreEqual(fromString.State.Count, fromBytes.State.Count); + } + + [TestMethod] + public void TestResponseWithoutResultStillCompletes() + { + RequestManager manager = new RequestManager(); + + RequestManager.XrplGRequest untyped = Pending(manager); + manager.HandleResponse($"{{\"id\":\"{untyped.Id:D}\",\"status\":\"success\",\"type\":\"response\",\"result\":null}}"); + JsonElement empty = (JsonElement)untyped.Promise.GetAwaiter().GetResult(); + Assert.AreEqual(JsonValueKind.Object, empty.ValueKind); + Assert.IsFalse(empty.TryGetProperty("state", out _)); + + RequestManager.XrplGRequest typed = Pending(manager); + manager.HandleResponse($"{{\"id\":\"{typed.Id:D}\",\"status\":\"success\",\"type\":\"response\"}}"); + LOLedgerData defaults = (LOLedgerData)typed.Promise.GetAwaiter().GetResult(); + Assert.IsNotNull(defaults); + Assert.IsNull(defaults.State); + } + + [TestMethod] + public void TestErrorStatusRejectsWithTheParsedErrorResponse() + { + RequestManager manager = new RequestManager(); + RequestManager.XrplGRequest pending = Pending(manager); + + manager.HandleResponse( + $"{{\"id\":\"{pending.Id:D}\",\"status\":\"error\",\"type\":\"response\"," + + "\"error\":\"lgrNotFound\",\"error_message\":\"ledgerNotFound\"}"); + + RippledException rippled = null; + try + { + pending.Promise.Wait(); + } + catch (AggregateException raised) + { + rippled = raised.InnerException as RippledException; + } + + Assert.IsNotNull(rippled, "the request should have been rejected with a RippledException"); + StringAssert.Contains(rippled.Message, "lgrNotFound"); + Assert.IsNotNull(rippled.Response, "the parsed error response must be attached"); + Assert.AreEqual("lgrNotFound", rippled.Response.Error); + Assert.AreEqual("ledgerNotFound", rippled.Response.ErrorMessage); + } + + /// + /// Guards the allocation budget of the response path. Before the result node was + /// deserialized directly, one response cost about 7.4 times its own byte length: a UTF-16 + /// copy of the message, a document over it, a second UTF-16 copy of the result, and a + /// second document over that. The direct path costs about half of it from a string and + /// about 1.7 times from UTF-8 bytes. The bound below sits between the two, well clear of + /// both, so it fails only if the double round-trip comes back. + /// + [TestMethod] + public void TestResponseParsingStaysWithinItsAllocationBudget() + { + const int Entries = 4096; + const int Rounds = 12; + + RequestManager manager = new RequestManager(); + + // Built once and reused, with only the id rewritten in place, so nothing the harness + // allocates lands inside the measured window. + RequestManager.XrplGRequest warmup = Pending(manager); + byte[] message = Encoding.UTF8.GetBytes(BuildLedgerDataMessage(warmup.Id, Entries)); + const int IdOffset = 7; // past {"id":" + manager.HandleResponse(message); + _ = warmup.Promise.GetAwaiter().GetResult(); + + GC.Collect(2, GCCollectionMode.Forced, blocking: true); + long before = GC.GetAllocatedBytesForCurrentThread(); + + for (int i = 0; i < Rounds; i++) + { + RequestManager.XrplGRequest pending = Pending(manager); + WriteId(message, IdOffset, pending.Id); + manager.HandleResponse(message); + JsonElement result = (JsonElement)pending.Promise.GetAwaiter().GetResult(); + Assert.AreEqual(Entries, result.GetProperty("state").GetArrayLength()); + } + + long allocated = GC.GetAllocatedBytesForCurrentThread() - before; + double perResponse = allocated / (double)Rounds; + double ratio = perResponse / message.Length; + + Console.WriteLine($"response {message.Length:N0} bytes, {perResponse / 1024 / 1024:F2} MB allocated per response ({ratio:F2}x)"); + + Assert.IsTrue( + ratio < 4.0, + $"response parsing allocated {ratio:F2}x the response size, budget is 4x " + + "(the pre-fix double round-trip cost about 7x here)"); + } + } +} diff --git a/Xrpl/Client/RequestManager.cs b/Xrpl/Client/RequestManager.cs index 2a18e1be..74c1b60d 100644 --- a/Xrpl/Client/RequestManager.cs +++ b/Xrpl/Client/RequestManager.cs @@ -55,6 +55,20 @@ public class XrplGRequest private readonly ConcurrentDictionary promisesAwaitingResponse = new ConcurrentDictionary(); private readonly JsonSerializerOptions serializerOptions = XrplJsonOptions.Default; + /// + /// Stands in for a missing result, matching what deserializing the literal + /// "{}" used to produce. + /// + private static readonly JsonElement EmptyResult = ParseEmptyObject(); + + private static JsonElement ParseEmptyObject() + { + using (JsonDocument document = JsonDocument.Parse("{}")) + { + return document.RootElement.Clone(); + } + } + public RequestManager() { } @@ -74,7 +88,7 @@ public void Resolve(Guid id, BaseResponse response) try { - object deserialized = JsonSerializer.Deserialize(response.Result?.ToString() ?? "{}", taskInfo.Type, serializerOptions); + object deserialized = DeserializeResult(response.Result, taskInfo.Type); CompleteWithResult(taskInfo, deserialized); this.DeletePromise(id, taskInfo); } @@ -86,6 +100,47 @@ public void Resolve(Guid id, BaseResponse response) } } + /// + /// Converts the result member of a response into the type the request was created + /// with. + /// + /// + /// The member arrives already parsed - is typed + /// , which System.Text.Json fills with a self-contained + /// . Rendering that element back to a string and parsing the + /// string a second time cost two more copies of the whole response per request: a UTF-16 + /// string at twice the byte length, and a second document on top of it. On a paged walk + /// of the ledger both copies are large-object-heap sized and both are pure waste, so the + /// element is deserialized directly instead, and handed straight through when the request + /// asked for the untyped node in the first place. + /// + private object DeserializeResult(object result, Type type) + { + JsonElement element; + + if (result is null) + { + element = EmptyResult; + } + else if (result is JsonElement parsed) + { + element = parsed; + } + else + { + // A response assembled by hand rather than parsed off the wire: there is no node + // to reuse, so this is still the only way in. + return JsonSerializer.Deserialize(result.ToString(), type, serializerOptions); + } + + if (type == typeof(JsonElement) || type == typeof(object)) + { + return element; + } + + return element.Deserialize(type, serializerOptions); + } + /// /// Rejects a pending request with the specified exception. /// Safe to call even if the promise no longer exists (e.g., already resolved). @@ -444,8 +499,25 @@ public XrplRequest CreateRequest( /// public (BaseResponse Response, bool Handled) HandleResponse(string message) { - var response = JsonSerializer.Deserialize(message, serializerOptions); + return HandleResponse(JsonSerializer.Deserialize(message, serializerOptions)); + } + + /// + /// Same as for a message still in its wire form. + /// + /// + /// Preferred on the socket path: transcoding the frame to a UTF-16 string first costs a + /// copy at twice the byte length of the message, which for a large response is a + /// large-object-heap allocation spent only to hand System.Text.Json something it converts + /// straight back to UTF-8. + /// + public (BaseResponse Response, bool Handled) HandleResponse(ReadOnlySpan utf8Message) + { + return HandleResponse(JsonSerializer.Deserialize(utf8Message, serializerOptions)); + } + private (BaseResponse Response, bool Handled) HandleResponse(ErrorResponse response) + { if (response.Id == null) { return (response, false); @@ -481,21 +553,13 @@ public XrplRequest CreateRequest( if (response.Status == "error" ) { - ErrorResponse errorResponse = null; - try - { - errorResponse = JsonSerializer.Deserialize(message, serializerOptions); - - } - catch (Exception e) - { - - } + // The message was already deserialized into an ErrorResponse above, so the error + // details are in hand - parsing it a second time only produced an equal copy. string detail = response.ErrorMessage ?? response.ErrorException; var errMessage = response.Error is null ? detail : $"{response.Error} - {detail}"; - var error = new RippledException(errMessage, errorResponse); + var error = new RippledException(errMessage, response); this.Reject(id, error); return (response, true); } diff --git a/Xrpl/Client/connection.cs b/Xrpl/Client/connection.cs index 84afbe80..12ab7e03 100644 --- a/Xrpl/Client/connection.cs +++ b/Xrpl/Client/connection.cs @@ -3,6 +3,7 @@ using System.Diagnostics; using System.IO; using System.Net.WebSockets; +using System.Text; using System.Text.Json; using System.Threading; using System.Threading.Channels; @@ -1060,7 +1061,10 @@ await errorHandler.Invoke( } }); - ws.OnMessageReceived(async (m, ws) => + // Bound to the binary callback rather than the string one: the frame is already UTF-8 + // and that is what the JSON reader wants, so the UTF-16 copy of every message - twice + // the byte length, on the large object heap for a big response - is never made. + ws.OnBinaryMessage(async (m, ws) => { try { @@ -2863,6 +2867,36 @@ private bool IsLikelyResponse(string message) return false; } + /// + /// over the raw frame, so the discriminator scan does + /// not force a UTF-16 copy of the message. Byte-wise scanning is equivalent here: the tokens + /// looked for are ASCII, and UTF-8 never encodes them inside a multi-byte sequence. + /// + private bool IsLikelyResponse(ReadOnlySpan utf8Message) + { + if (utf8Message.Length < 10) + return false; + + int firstBrace = utf8Message.IndexOf((byte)'{'); + if (firstBrace < 0 || firstBrace + 10 >= utf8Message.Length) + return false; + + ReadOnlySpan rest = utf8Message.Slice(firstBrace + 1); + int idIndex = rest.IndexOf("\"id\""u8); + if (idIndex < 0) + return false; + + int checkEnd = Math.Min(rest.Length, idIndex + 10); + for (int i = idIndex + 4; i < checkEnd; i++) + { + byte c = rest[i]; + if (c == (byte)':') return true; // This is a response + if (c != (byte)' ' && c != (byte)'\t' && c != (byte)'\n' && c != (byte)'\r') break; + } + + return false; + } + /// /// Starts the background message processor for stream messages. /// Creates a new session-bound channel and processor task. @@ -3127,13 +3161,39 @@ private async Task ProcessStreamMessageAsync(string message) /// concurrently from the ThreadPool. Handler implementations MUST be thread-safe /// or marshal to their own synchronization context (e.g., UI thread). /// - private async Task IOnMessageFastPath(string message) + private Task IOnMessageFastPath(string message) + { + return IOnMessageFastPath(message, null); + } + + /// + /// Overload for a message still in its wire form, used by the socket callback. See + /// for why the bytes are kept as they are. + /// + private Task IOnMessageFastPath(byte[] utf8Message) + { + return IOnMessageFastPath(null, utf8Message); + } + + /// + /// Exactly one of and carries the + /// message; the other is null. + /// + /// + /// A response is parsed straight out of when it is the one + /// present, so the UTF-16 copy of the message - twice its byte length - is never made for the + /// common case. Everything that genuinely needs text (stream messages, the warning and error + /// callbacks) asks for it through Text(), which materializes it once and only then. + /// + private async Task IOnMessageFastPath(string message, byte[] utf8Message) { lastActivityTime = DateTime.UtcNow; + string Text() => message ??= Encoding.UTF8.GetString(utf8Message); + // Scan message for "id" property to detect response messages - var isResponse = IsLikelyResponse(message); - + var isResponse = utf8Message is null ? IsLikelyResponse(message) : IsLikelyResponse(utf8Message); + if (isResponse) { // This is a response (including ping/pong) - process immediately with full parsing @@ -3141,20 +3201,23 @@ private async Task IOnMessageFastPath(string message) BaseResponse data; bool handled; try - { + { // FIRST: Handle response immediately to unblock any waiting requests (like ping) // This is the most time-critical operation - (data, handled) = requestManager.HandleResponse(message); + (data, handled) = utf8Message is null + ? requestManager.HandleResponse(message) + : requestManager.HandleResponse(utf8Message); } catch (Exception error) { var errInfo = XrplErrorClassifier.Classify(error); + var capturedText = Text(); // Fire-and-forget for error callback - don't block _ = Task.Run(async () => { if (OnError is not null) { - await OnError.Invoke(error: "error", errorMessage: "badMessage", errInfo.UserMessage, message); + await OnError.Invoke(error: "error", errorMessage: "badMessage", errInfo.UserMessage, capturedText); } }); return; @@ -3164,7 +3227,7 @@ private async Task IOnMessageFastPath(string message) { // Message has "id" but no matching pending request — this is an async // follow-up (e.g. path_find updates). Route to stream processing. - EnqueueStreamMessage(message); + EnqueueStreamMessage(Text()); return; } @@ -3173,7 +3236,7 @@ private async Task IOnMessageFastPath(string message) if (data.Warning != null || data.Warnings is { Count: > 0 }) { var capturedData = data; - var capturedMessage = message; + var capturedMessage = Text(); _ = Task.Run(async () => { if (capturedData.Warning != null && OnWarning is not null) @@ -3192,7 +3255,7 @@ private async Task IOnMessageFastPath(string message) { // This is a stream message (no "id") - process asynchronously // to avoid blocking the receive loop and causing ping timeouts - EnqueueStreamMessage(message); + EnqueueStreamMessage(Text()); } } diff --git a/Xrpl/Xrpl.csproj b/Xrpl/Xrpl.csproj index 5cc56773..1c78a156 100644 --- a/Xrpl/Xrpl.csproj +++ b/Xrpl/Xrpl.csproj @@ -14,7 +14,7 @@ Apache-2.0 https://github.com/StaticBit-io/XrplCSharp XrplCSharp - 10.11.1.0 + 10.12.0.0 From 9708967df2d33399246944f995b1c28f1849e989 Mon Sep 17 00:00:00 2001 From: Aleksandr Platonenkov Date: Sun, 16 Aug 2026 14:22:00 -0300 Subject: [PATCH 2/5] =?UTF-8?q?test(client):=20=D1=80=D0=B5=D0=B3=D1=80?= =?UTF-8?q?=D0=B5=D1=81=D1=81=D0=B8=D1=8F=20=D0=BD=D0=B0=20=D1=80=D0=B0?= =?UTF-8?q?=D0=B7=D0=B1=D0=BE=D1=80=20=D0=BE=D1=82=D0=B2=D0=B5=D1=82=D0=B0?= =?UTF-8?q?=20=D0=B8=20=D0=B1=D0=B5=D0=BD=D1=87=D0=BC=D0=B0=D1=80=D0=BA=20?= =?UTF-8?q?=D0=BD=D0=B5=D1=82=D0=B8=D0=BF=D0=B8=D0=B7=D0=B8=D1=80=D0=BE?= =?UTF-8?q?=D0=B2=D0=B0=D0=BD=D0=BD=D0=BE=D0=B3=D0=BE=20=D0=BE=D0=B1=D1=85?= =?UTF-8?q?=D0=BE=D0=B4=D0=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit TestUResponseParsing фиксирует поведение, которое обязано было пережить правку: отданный наружу JsonElement самодостаточен и читается после принудительной сборки gen2, типизированная модель даёт те же значения, перегрузки на string и UTF-8 совпадают, отсутствующий result по-прежнему завершает запрос, а status "error" по-прежнему отклоняет его с разобранным ErrorResponse. Плюс бюджет аллокаций в 4x от размера ответа — счётчик берётся потоковый, иначе параллельный по классам прогон тестов ломает замер. BenchmarkSequentialPagingUntyped гоняет тот же обход через GRequest — путь потребителя, которому нужны сырые объекты леджера. Как и соседний бенчмарк, вне фильтров TestU/TestI. --- .../Client/BenchmarkLedgerDataCrawl.cs | 76 +++++++++++++++++++ 1 file changed, 76 insertions(+) diff --git a/Tests/Xrpl.Tests/Client/BenchmarkLedgerDataCrawl.cs b/Tests/Xrpl.Tests/Client/BenchmarkLedgerDataCrawl.cs index dc86527a..18521c80 100644 --- a/Tests/Xrpl.Tests/Client/BenchmarkLedgerDataCrawl.cs +++ b/Tests/Xrpl.Tests/Client/BenchmarkLedgerDataCrawl.cs @@ -5,7 +5,11 @@ using System.Diagnostics; using System.Threading.Tasks; +using System.Text.Json; + using Xrpl.Client; +using Xrpl.Models.Common; +using Xrpl.Models.Methods; namespace Xrpl.Tests.ClientLib { @@ -112,6 +116,78 @@ public async Task BenchmarkSequentialPaging() Console.WriteLine("decile profile (ms/page): " + string.Join(" | ", DecileProfile(pageMs))); } + /// + /// Same crawl through GRequest<JsonElement, …> — the path a consumer takes when + /// it needs the raw ledger objects, because the typed LOLedgerData.State drops + /// fields the models do not know. This is where the response path's own cost shows up + /// undiluted: nothing is materialized into a model, so what is measured is receive, + /// route and parse. + /// + [TestMethod] + public async Task BenchmarkSequentialPagingUntyped() + { + int pages = EnvInt("CRAWL_PAGES", 2000); + int payloadBytes = EnvInt("CRAWL_PAYLOAD_BYTES", 2 * 1024 * 1024); + int fragments = EnvInt("CRAWL_FRAGMENTS", 32); + + using PagedResponseServer server = new PagedResponseServer(payloadBytes, fragments); + using XrplClient client = new XrplClient(server.Url); + + await client.Connect().ConfigureAwait(false); + await RequestUntypedPageAsync(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 gen2Before = GC.CollectionCount(2); + long lohBefore = LohBytes(); + long startTicks = Stopwatch.GetTimestamp(); + + for (int i = 0; i < pages; i++) + { + await RequestUntypedPageAsync(client).ConfigureAwait(false); + } + + double totalSeconds = (Stopwatch.GetTimestamp() - startTicks) / (double)Stopwatch.Frequency; + long allocated = GC.GetTotalAllocatedBytes(precise: true) - allocatedBefore; + long lohAfter = LohBytes(); + + await client.Disconnect().ConfigureAwait(false); + + Console.WriteLine("=== ledger_data crawl benchmark (untyped JsonElement result) ==="); + Console.WriteLine($"pages : {pages}"); + Console.WriteLine($"payload : {payloadBytes / 1024.0 / 1024.0:F2} MiB"); + Console.WriteLine($"total time : {totalSeconds:F2} s ({pages / totalSeconds:F2} pages/s)"); + Console.WriteLine($"allocated per page: {allocated / (double)pages / 1024.0 / 1024.0:F1} MiB " + + $"({allocated / (double)pages / payloadBytes:F1}x payload)"); + Console.WriteLine($"gen2 collections : {GC.CollectionCount(2) - gen2Before}"); + Console.WriteLine($"LOH size : {lohBefore / 1024.0 / 1024.0:F1} -> {lohAfter / 1024.0 / 1024.0:F1} MiB"); + } + + private static async Task RequestUntypedPageAsync(XrplClient client) + { + JsonElement result = await client + .GRequest(new LedgerDataRequest + { + LedgerIndex = new LedgerIndex(96000000), + Binary = true, + Limit = 2048 + }) + .ConfigureAwait(false); + + if (result.ValueKind != JsonValueKind.Object || !result.TryGetProperty("state", out JsonElement state)) + { + throw new InvalidOperationException("empty ledger_data response"); + } + + if (state.GetArrayLength() == 0) + { + throw new InvalidOperationException("ledger_data page carried no objects"); + } + } + private static async Task RequestPageAsync(XrplClient client) { Dictionary request = new Dictionary From 191a6cb810b5d85d177a8c6980daa96ce753ca1f Mon Sep 17 00:00:00 2001 From: Aleksandr Platonenkov Date: Sun, 16 Aug 2026 14:56:31 -0300 Subject: [PATCH 3/5] =?UTF-8?q?fix(client):=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=82=D1=87=D1=91=D1=82=20?= =?UTF-8?q?=D0=BE=D0=B1=20=D0=BE=D1=88=D0=B8=D0=B1=D0=BA=D0=B5=20=D1=80?= =?UTF-8?q?=D0=B0=D0=B7=D0=B1=D0=BE=D1=80=D0=B0=20=D0=B8=20=D0=B7=D0=B0?= =?UTF-8?q?=D1=84=D0=B8=D0=BA=D1=81=D0=B8=D1=80=D0=BE=D0=B2=D0=B0=D1=82?= =?UTF-8?q?=D1=8C=20UTF-8=20=D0=BF=D1=83=D1=82=D1=8C=20=D1=82=D0=B5=D1=81?= =?UTF-8?q?=D1=82=D0=BE=D0=BC?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit По итогам ревью собственного PR, четыре пункта. Отчёт об ошибке разбора больше не падает вместе с тем, что его вызвало. Ответ, который не разобрался, — это чаще всего кончившаяся куча, а материализация сообщения для OnError была на этом пути самой крупной аллокацией: второй OOM внутри catch гасился колбэком сокета, и потребитель не видел ничего. Текст строится только когда обработчик подписан, а OutOfMemoryException при его построении подменяется литеральной заглушкой. Connection.OnMessage(null) снова уходит в OnError, а не бросает ArgumentNullException из точки входа: Text() больше не разыменовывает utf8Message без проверки. TestSocketPathKeepsResponsesInTheirWireForm гоняет 20 страниц через Connection по настоящему сокету и держит бюджет 4x от полезной нагрузки. Ни один существующий тест не видел, какую перегрузку выбирает клиент, — возврат на OnMessageReceived прошёл бы молча. Замер разделяет варианты с запасом с обеих сторон: 2.18x как есть, 4.84x со строковым колбэком. Счётчик здесь процессный (аллокации идут в цикле приёма), поэтому тест вынесен из параллельного прогона, а PagedResponseServer переиспользует один кадр на соединение и переписывает id на месте, чтобы сервер не попадал в замер клиента. Диагностика в колбэке сокета называла OnMessageReceived, хотя подписка давно через OnBinaryMessage. --- CHANGES.md | 3 +- .../Xrpl.Tests/Client/PagedResponseServer.cs | 40 +++++++++-- .../Xrpl.Tests/Client/TestUResponseParsing.cs | 72 +++++++++++++++++++ Xrpl/Client/connection.cs | 41 ++++++++++- 4 files changed, 148 insertions(+), 8 deletions(-) diff --git a/CHANGES.md b/CHANGES.md index 26ffd8d0..62082652 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -5,9 +5,10 @@ * **Every response was parsed twice and copied to UTF-16 twice** — the cost of reading a response, measured rather than reasoned about. `RequestManager.Resolve` did `JsonSerializer.Deserialize(response.Result?.ToString() ?? "{}", taskInfo.Type, ...)`. `BaseResponse.Result` is typed `object`, which System.Text.Json fills with a `JsonElement` that already owns a private copy of the `result` bytes — so `.ToString()` rendered that element back into a UTF-16 string and the serializer parsed the string a second time. On a `ledger_data` page at `limit=2048` (~1 MB) the four stages measured, per response, at: 1.97 MB for the UTF-16 copy of the message, 1.68 MB for the document built over it, 1.97 MB for the UTF-16 copy of the `result`, 1.68 MB for the second document — **7.30 MB, 7.42x the response size**, all four allocations past the 85 KB large-object threshold. Both halves are now gone: * `DeserializeResult` works off the parsed node: `element.Deserialize(type, options)` for a typed model, and the element itself when the request asked for `JsonElement` or `object`, which is what a consumer that needs the raw ledger objects asks for (the typed `LOLedgerData.State` drops unknown fields). A `BaseResponse` assembled by hand rather than parsed off the wire keeps the old string path. Behaviour is otherwise unchanged, including a missing or JSON-`null` `result`, which still yields what deserializing `"{}"` yielded * the socket path carries the frame as it arrived. `Connection` binds `OnBinaryMessage` instead of `OnMessageReceived`, `IsLikelyResponse` and `RequestManager.HandleResponse` have `ReadOnlySpan` overloads, and the UTF-16 string is materialized — once, lazily — only for what genuinely needs text: stream messages and the `OnWarning`/`OnServerWarning`/`OnError` callbacks. The `string` overloads stay for `Connection.OnMessage(string)` and for external callers + * the failure report survives the failure. A response that will not parse is most often a heap that has just run out, and materializing the message for `OnError` is then the largest allocation left on the path — if it throws, the notification is lost inside the handler and the consumer sees silence. The text is now built only when a handler is attached, and an `OutOfMemoryException` while building it falls back to a literal placeholder so the classification still goes out. `Connection.OnMessage(null)` also keeps its old route through `OnError` instead of throwing `ArgumentNullException` out of the entry point * measured end to end against a local WebSocket server, 600 `ledger_data` pages of ~1 MB: **8.32 → 2.68 MB allocated per response** (8.46x → 2.72x the payload), 11.49 → 7.96 ms per response, 87 → 126 responses/s, peak managed heap 42.4 → 20.8 MB, peak LOH 39.2 → 17.3 MB, peak working set 203.4 → 71.8 MB. Under a lowered `DOTNET_GCHeapHardLimit` the pre-fix path reproduced the production failure exactly — `XrplException: Failed to deserialize response for request : Exception of type 'System.OutOfMemoryException' was thrown`, with `JsonElement.ToString()` at the top of the inner stack — at a ceiling the fixed path completes 15/15 pages under * the win is not specific to `ledger_data` or to `JsonElement`: the second parse was on the path of every command. The repo's own `BenchmarkLedgerDataCrawl`, which goes through `Request` → `Dictionary`, drops from 22.9 to 14.8 MiB allocated per 2 MiB page (11.4x → 7.4x) with LOH ending at 38.2 instead of 115.4 MiB — it stays above the `JsonElement` figure because building a `Dictionary` boxes every value, which this change does not address - * `TestUResponseParsing` pins the behaviour that had to survive — the untyped node handed through is self-contained and readable after a forced gen2 collection, a typed model deserializes to the same values, the `string` and UTF-8 overloads agree, a missing `result` still completes, an `error` status still rejects with the parsed `ErrorResponse` attached — and holds the allocation budget at 4x the response size, measured per thread so the class-parallel test run cannot perturb it + * `TestUResponseParsing` pins the behaviour that had to survive — the untyped node handed through is self-contained and readable after a forced gen2 collection, a typed model deserializes to the same values, the `string` and UTF-8 overloads agree, a missing `result` still completes, an `error` status still rejects with the parsed `ErrorResponse` attached, a null message does not throw out of the entry point — and holds two allocation budgets at 4x the response size. The first measures `RequestManager` alone, per thread so the class-parallel run cannot perturb it (1.89x now). The second runs 20 pages through `Connection` over a real socket, because nothing else in the suite can see *which* overload the client picks: it reads the process-wide counter and is therefore kept out of the parallel pass, and it separates the two paths with room on both sides — 2.18x as bound, 4.84x with the string callback bound instead. `PagedResponseServer` reuses one response frame per connection and rewrites the id in place so the server contributes nothing to what the client is measured on * **`error` responses were deserialized a third time** — the `status == "error"` branch of `HandleResponse` re-parsed the whole message into an `ErrorResponse` inside a `try`/`catch` that swallowed everything, to build the exception's `Response`. The message had already been deserialized into an `ErrorResponse` at the top of the same method; the second parse only produced an equal copy, and on a large error payload it was a second large-object allocation on a path that is already failing ## 10.11.1.0 08/13/2026 diff --git a/Tests/Xrpl.Tests/Client/PagedResponseServer.cs b/Tests/Xrpl.Tests/Client/PagedResponseServer.cs index ce23bd38..8b5ff012 100644 --- a/Tests/Xrpl.Tests/Client/PagedResponseServer.cs +++ b/Tests/Xrpl.Tests/Client/PagedResponseServer.cs @@ -15,6 +15,9 @@ namespace Xrpl.Tests /// internal sealed class PagedResponseServer : WebSocketTestServerBase { + /// The id `Connection` sends: a quoted GUID in `D` format, always 38 characters. + private const int QuotedGuidLength = 38; + private readonly int _fragments; private readonly string _resultBody; private int _served; @@ -79,6 +82,9 @@ private static void AppendHex(StringBuilder builder, int seed, int length) protected override async Task ServeAsync(NetworkStream stream) { + // One frame buffer per connection, so two clients cannot rewrite each other's id. + byte[]? frame = null; + while (!Token.IsCancellationRequested) { string? message = await ReadTextFrameAsync(stream).ConfigureAwait(false); @@ -88,15 +94,41 @@ protected override async Task ServeAsync(NetworkStream stream) } string id = ExtractId(message); - string envelope = "{\"id\":" + id + ",\"status\":\"success\",\"type\":\"response\",\"result\":" + - _resultBody + "}"; - await WriteFragmentedMessageAsync(stream, Encoding.UTF8.GetBytes(envelope), _fragments) - .ConfigureAwait(false); + await WriteFragmentedMessageAsync(stream, BuildFrame(id, ref frame), _fragments).ConfigureAwait(false); Interlocked.Increment(ref _served); } } + /// + /// Builds the response frame for . The usual case - a quoted GUID, the + /// only form Connection sends - reuses one buffer and rewrites the id in place, so + /// the server contributes nothing to what a caller measures on the client side. Anything + /// else falls back to assembling the envelope. + /// + private byte[] BuildFrame(string id, ref byte[]? reusable) + { + if (id.Length != QuotedGuidLength) + { + return Encoding.UTF8.GetBytes(Envelope(id)); + } + + reusable ??= Encoding.UTF8.GetBytes(Envelope("\"" + new string('0', QuotedGuidLength - 2) + "\"")); + + int idOffset = "{\"id\":".Length; + for (int i = 0; i < QuotedGuidLength; i++) + { + reusable[idOffset + i] = (byte)id[i]; + } + + return reusable; + } + + private string Envelope(string id) + { + return "{\"id\":" + id + ",\"status\":\"success\",\"type\":\"response\",\"result\":" + _resultBody + "}"; + } + /// Pulls the JSON string value of the request's "id" property. private static string ExtractId(string message) { diff --git a/Tests/Xrpl.Tests/Client/TestUResponseParsing.cs b/Tests/Xrpl.Tests/Client/TestUResponseParsing.cs index dc2fdac1..9009781c 100644 --- a/Tests/Xrpl.Tests/Client/TestUResponseParsing.cs +++ b/Tests/Xrpl.Tests/Client/TestUResponseParsing.cs @@ -216,5 +216,77 @@ public void TestResponseParsingStaysWithinItsAllocationBudget() $"response parsing allocated {ratio:F2}x the response size, budget is 4x " + "(the pre-fix double round-trip cost about 7x here)"); } + + /// + /// The same budget one level up, over a real socket, because the budget above cannot see + /// which overload chooses. Binding the string callback again + /// would put a UTF-16 copy of every frame back on the path and nothing else in the suite + /// would notice. + /// + /// + /// Allocations here happen on the receive loop's thread, so this has to read the + /// process-wide counter, which is why the test is kept out of the parallel pass. + /// + [TestMethod] + [DoNotParallelize] + public async Task TestSocketPathKeepsResponsesInTheirWireForm() + { + const int Pages = 20; + const int PayloadBytes = 1024 * 1024; + + using PagedResponseServer server = new PagedResponseServer(PayloadBytes, fragments: 8); + using XrplClient client = new XrplClient(server.Url); + + await client.Connect().ConfigureAwait(false); + await CrawlPageAsync(client).ConfigureAwait(false); + + GC.Collect(2, GCCollectionMode.Forced, blocking: true); + GC.WaitForPendingFinalizers(); + GC.Collect(2, GCCollectionMode.Forced, blocking: true); + + long before = GC.GetTotalAllocatedBytes(precise: true); + + for (int i = 0; i < Pages; i++) + { + await CrawlPageAsync(client).ConfigureAwait(false); + } + + long allocated = GC.GetTotalAllocatedBytes(precise: true) - before; + await client.Disconnect().ConfigureAwait(false); + + double ratio = allocated / (double)Pages / PayloadBytes; + Console.WriteLine($"socket path: {allocated / (double)Pages / 1024 / 1024:F2} MB allocated per page ({ratio:F2}x)"); + + Assert.IsTrue( + ratio < 4.0, + $"the socket path allocated {ratio:F2}x the payload per page, budget is 4x " + + "(binding the string callback instead of the binary one costs about 2x more)"); + } + + /// + /// The string entry point is public and used by tests and consumers that feed messages in + /// by hand. A null there travelled down to the stream processor and was reported through + /// OnError; carrying the frame as bytes must not turn that into a throw out of the + /// method itself. + /// + [TestMethod] + public async Task TestNullMessageDoesNotThrowOutOfTheEntryPoint() + { + Connection connection = new Connection("ws://127.0.0.1:1/"); + + await connection.OnMessage(null).ConfigureAwait(false); + } + + private static async Task CrawlPageAsync(XrplClient client) + { + JsonElement page = await client + .GRequest(new LedgerDataRequest { Binary = true, Limit = 2048 }) + .ConfigureAwait(false); + + if (page.GetProperty("state").GetArrayLength() == 0) + { + throw new InvalidOperationException("ledger_data page carried no objects"); + } + } } } diff --git a/Xrpl/Client/connection.cs b/Xrpl/Client/connection.cs index 12ab7e03..fba830d4 100644 --- a/Xrpl/Client/connection.cs +++ b/Xrpl/Client/connection.cs @@ -1074,7 +1074,7 @@ await errorHandler.Invoke( } catch (Exception ex) { - Debug.WriteLine($"{DateTime.Now}OnMessageReceived callback error: {ex.Message}"); + Debug.WriteLine($"{DateTime.Now}OnBinaryMessage callback error: {ex.Message}"); } }); ws.OnDisconnect(async (closeStatus, closeDescription, closingSocket) => @@ -3175,6 +3175,12 @@ private Task IOnMessageFastPath(byte[] utf8Message) return IOnMessageFastPath(null, utf8Message); } + /// + /// Sent to in place of a message that could not be turned into text. + /// A literal, so reporting the failure needs no allocation of its own. + /// + private const string UnavailableMessageText = ""; + /// /// Exactly one of and carries the /// message; the other is null. @@ -3189,7 +3195,17 @@ private async Task IOnMessageFastPath(string message, byte[] utf8Message) { lastActivityTime = DateTime.UtcNow; - string Text() => message ??= Encoding.UTF8.GetString(utf8Message); + // Null in, null out: the string entry point is public, and a null message used to travel + // down to the stream processor and be reported through OnError rather than throw here. + string Text() + { + if (message is null && utf8Message is not null) + { + message = Encoding.UTF8.GetString(utf8Message); + } + + return message; + } // Scan message for "id" property to detect response messages var isResponse = utf8Message is null ? IsLikelyResponse(message) : IsLikelyResponse(utf8Message); @@ -3211,7 +3227,26 @@ private async Task IOnMessageFastPath(string message, byte[] utf8Message) catch (Exception error) { var errInfo = XrplErrorClassifier.Classify(error); - var capturedText = Text(); + if (OnError is null) + { + return; + } + + // The report has to survive whatever produced it. A response that fails to parse + // is most often a heap that has just run out, and a UTF-16 copy of the whole + // message is the largest allocation left on this path - if it cannot be had, the + // classification still goes out rather than the notification being lost to a + // second failure inside the handler. + string capturedText; + try + { + capturedText = Text(); + } + catch (OutOfMemoryException) + { + capturedText = UnavailableMessageText; + } + // Fire-and-forget for error callback - don't block _ = Task.Run(async () => { From 89fbc2c504037a816253ee3d1c7a96e750c5c591 Mon Sep 17 00:00:00 2001 From: Aleksandr Platonenkov Date: Sun, 16 Aug 2026 15:12:49 -0300 Subject: [PATCH 4/5] =?UTF-8?q?fix(client):=20=D0=BD=D0=B5=20=D1=81=D1=82?= =?UTF-8?q?=D1=80=D0=BE=D0=B8=D1=82=D1=8C=20=D1=82=D0=B5=D0=BA=D1=81=D1=82?= =?UTF-8?q?=20=D0=BF=D1=80=D0=B5=D0=B4=D1=83=D0=BF=D1=80=D0=B5=D0=B6=D0=B4?= =?UTF-8?q?=D0=B5=D0=BD=D0=B8=D0=B9=20=D0=B1=D0=B5=D0=B7=20=D0=BF=D0=BE?= =?UTF-8?q?=D0=B4=D0=BF=D0=B8=D1=81=D1=87=D0=B8=D0=BA=D0=BE=D0=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Обе находки CodeRabbit по PR #94. Диспетчеризация предупреждений материализовала UTF-16 копию ответа, как только в нём был warning, ещё до проверки, подписан ли хоть один из OnWarning и OnServerWarning. rippled вешает warning на ответы под нагрузкой и на reporting-сервере, то есть на таком узле это ровно та аллокация, которую PR убирает, — обратно на каждую страницу. Замер: предупреждения на всех 20 страницах, подписчиков нет — 4.28x -> 2.08x. Тест на null-сообщение проверял только отсутствие исключения и прошёл бы, если бы сообщение молча выбрасывалось. Теперь он подписывается на OnError и дожидается badMessage — то есть фиксирует именно тот маршрут, который правка восстанавливала. Путь предупреждений в тестах не был покрыт вообще, а условие я менял, поэтому добавлены обе стороны: TestWarningsStillReachTheirCallbacks (ответ с warning и warnings доходит до обоих колбэков) и TestUnsubscribedWarningsCostNothing (предупреждения без подписчиков не возвращают копию сообщения). У PagedResponseServer появился флаг withWarnings, по умолчанию выключенный. Пороги сокетных бюджетов подтянуты с 4x до 3x: 4.28x у регрессии стояло слишком близко к границе, теперь запас с обеих сторон — 2.08-2.19x как есть против 4.28-4.84x при откате. --- CHANGES.md | 1 + .../Xrpl.Tests/Client/PagedResponseServer.cs | 15 ++- .../Xrpl.Tests/Client/TestUResponseParsing.cs | 112 ++++++++++++++++-- Xrpl/Client/connection.cs | 11 +- 4 files changed, 128 insertions(+), 11 deletions(-) diff --git a/CHANGES.md b/CHANGES.md index 62082652..1f4e3417 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -5,6 +5,7 @@ * **Every response was parsed twice and copied to UTF-16 twice** — the cost of reading a response, measured rather than reasoned about. `RequestManager.Resolve` did `JsonSerializer.Deserialize(response.Result?.ToString() ?? "{}", taskInfo.Type, ...)`. `BaseResponse.Result` is typed `object`, which System.Text.Json fills with a `JsonElement` that already owns a private copy of the `result` bytes — so `.ToString()` rendered that element back into a UTF-16 string and the serializer parsed the string a second time. On a `ledger_data` page at `limit=2048` (~1 MB) the four stages measured, per response, at: 1.97 MB for the UTF-16 copy of the message, 1.68 MB for the document built over it, 1.97 MB for the UTF-16 copy of the `result`, 1.68 MB for the second document — **7.30 MB, 7.42x the response size**, all four allocations past the 85 KB large-object threshold. Both halves are now gone: * `DeserializeResult` works off the parsed node: `element.Deserialize(type, options)` for a typed model, and the element itself when the request asked for `JsonElement` or `object`, which is what a consumer that needs the raw ledger objects asks for (the typed `LOLedgerData.State` drops unknown fields). A `BaseResponse` assembled by hand rather than parsed off the wire keeps the old string path. Behaviour is otherwise unchanged, including a missing or JSON-`null` `result`, which still yields what deserializing `"{}"` yielded * the socket path carries the frame as it arrived. `Connection` binds `OnBinaryMessage` instead of `OnMessageReceived`, `IsLikelyResponse` and `RequestManager.HandleResponse` have `ReadOnlySpan` overloads, and the UTF-16 string is materialized — once, lazily — only for what genuinely needs text: stream messages and the `OnWarning`/`OnServerWarning`/`OnError` callbacks. The `string` overloads stay for `Connection.OnMessage(string)` and for external callers + * the warning callbacks no longer pay for listeners that are not there. rippled attaches `warning`/`warnings` to responses under load and on a reporting-mode server, and the dispatch built the UTF-16 text for them before checking whether `OnWarning`/`OnServerWarning` were subscribed — on such a server that is the removed allocation, back on every page. Measured with warnings on all 20 pages and nothing subscribed: 4.28x → 2.08x * the failure report survives the failure. A response that will not parse is most often a heap that has just run out, and materializing the message for `OnError` is then the largest allocation left on the path — if it throws, the notification is lost inside the handler and the consumer sees silence. The text is now built only when a handler is attached, and an `OutOfMemoryException` while building it falls back to a literal placeholder so the classification still goes out. `Connection.OnMessage(null)` also keeps its old route through `OnError` instead of throwing `ArgumentNullException` out of the entry point * measured end to end against a local WebSocket server, 600 `ledger_data` pages of ~1 MB: **8.32 → 2.68 MB allocated per response** (8.46x → 2.72x the payload), 11.49 → 7.96 ms per response, 87 → 126 responses/s, peak managed heap 42.4 → 20.8 MB, peak LOH 39.2 → 17.3 MB, peak working set 203.4 → 71.8 MB. Under a lowered `DOTNET_GCHeapHardLimit` the pre-fix path reproduced the production failure exactly — `XrplException: Failed to deserialize response for request : Exception of type 'System.OutOfMemoryException' was thrown`, with `JsonElement.ToString()` at the top of the inner stack — at a ceiling the fixed path completes 15/15 pages under * the win is not specific to `ledger_data` or to `JsonElement`: the second parse was on the path of every command. The repo's own `BenchmarkLedgerDataCrawl`, which goes through `Request` → `Dictionary`, drops from 22.9 to 14.8 MiB allocated per 2 MiB page (11.4x → 7.4x) with LOH ending at 38.2 instead of 115.4 MiB — it stays above the `JsonElement` figure because building a `Dictionary` boxes every value, which this change does not address diff --git a/Tests/Xrpl.Tests/Client/PagedResponseServer.cs b/Tests/Xrpl.Tests/Client/PagedResponseServer.cs index 8b5ff012..1d682772 100644 --- a/Tests/Xrpl.Tests/Client/PagedResponseServer.cs +++ b/Tests/Xrpl.Tests/Client/PagedResponseServer.cs @@ -20,14 +20,20 @@ internal sealed class PagedResponseServer : WebSocketTestServerBase private readonly int _fragments; private readonly string _resultBody; + private readonly bool _withWarnings; 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) + /// + /// Attach warning and warnings to every response, the way rippled does under + /// load and on a reporting-mode server. + /// + public PagedResponseServer(int approximatePayloadBytes, int fragments, bool withWarnings = false) { _fragments = Math.Max(1, fragments); _resultBody = BuildResultBody(approximatePayloadBytes); + _withWarnings = withWarnings; StartAccepting(); } @@ -126,7 +132,12 @@ private byte[] BuildFrame(string id, ref byte[]? reusable) private string Envelope(string id) { - return "{\"id\":" + id + ",\"status\":\"success\",\"type\":\"response\",\"result\":" + _resultBody + "}"; + string warnings = _withWarnings + ? ",\"warning\":\"load\",\"warnings\":[{\"id\":1001,\"message\":\"This is a reporting server.\"}]" + : string.Empty; + + return "{\"id\":" + id + ",\"status\":\"success\",\"type\":\"response\",\"result\":" + _resultBody + + warnings + "}"; } /// Pulls the JSON string value of the request's "id" property. diff --git a/Tests/Xrpl.Tests/Client/TestUResponseParsing.cs b/Tests/Xrpl.Tests/Client/TestUResponseParsing.cs index 9009781c..5352021f 100644 --- a/Tests/Xrpl.Tests/Client/TestUResponseParsing.cs +++ b/Tests/Xrpl.Tests/Client/TestUResponseParsing.cs @@ -258,23 +258,121 @@ public async Task TestSocketPathKeepsResponsesInTheirWireForm() Console.WriteLine($"socket path: {allocated / (double)Pages / 1024 / 1024:F2} MB allocated per page ({ratio:F2}x)"); Assert.IsTrue( - ratio < 4.0, - $"the socket path allocated {ratio:F2}x the payload per page, budget is 4x " + - "(binding the string callback instead of the binary one costs about 2x more)"); + ratio < 3.0, + $"the socket path allocated {ratio:F2}x the payload per page, budget is 3x " + + "(2.18x as bound, 4.84x with the string callback bound instead)"); + } + + /// + /// A response carrying warning/warnings still reaches both callbacks. The + /// text they are handed is now built only when one of them is subscribed, so this is the + /// side of that condition that must not have been broken. + /// + [TestMethod] + public async Task TestWarningsStillReachTheirCallbacks() + { + using PagedResponseServer server = new PagedResponseServer(64 * 1024, fragments: 1, withWarnings: true); + using XrplClient client = new XrplClient(server.Url); + + TaskCompletionSource warning = new TaskCompletionSource( + TaskCreationOptions.RunContinuationsAsynchronously); + TaskCompletionSource serverWarnings = new TaskCompletionSource( + TaskCreationOptions.RunContinuationsAsynchronously); + + await client.Connect().ConfigureAwait(false); + + client.connection.OnWarning += (text, message) => + { + warning.TrySetResult(text); + return Task.CompletedTask; + }; + + client.connection.OnServerWarning += (warnings, message) => + { + serverWarnings.TrySetResult(warnings.Count); + return Task.CompletedTask; + }; + + await CrawlPageAsync(client).ConfigureAwait(false); + + Task both = Task.WhenAll(warning.Task, serverWarnings.Task); + Task finished = await Task.WhenAny(both, Task.Delay(TimeSpan.FromSeconds(10))).ConfigureAwait(false); + await client.Disconnect().ConfigureAwait(false); + + Assert.AreSame(both, finished, "a warned response did not reach OnWarning/OnServerWarning"); + Assert.AreEqual("load", await warning.Task.ConfigureAwait(false)); + Assert.AreEqual(1, await serverWarnings.Task.ConfigureAwait(false)); + } + + /// + /// And the other side of it: warnings on every page with nothing subscribed must not put + /// the UTF-16 copy of each response back on the path. + /// + /// Reads the process-wide counter, so it stays out of the parallel pass. + [TestMethod] + [DoNotParallelize] + public async Task TestUnsubscribedWarningsCostNothing() + { + const int Pages = 20; + const int PayloadBytes = 1024 * 1024; + + using PagedResponseServer server = new PagedResponseServer(PayloadBytes, fragments: 8, withWarnings: true); + using XrplClient client = new XrplClient(server.Url); + + await client.Connect().ConfigureAwait(false); + await CrawlPageAsync(client).ConfigureAwait(false); + + GC.Collect(2, GCCollectionMode.Forced, blocking: true); + GC.WaitForPendingFinalizers(); + GC.Collect(2, GCCollectionMode.Forced, blocking: true); + + long before = GC.GetTotalAllocatedBytes(precise: true); + + for (int i = 0; i < Pages; i++) + { + await CrawlPageAsync(client).ConfigureAwait(false); + } + + long allocated = GC.GetTotalAllocatedBytes(precise: true) - before; + await client.Disconnect().ConfigureAwait(false); + + double ratio = allocated / (double)Pages / PayloadBytes; + Console.WriteLine($"warned pages, no subscribers: {allocated / (double)Pages / 1024 / 1024:F2} MB per page ({ratio:F2}x)"); + + Assert.IsTrue( + ratio < 3.0, + $"warned responses allocated {ratio:F2}x the payload per page with nothing subscribed, " + + "budget is 3x (2.08x as bound, 4.28x when the text is built regardless of subscribers)"); } /// /// The string entry point is public and used by tests and consumers that feed messages in - /// by hand. A null there travelled down to the stream processor and was reported through - /// OnError; carrying the frame as bytes must not turn that into a throw out of the - /// method itself. + /// by hand. A null there travelled down to the stream processor and came back out through + /// OnError as a badMessage; carrying the frame as bytes must not turn that + /// into a throw out of the method itself, and must not turn it into silence either. /// [TestMethod] - public async Task TestNullMessageDoesNotThrowOutOfTheEntryPoint() + public async Task TestNullMessageIsStillReportedThroughOnError() { Connection connection = new Connection("ws://127.0.0.1:1/"); + TaskCompletionSource reported = new TaskCompletionSource( + TaskCreationOptions.RunContinuationsAsynchronously); + + connection.OnError += (error, errorMessage, message, data) => + { + reported.TrySetResult(errorMessage); + return Task.CompletedTask; + }; + // Must not throw: the entry point is public and a null used to be routed, not raised. await connection.OnMessage(null).ConfigureAwait(false); + + // The routing itself is fire-and-forget, so the report arrives after the call returns. + Task completed = await Task.WhenAny(reported.Task, Task.Delay(TimeSpan.FromSeconds(5))) + .ConfigureAwait(false); + + Assert.AreSame(reported.Task, completed, "a null message was dropped instead of being reported"); + Assert.AreEqual("badMessage", await reported.Task.ConfigureAwait(false)); } private static async Task CrawlPageAsync(XrplClient client) diff --git a/Xrpl/Client/connection.cs b/Xrpl/Client/connection.cs index fba830d4..e330e7fb 100644 --- a/Xrpl/Client/connection.cs +++ b/Xrpl/Client/connection.cs @@ -3267,8 +3267,15 @@ string Text() } // THEN: Handle warnings and errors in background (fire-and-forget) - // These are informational and should not delay response processing - if (data.Warning != null || data.Warnings is { Count: > 0 }) + // These are informational and should not delay response processing. + // Materialize the text only when something is actually listening: rippled attaches a + // warning to every response under load and on a reporting-mode server, so building a + // UTF-16 copy for a callback nobody registered would put back, page after page, + // exactly the allocation this path exists to avoid. + bool warningNeedsText = (data.Warning != null && OnWarning is not null) + || (data.Warnings is { Count: > 0 } && OnServerWarning is not null); + + if (warningNeedsText) { var capturedData = data; var capturedMessage = Text(); From 788693929d0671d7d7a8d2a06cfa594bbaf0ba58 Mon Sep 17 00:00:00 2001 From: Aleksandr Platonenkov Date: Sun, 16 Aug 2026 15:23:23 -0300 Subject: [PATCH 5/5] =?UTF-8?q?test(client):=20=D0=BF=D1=80=D0=BE=D0=B2?= =?UTF-8?q?=D0=B5=D1=80=D1=8F=D1=82=D1=8C=20=D0=B8=20=D1=82=D0=B5=D0=BA?= =?UTF-8?q?=D1=81=D1=82,=20=D0=BA=D0=BE=D1=82=D0=BE=D1=80=D1=8B=D0=B9=20?= =?UTF-8?q?=D0=BF=D0=BE=D0=BB=D1=83=D1=87=D0=B0=D1=8E=D1=82=20=D0=BA=D0=BE?= =?UTF-8?q?=D0=BB=D0=B1=D1=8D=D0=BA=D0=B8=20=D0=BF=D1=80=D0=B5=D0=B4=D1=83?= =?UTF-8?q?=D0=BF=D1=80=D0=B5=D0=B6=D0=B4=D0=B5=D0=BD=D0=B8=D0=B9?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit По итогам селф-ревью коммита 89fbc2c. TestWarningsStillReachTheirCallbacks сверял значения предупреждений, но игнорировал аргумент message — то есть прошёл бы, если бы колбэкам передавали null или заглушку UnavailableMessageText. А решение о том, строить ли этот текст, как раз и принимает условие, которое тест охраняет. Теперь оба лямбда-обработчика захватывают message, и он проверяется на непустоту и на содержимое ответа. --- .../Xrpl.Tests/Client/TestUResponseParsing.cs | 29 ++++++++++++++----- 1 file changed, 21 insertions(+), 8 deletions(-) diff --git a/Tests/Xrpl.Tests/Client/TestUResponseParsing.cs b/Tests/Xrpl.Tests/Client/TestUResponseParsing.cs index 5352021f..ff3c5742 100644 --- a/Tests/Xrpl.Tests/Client/TestUResponseParsing.cs +++ b/Tests/Xrpl.Tests/Client/TestUResponseParsing.cs @@ -274,22 +274,22 @@ public async Task TestWarningsStillReachTheirCallbacks() using PagedResponseServer server = new PagedResponseServer(64 * 1024, fragments: 1, withWarnings: true); using XrplClient client = new XrplClient(server.Url); - TaskCompletionSource warning = new TaskCompletionSource( - TaskCreationOptions.RunContinuationsAsynchronously); - TaskCompletionSource serverWarnings = new TaskCompletionSource( - TaskCreationOptions.RunContinuationsAsynchronously); + TaskCompletionSource<(string Warning, string Message)> warning = + new TaskCompletionSource<(string, string)>(TaskCreationOptions.RunContinuationsAsynchronously); + TaskCompletionSource<(int Count, string Message)> serverWarnings = + new TaskCompletionSource<(int, string)>(TaskCreationOptions.RunContinuationsAsynchronously); await client.Connect().ConfigureAwait(false); client.connection.OnWarning += (text, message) => { - warning.TrySetResult(text); + warning.TrySetResult((text, message)); return Task.CompletedTask; }; client.connection.OnServerWarning += (warnings, message) => { - serverWarnings.TrySetResult(warnings.Count); + serverWarnings.TrySetResult((warnings.Count, message)); return Task.CompletedTask; }; @@ -300,8 +300,21 @@ public async Task TestWarningsStillReachTheirCallbacks() await client.Disconnect().ConfigureAwait(false); Assert.AreSame(both, finished, "a warned response did not reach OnWarning/OnServerWarning"); - Assert.AreEqual("load", await warning.Task.ConfigureAwait(false)); - Assert.AreEqual(1, await serverWarnings.Task.ConfigureAwait(false)); + + (string Warning, string Message) warned = await warning.Task.ConfigureAwait(false); + (int Count, string Message) served = await serverWarnings.Task.ConfigureAwait(false); + + Assert.AreEqual("load", warned.Warning); + Assert.AreEqual(1, served.Count); + + // The message the callbacks are handed is the point of the condition guarding it: they + // must get the response text, not null and not the out-of-memory placeholder. + foreach (string text in new[] { warned.Message, served.Message }) + { + Assert.IsNotNull(text, "the warning callbacks were handed no message"); + StringAssert.Contains(text, "\"warning\":\"load\""); + StringAssert.Contains(text, "\"state\":["); + } } ///