From 7c6ee29bc1a78c3ffcd638d23e33a96e48c0fbb6 Mon Sep 17 00:00:00 2001
From: Aleksandr <44946855+Platonenkov@users.noreply.github.com>
Date: Sun, 6 Sep 2026 22:12:19 +0000
Subject: [PATCH] fix(connection): retire the session before sweeping the
pending requests (#178)
* fix(connection): retire the session before sweeping the pending requests
A request created while the client was retiring its connection was written
to the socket being retired and hung until RequestTimeout (issue #177).
Every retirement path swept the pending requests with
RejectAllWithCancellation() and only afterwards cleared `ws`. The sweep
resumes consumer continuations - inline on the sweeping thread when there
is no synchronization context, since RequestManager's completion sources
are not created with RunContinuationsAsynchronously - and a consumer that
issued its next request from there passed ShouldBeConnected() on the old
socket, was sent into it after the sweep, and was never completed. Moving
the clear ahead of the first await, as the issue suggested, does not close
that window: the continuation runs inside the sweep itself.
Clear `ws` and mark the session retiring before the sweep on all five
paths: ChangeServer, the ping/network fast reconnect, Disconnect,
DisconnectAndWaitAsync, and the failed-OnConnected-handler path, which the
issue did not list. `ws` is now volatile - it is the gate every request
reads lock-free. A DisconnectedException out of the send no longer leaves
the just-created request pending for RequestTimeout; it is rejected like an
encoding failure and surfaces through the same await.
The regression test pins the ordering with an ExecuteSynchronously
continuation, which is the shape the sweep hands a consumer on a thread
pool: it asserts the socket the follow-up saw is not the retired one, and
that it completes within seconds rather than after RequestTimeout.
* fix(connection): let the retirement paths yield instead of block, and settle a fast reconnect once
Three defects seen on the Blazor WebAssembly stand while verifying #177,
all on the paths that retire a session.
Stopping the message processor blocked the calling thread on the reader
task for up to two seconds. On a single-threaded host the reader's
continuation needs that very thread to observe the completed channel, so
the block never returned early: every ChangeServer, fast reconnect and
Disconnect stalled the UI for the full two seconds. The stop is now
detach-under-lock plus an awaited exit with the same cap, taken after the
request sweep so consumers are released first. Measured on the stand:
2000 ms to 5-164 ms per switch.
RetireCurrentSessionAndReconnectAsync is awaited from inside the ping
check, and only from there, yet it waited for the ping to finish - its
own - and so ran out WaitForPingToFinishAsync's three-second timeout on
every ping-triggered reconnect. An AsyncLocal set by the ping check lets
the wait recognise the ping it runs in. RestoringConnection to the
session-ended notification: 6 s to 20 ms.
When the fast reconnect's own attempt failed at the socket, the failure
callback started the reconnect loop on the same cancellation source; the
loop connected first, OnceOpen retired the source, and the fast
reconnect's wait came back cancelled. Its catch read that as a failure:
it reported RestoringConnection on a connected client and started a
second loop, whose first attempt retired the live socket and opened
another. A connected client is now recognised as settled in that catch,
and the loop itself no longer retires a socket that is open when its
turn comes.
Both timing defects and the double reconnect are pinned by tests on a
silent-on-ping server: the first measures the retirement gap, the second
takes the server down at RestoringConnection so the sequence falls to
the loop, brings a replacement up on the same port, and requires exactly
one connection afterwards.
* chore(release): Xrpl 11.3.2.0
Patch: the connection fixes for #177 and the three defects found while
verifying it change no contract. Only the Xrpl package moves - the base
packages are untouched since 11.3.1.0 and keep their published versions.
* fix(connection): scope the ping-check marker to its connection, and let the switch tests fail as themselves
Review follow-ups on #178.
The AsyncLocal that lets WaitForPingToFinishAsync recognise the ping it
runs in was static. Its value follows the execution context, so a
consumer's OnPing handler that awaited another connection would carry
the flag into that connection and let it skip waiting for its own ping.
It is an instance field now.
The regression tests judged the follow-up request inside a try whose
catch-all also caught the assertion failures, replacing the original
message with an exception-type mismatch. The outcome is captured first
and judged afterwards. The IsNotNull on the response went with it:
XrplResponse is a struct, and completion within the bound is all a
success has to show.
* fix(connection): a socket a user Disconnect() took cannot come back, and a stale handler failure keeps its hands off the live connection
Review follow-ups on #178, all on code adjacent to the retirement paths.
OnceOpen checked only the session. Disconnect() does not retire the
session - OnceClose is what announces a user disconnect - so a handshake
that completed while Disconnect() was running had its OnConnect callback
install the socket, clear the intentional-disconnect tracking, and hand
the close that followed to the reconnect loop as a network drop. The
socket's own marks are the signal now: a socket Disconnect() took, or a
client that is permanently disconnected, is refused.
OnConnectHandlerFailedAsync swept the pending requests, stopped the
message processor and the ping timer before asking whether the failed
socket was still the current one. When a newer connection had replaced
it, that cleanup hit the newer connection. The ownership check comes
first; a stale callback closes its own socket and steps aside.
DisconnectAndWaitAsync returned "already disconnected" at once when a
concurrent call had taken the socket, while that call was still closing
it. It now waits on the disconnect in progress, within its own timeout.
* fix(connection): a stale OnConnected failure is not the current connection's failure
Review follow-up on #178. The ownership check in OnConnectHandlerFailedAsync
sat after the give-up branch, so a delayed failure from a socket that a
newer connection had already replaced still counted towards giving up
and, with StopAfterMaxAttempts on, ran RejectAll and Disconnect() against
the live connection. The check comes first now, read-only: a stale
callback closes its own socket and returns before anything is counted
or torn down. The clearing check under the same lock stays where it was,
for a replacement that lands during the consumer notifications in between.
* fix(connection): the give-up teardown rechecks ownership after notifying
Review follow-up on #178. In OnConnectHandlerFailedAsync the give-up
branch notifies the consumer before it rejects the requests and calls
Disconnect() - deliberately, so the reason reaches the consumer before
Disconnect() moves the state itself. That notification runs consumer
code: a handler that answers "gave up" with a ChangeServer has already
taken the failed socket out of ws and is opening another, and the
teardown would have hit that connection. Ownership is rechecked after
the notification; the three checks in the method now share one helper.
* fix(connection): the handler-failure notification goes out before the socket is taken
Review follow-up on #178. Moving the RestoringConnection notification
behind the clearing ownership check, as the previous round did, left the
consumer's handler running after this path had already taken the socket
out of ws: a ChangeServer from that handler installed a replacement, and
the cleanup that followed - ping timer, sweep, processor, reconnect
state - hit it, with no way to tell from ws alone. The notification now
precedes the clear, so the ownership check that comes with the clear is
the one that sees what the handler did: a ChangeServer retires the socket
itself, the check finds it gone, and this path closes its socket and
steps aside.
---
CHANGES.md | 15 +
.../Client/TestUFastReconnectSettling.cs | 271 ++++++++++++
.../Client/TestURequestDuringServerSwitch.cs | 322 +++++++++++++++
Xrpl/Client/connection.cs | 386 +++++++++++++-----
Xrpl/Xrpl.csproj | 2 +-
5 files changed, 886 insertions(+), 110 deletions(-)
create mode 100644 Tests/Xrpl.Tests/Client/TestUFastReconnectSettling.cs
create mode 100644 Tests/Xrpl.Tests/Client/TestURequestDuringServerSwitch.cs
diff --git a/CHANGES.md b/CHANGES.md
index f703f231..e9009d48 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -1,5 +1,20 @@
# Changes
+## 11.3.2.0 06/09/2026
+
+* **A request issued while the client is switching servers no longer hangs until `RequestTimeout`** (#177). Every path that retires a connection - `ChangeServer`, the ping-triggered fast reconnect, `Disconnect`, `DisconnectAndWaitAsync`, and the path taken when an `OnConnected` handler fails - rejected the pending requests first and cleared the socket reference afterwards. The rejection resumes the consumer, and a consumer that issues its next request from there - the second value of a page load, read from the response handler of the first - found the retired socket still installed, passed the connectivity check on it, and was written into it after the sweep that would have rejected it. Nothing completed it: the sweep had run, and a failed send is report-only. Forty seconds later it timed out, with the connection healthy for thirty-nine of them.
+ * the fix clears the socket reference before the sweep. Moving it ahead of the first `await`, as the issue proposed, is not enough: `RequestManager` builds its completion sources without `RunContinuationsAsynchronously`, so on a thread pool the consumer's continuation runs inline, inside the sweep itself, before any await. A fifth retirement path the issue did not list is covered too
+ * `ImmediateFail` now refuses such a request at once with `NotConnectedException`; `WaitForConnection` carries it over to the new connection. A request whose send fails because the connection went away between the check and the send is rejected rather than left pending for `RequestTimeout`
+ * for consumers: retry logic that recognised this failure by the `TimeoutException` it used to produce now sees `NotConnectedException` (or `OperationCanceledException`, for a request that was in flight when the switch began) immediately. Classify by type rather than by message
+ * the regression test hands the sweep a continuation that runs synchronously and asserts which socket the follow-up saw
+
+* **`ChangeServer`, the fast reconnect and `Disconnect` no longer stall a single-threaded host for two seconds.** Stopping the stream-message processor blocked the calling thread on its reader task, with a two-second cap. On Blazor WebAssembly the reader's continuation needs the very thread that was blocked, so the cap was always reached: every server switch froze the UI for two seconds, and a wake-from-background reconnect is such a switch. The stop is awaited now, after the request sweep, and the reader is gone in milliseconds. Measured on the WebAssembly test client: 2000 ms to 5-390 ms per switch.
+
+* **A ping-triggered reconnect no longer waits three seconds for itself.** `RetireCurrentSessionAndReconnectAsync` runs inside the ping check that calls it, and waited for the ping to finish before retiring the session - its own, whose flag could not clear until it returned. Every reconnect the health check started paid the full `WaitForPingToFinishAsync` timeout before announcing that the session had ended. The wait now recognises the ping it runs in. `RestoringConnection` to `OnSessionEnded` on the stand: 6 s to 20 ms.
+
+* **A reconnect the loop finished is not reconnected a second time.** When the fast reconnect's own attempt failed at the socket, the failure callback started the reconnect loop on the same cancellation source; the loop connected first, `OnceOpen` retired the source, and the fast reconnect's wait came back cancelled. Its catch read that as a failure: it reported `RestoringConnection` on a client that was connected and started a second loop, whose first attempt retired the live socket and opened another. Consumers saw two `OnConnected` per recovery, with a spurious `RestoringConnection` between them, and restored their subscriptions twice. A connected client is now recognised as settled, and the loop no longer retires a socket that is open when its turn comes.
+ * pinned by a test that takes the server down at `RestoringConnection`, so the sequence falls to the loop, brings a replacement up on the same port and requires exactly one connection afterwards
+
## 11.3.1.0 05/09/2026
* **`FundWallet` no longer reports success for a wallet the faucet never paid** (#174). The starting balance was read before the faucet was asked, and any failure to read it left it at zero; the wait then asked whether the balance had risen above that zero, so a wallet that already held funds satisfied the test and the call returned `Funded` with the balance the account had all along. Found three times independently while cold-reviewing the previous release, by three reviewers on two models.
diff --git a/Tests/Xrpl.Tests/Client/TestUFastReconnectSettling.cs b/Tests/Xrpl.Tests/Client/TestUFastReconnectSettling.cs
new file mode 100644
index 00000000..e8fefa72
--- /dev/null
+++ b/Tests/Xrpl.Tests/Client/TestUFastReconnectSettling.cs
@@ -0,0 +1,271 @@
+using Microsoft.VisualStudio.TestTools.UnitTesting;
+
+using System;
+using System.Collections.Generic;
+using System.Diagnostics;
+using System.Threading.Tasks;
+
+using Xrpl.Client;
+
+namespace Xrpl.Tests
+{
+ ///
+ /// The fast-reconnect path (RetireCurrentSessionAndReconnectAsync) runs inside the ping
+ /// check that triggers it. Two things followed from that and went unnoticed because the path
+ /// only shows its timing on a live node.
+ ///
+ ///
+ ///
+ /// It waited for the ping to finish - its own ping - and so waited out the whole
+ /// WaitForPingToFinishAsync timeout (3 s) on every ping-triggered reconnect.
+ ///
+ ///
+ /// And when its own connection attempt failed at the socket, the failure callback started the
+ /// reconnect loop on the same cancellation source; the loop connected first, OnceOpen
+ /// retired the source, the fast reconnect's wait came back cancelled, and its catch read that
+ /// as a failure: it started a second loop, whose first attempt retired the live socket and
+ /// opened another. One reconnect became two, with RestoringConnection reported on a
+ /// client that was connected.
+ ///
+ ///
+ [TestClass]
+ public class TestUFastReconnectSettling
+ {
+ private XrplClient _client;
+
+ private static Dictionary ServerInfoResponse() => new Dictionary
+ {
+ { "type", "response" },
+ { "status", "success" },
+ { "result", new Dictionary
+ {
+ { "info", new Dictionary
+ {
+ { "build_version", "test-mock" },
+ { "complete_ledgers", "1-1" },
+ { "server_state", "full" },
+ }
+ },
+ }
+ },
+ };
+
+ private static Dictionary EmptyResponse() => new Dictionary
+ {
+ { "type", "response" },
+ { "status", "success" },
+ { "result", new Dictionary() },
+ };
+
+ ///
+ /// The two knobs that make the inactivity path reachable in under a second, plus a reconnect
+ /// backoff short enough for a whole failed-then-succeeded sequence to fit in a test.
+ ///
+ private static XrplClient.ClientOptions FastReconnectOptions() => new XrplClient.ClientOptions
+ {
+ RequestPolicy = RequestFailurePolicy.ImmediateFail,
+ ReconnectBaseDelay = TimeSpan.FromMilliseconds(100),
+ ReconnectMaxDelay = TimeSpan.FromMilliseconds(500),
+ ConnectionAcquisitionTimeout = TimeSpan.FromSeconds(5),
+ ConnectionAttemptTimeout = TimeSpan.FromSeconds(3),
+ UseCustomPing = true,
+ HealthCheckInterval = TimeSpan.FromMilliseconds(200),
+ InactivityTimeout = TimeSpan.FromMilliseconds(500),
+ };
+
+ [TestCleanup]
+ public async Task MyTestCleanup()
+ {
+ if (_client != null)
+ {
+ try
+ {
+ await _client.Disconnect();
+ }
+ catch
+ {
+ // Cleanup is not an assertion.
+ }
+
+ _client = null;
+ }
+ }
+
+ private static async Task WaitUntilAsync(Func condition, TimeSpan timeout, string what)
+ {
+ Stopwatch clock = Stopwatch.StartNew();
+ while (!condition())
+ {
+ Assert.IsTrue(clock.Elapsed < timeout, $"Timed out after {timeout.TotalSeconds:F0}s waiting for: {what}");
+ await Task.Delay(50);
+ }
+ }
+
+ ///
+ /// From the moment the health check hands a silent connection to the fast-reconnect path
+ /// to the moment the old session is announced as ended, nothing has to wait for anything:
+ /// the requests are swept, the socket is retired in the background. Three seconds in that
+ /// gap is WaitForPingToFinishAsync timing out on the ping this path runs inside.
+ ///
+ [TestMethod]
+ public async Task TestFastReconnectDoesNotWaitOutItsOwnPingTimeout()
+ {
+ using SilentOnPingServer server = new SilentOnPingServer();
+ _client = new XrplClient(server.Url, FastReconnectOptions());
+
+ Stopwatch clock = Stopwatch.StartNew();
+ long restoringAt = -1;
+ long sessionEndedAt = -1;
+ object gate = new object();
+
+ _client.OnConnectionStatus += info =>
+ {
+ if (info.ConnectionState != XrpConnectionState.RestoringConnection)
+ {
+ return;
+ }
+
+ lock (gate)
+ {
+ if (restoringAt < 0)
+ {
+ restoringAt = clock.ElapsedMilliseconds;
+ }
+ }
+ };
+
+ _client.OnSessionEnded += (reason, description) =>
+ {
+ lock (gate)
+ {
+ if (sessionEndedAt < 0)
+ {
+ sessionEndedAt = clock.ElapsedMilliseconds;
+ }
+ }
+
+ return Task.CompletedTask;
+ };
+
+ await _client.Connect();
+ Assert.IsTrue(_client.connection.IsConnected(), "Precondition: the client must be connected.");
+
+ await WaitUntilAsync(
+ () => { lock (gate) { return sessionEndedAt >= 0; } },
+ TimeSpan.FromSeconds(15),
+ "the silent connection to be retired by the fast-reconnect path");
+
+ long gap;
+ lock (gate)
+ {
+ Assert.IsTrue(restoringAt >= 0, "RestoringConnection must have been reported before the session ended.");
+ gap = sessionEndedAt - restoringAt;
+ }
+
+ Assert.IsTrue(
+ gap < 2000,
+ $"Retiring the session took {gap}ms after RestoringConnection was reported. The fast-reconnect " +
+ "path waited out WaitForPingToFinishAsync's timeout for the ping check it is itself running in.");
+ }
+
+ ///
+ /// A fast reconnect whose own attempt fails at the socket hands the sequence to the reconnect
+ /// loop on the same source. When that loop connects, the fast reconnect is done - it must not
+ /// read the cancellation of its wait as a failure and start reconnecting a connected client.
+ ///
+ [TestMethod]
+ public async Task TestReconnectLoopSettlingAFastReconnectDoesNotReconnectAgain()
+ {
+ SilentOnPingServer silentServer = new SilentOnPingServer();
+ CreateMockRippled replacement = null;
+ int port = silentServer.Port;
+
+ try
+ {
+ _client = new XrplClient(silentServer.Url, FastReconnectOptions());
+
+ TaskCompletionSource loopIsDelaying = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
+ object gate = new object();
+ bool retired = false;
+ int connectedAfterRetire = 0;
+ WebSocketClient socketAfterRetire = null;
+ List restoringAfterConnected = new List();
+
+ _client.OnConnectionStatus += info =>
+ {
+ lock (gate)
+ {
+ switch (info.ConnectionState)
+ {
+ case XrpConnectionState.RestoringConnection when !retired:
+ // The health check just handed the silent connection to the
+ // fast-reconnect path, and this handler runs before its attempt
+ // starts. Take the server down now, so that attempt fails at the
+ // socket and the sequence falls to the reconnect loop.
+ retired = true;
+ silentServer.Dispose();
+ break;
+
+ case XrpConnectionState.RestoringConnection when connectedAfterRetire > 0:
+ restoringAfterConnected.Add(info.Message);
+ break;
+
+ case XrpConnectionState.RestoringConnection
+ when info.Message.StartsWith("Reconnecting in", StringComparison.Ordinal):
+ // The loop's first, immediate attempt failed too; it is now waiting
+ // out a backoff, which is the window to bring a server up in.
+ loopIsDelaying.TrySetResult(true);
+ break;
+
+ case XrpConnectionState.Connected when retired:
+ connectedAfterRetire++;
+ socketAfterRetire ??= _client.connection.ws;
+ break;
+ }
+ }
+ };
+
+ await _client.Connect();
+ Assert.IsTrue(_client.connection.IsConnected(), "Precondition: the client must be connected.");
+
+ Task finished = await Task.WhenAny(loopIsDelaying.Task, Task.Delay(TimeSpan.FromSeconds(15)));
+ Assert.AreSame(loopIsDelaying.Task, finished, "The reconnect loop never reached a delayed attempt.");
+
+ replacement = new CreateMockRippled(port) { suppressOutput = true };
+ replacement.AddResponse("server_info", ServerInfoResponse());
+ replacement.AddResponse("ping", EmptyResponse());
+ replacement.Start();
+
+ await WaitUntilAsync(
+ () => { lock (gate) { return connectedAfterRetire > 0; } },
+ TimeSpan.FromSeconds(15),
+ "the reconnect loop to connect to the replacement server");
+
+ // Long enough for a second loop's first attempt (CalcBackoff(1) = 2 x base delay)
+ // to have retired the socket and reconnected, had one been started.
+ await Task.Delay(TimeSpan.FromSeconds(3));
+
+ lock (gate)
+ {
+ Assert.AreEqual(
+ 0,
+ restoringAfterConnected.Count,
+ "RestoringConnection was reported on a connected client: " + string.Join(" | ", restoringAfterConnected));
+ Assert.AreEqual(
+ 1,
+ connectedAfterRetire,
+ "The client connected more than once: the fast reconnect started a second loop after the first one had already connected.");
+ Assert.AreSame(
+ socketAfterRetire,
+ _client.connection.ws,
+ "The socket the loop opened was retired and replaced by a reconnect nobody needed.");
+ }
+ }
+ finally
+ {
+ replacement?.Stop();
+ silentServer.Dispose();
+ }
+ }
+ }
+}
diff --git a/Tests/Xrpl.Tests/Client/TestURequestDuringServerSwitch.cs b/Tests/Xrpl.Tests/Client/TestURequestDuringServerSwitch.cs
new file mode 100644
index 00000000..e157701f
--- /dev/null
+++ b/Tests/Xrpl.Tests/Client/TestURequestDuringServerSwitch.cs
@@ -0,0 +1,322 @@
+using Microsoft.VisualStudio.TestTools.UnitTesting;
+
+using System;
+using System.Collections.Generic;
+
+using System.Threading.Tasks;
+
+using Xrpl.Client;
+
+using TimeoutException = Xrpl.Client.Exceptions.TimeoutException;
+
+namespace Xrpl.Tests
+{
+ ///
+ /// Regression tests for issue #177 - a request created while the client is retiring its
+ /// connection is written to the socket that is being retired and nothing ever completes it.
+ ///
+ ///
+ ///
+ /// Every retirement path (ChangeServer, the ping/network fast reconnect,
+ /// Disconnect, DisconnectAndWaitAsync, the failed-OnConnected-handler path) used
+ /// to sweep the pending requests with RejectAllWithCancellation() and only afterwards
+ /// clear ws. Between those two points ShouldBeConnected() still reported the
+ /// retired socket as usable, so a request issued there passed the connectivity check and was
+ /// sent into a socket that was on its way out. The sweep had already run, so nothing rejected
+ /// it; the send is async void and report-only, so a failure did not reject it either.
+ /// The caller waited out the whole RequestTimeout. The fix clears ws before the
+ /// sweep on every one of those paths.
+ ///
+ ///
+ /// The window is not a thread race that has to be won by luck. RequestManager builds its
+ /// without
+ /// , so on a thread pool
+ /// context the consumer's continuation runs inline, inside the sweep itself - which is
+ /// what these tests exercise. Under a single-threaded synchronization context (Blazor
+ /// WebAssembly, where this was observed) the continuation is posted instead and lands on the
+ /// first real yield of the retirement path, still ahead of the ws clear whenever a ping
+ /// is in flight.
+ ///
+ ///
+ [TestClass]
+ public class TestURequestDuringServerSwitch
+ {
+ ///
+ /// Deliberately short so a request that falls into the window fails the test in seconds
+ /// instead of the 40s production default - the assertions below are about the request not
+ /// waiting this out at all.
+ ///
+ private static readonly TimeSpan RequestTimeout = TimeSpan.FromSeconds(8);
+
+ ///
+ /// What a correctly handled request may take: it either fails fast or is carried over to
+ /// the new connection. Comfortably below so the two outcomes
+ /// cannot be confused.
+ ///
+ private static readonly TimeSpan AcceptableBound = TimeSpan.FromSeconds(4);
+
+ private CreateMockRippled _firstRippled;
+ private CreateMockRippled _secondRippled;
+ private XrplClient _client;
+ private int _firstPort;
+ private int _secondPort;
+
+ private static Dictionary ServerInfoResponse() => new Dictionary
+ {
+ { "type", "response" },
+ { "status", "success" },
+ { "result", new Dictionary
+ {
+ { "info", new Dictionary
+ {
+ { "build_version", "test-mock" },
+ { "complete_ledgers", "1-1" },
+ { "server_state", "full" },
+ }
+ },
+ }
+ },
+ };
+
+ private static Dictionary AccountInfoResponse() => new Dictionary
+ {
+ { "type", "response" },
+ { "status", "success" },
+ { "result", new Dictionary
+ {
+ { "account_data", new Dictionary
+ {
+ { "Account", "rTestAccountForIssue177000000000000" },
+ { "Balance", "1000000" },
+ { "Sequence", 1 },
+ }
+ },
+ }
+ },
+ };
+
+ private static Dictionary AccountInfoRequest() => new Dictionary
+ {
+ { "command", "account_info" },
+ { "account", "rTestAccountForIssue177000000000000" },
+ };
+
+ [TestInitialize]
+ public void MyTestInitialize()
+ {
+ _firstPort = TestUtils.GetFreePort();
+ _secondPort = TestUtils.GetFreePort();
+
+ // The first node answers server_info but sits on account_info far longer than the test
+ // runs: that is how a request is kept pending so the retirement sweep has something to
+ // reject, and the rejection is what hands control back to the consumer.
+ _firstRippled = new CreateMockRippled(_firstPort) { suppressOutput = true };
+ _firstRippled.AddResponse("server_info", ServerInfoResponse());
+ _firstRippled.AddDelayedResponse("account_info", AccountInfoResponse(), TimeSpan.FromMinutes(5));
+ _firstRippled.Start();
+
+ // The second node answers everything at once - a request carried over to it must come
+ // back quickly.
+ _secondRippled = new CreateMockRippled(_secondPort) { suppressOutput = true };
+ _secondRippled.AddResponse("server_info", ServerInfoResponse());
+ _secondRippled.AddResponse("account_info", AccountInfoResponse());
+ _secondRippled.Start();
+ }
+
+ [TestCleanup]
+ public async Task MyTestCleanup()
+ {
+ if (_client != null)
+ {
+ try
+ {
+ await _client.Disconnect();
+ }
+ catch
+ {
+ // The test may have left the client mid-switch; cleanup is not an assertion.
+ }
+
+ _client = null;
+ }
+
+ _firstRippled?.Stop();
+ _secondRippled?.Stop();
+ }
+
+ private XrplClient CreateClient(int port) => new XrplClient(
+ $"ws://127.0.0.1:{port}",
+ new XrplClient.ClientOptions
+ {
+ RequestTimeout = RequestTimeout,
+ RequestPolicy = RequestFailurePolicy.WaitForConnection,
+ ConnectionAcquisitionTimeout = TimeSpan.FromSeconds(5),
+ ConnectionAttemptTimeout = TimeSpan.FromSeconds(5),
+ ReconnectBaseDelay = TimeSpan.FromMilliseconds(100),
+ ReconnectMaxDelay = TimeSpan.FromSeconds(1),
+ UseCustomPing = false,
+ UseCheckHealth = false,
+ });
+
+ ///
+ /// Starts a request that the first node will not answer and returns once it is actually
+ /// pending in the request manager, so the retirement sweep is guaranteed to find it.
+ ///
+ private async Task>>> StartPendingRequestAsync()
+ {
+ Task>> pending =
+ _client.connection.Request(AccountInfoRequest());
+
+ // Give the send a moment to leave; nothing observable marks "in flight", and the
+ // request only has to exist as a promise for the sweep to reach it.
+ await Task.Delay(300);
+
+ Assert.IsFalse(
+ pending.IsCompleted,
+ "Precondition: the first node must leave this request unanswered.");
+
+ return pending;
+ }
+
+ ///
+ /// A request issued from the rejection continuation of a swept request - the position the
+ /// issue describes - must not be written into the socket ChangeServer is retiring.
+ ///
+ [TestMethod]
+ public async Task TestRequestIssuedWhileChangingServerDoesNotHang()
+ {
+ _client = CreateClient(_firstPort);
+ await _client.Connect();
+ Assert.IsTrue(_client.connection.IsConnected(), "Precondition: connected to the first node.");
+
+ Task>> pending = await StartPendingRequestAsync();
+
+ // The socket the client is about to retire. Captured so the test can say which socket
+ // the follow-up request actually saw, rather than only that it hung.
+ WebSocketClient retiredSocket = _client.connection.ws;
+ Assert.IsNotNull(retiredSocket, "Precondition: a live socket to retire.");
+
+ Task>> followUp = null;
+ WebSocketClient socketSeenByFollowUp = null;
+
+ // ExecuteSynchronously, not an await: the continuation has to run on the thread that
+ // completes the promise, which is the thread inside the retirement sweep. That is the
+ // consumer shape the issue reports - a second value read from the response handler of
+ // the first - and it is what puts the follow-up request in the window.
+ Task continuation = pending.ContinueWith(
+ _ =>
+ {
+ socketSeenByFollowUp = _client.connection.ws;
+ followUp = _client.connection.Request(AccountInfoRequest());
+ },
+ TaskContinuationOptions.ExecuteSynchronously);
+
+ await _client.connection.ChangeServer($"ws://127.0.0.1:{_secondPort}");
+
+ await continuation;
+ Assert.IsNotNull(followUp, "The rejection of the first request must have issued a follow-up.");
+
+ Assert.AreNotSame(
+ retiredSocket,
+ socketSeenByFollowUp,
+ "A request issued during the switch still saw the socket being retired as the active one - " +
+ "the connectivity check it passed was about a socket that was already on its way out.");
+
+ Task finished = await Task.WhenAny(followUp, Task.Delay(AcceptableBound));
+
+ Assert.AreSame(
+ followUp,
+ finished,
+ $"A request issued while ChangeServer was retiring the old socket is still pending after " +
+ $"{AcceptableBound.TotalSeconds:F0}s. It was written to the retired socket and nothing will " +
+ $"complete it before RequestTimeout ({RequestTimeout.TotalSeconds:F0}s) expires.");
+
+ // Either outcome is correct: sent on the new connection (it completed within the bound,
+ // which is all a success has to show), or refused outright. What is not correct is
+ // waiting out RequestTimeout. A failure is captured first and judged afterwards, so an
+ // assertion failure is reported as itself rather than caught here.
+ Exception failure = null;
+ try
+ {
+ await followUp;
+ }
+ catch (Exception error)
+ {
+ failure = error;
+ }
+
+ if (failure is not null)
+ {
+ Assert.IsNotInstanceOfType(
+ failure,
+ "The follow-up request waited out RequestTimeout instead of being handled.");
+ Assert.IsInstanceOfType(
+ failure,
+ $"A refused follow-up must say the client is not connected, not fail with {failure.GetType().Name}: {failure.Message}");
+ }
+ }
+
+ ///
+ /// The same window on the user disconnect path. Here there is no new connection to carry
+ /// the request over to, so the only correct outcome is an immediate refusal.
+ ///
+ ///
+ /// Not regression coverage for the ordering - this passed before the fix too. On this
+ /// path the socket is really closed, so its close callback runs the second sweep in
+ /// OnceClose, which happened to catch the follow-up. ChangeServer filters that
+ /// callback out as a retiring session, which is why only the test above turned red. Kept
+ /// so that the refusal stays immediate if that incidental second sweep ever goes away.
+ ///
+ [TestMethod]
+ public async Task TestRequestIssuedWhileDisconnectingFailsFast()
+ {
+ _client = CreateClient(_firstPort);
+ await _client.Connect();
+ Assert.IsTrue(_client.connection.IsConnected(), "Precondition: connected to the first node.");
+
+ Task>> pending = await StartPendingRequestAsync();
+
+ Task>> followUp = null;
+
+ Task continuation = pending.ContinueWith(
+ _ => { followUp = _client.connection.Request(AccountInfoRequest()); },
+ TaskContinuationOptions.ExecuteSynchronously);
+
+ await _client.Disconnect();
+ _client = null; // Disconnected already; keep cleanup from doing it twice.
+
+ await continuation;
+ Assert.IsNotNull(followUp, "The rejection of the first request must have issued a follow-up.");
+
+ Task finished = await Task.WhenAny(followUp, Task.Delay(AcceptableBound));
+
+ Assert.AreSame(
+ followUp,
+ finished,
+ $"A request issued while Disconnect() was retiring the socket is still pending after " +
+ $"{AcceptableBound.TotalSeconds:F0}s - it went into the socket being closed and waits out " +
+ $"RequestTimeout ({RequestTimeout.TotalSeconds:F0}s).");
+
+ // Captured first, judged afterwards - see the test above.
+ Exception failure = null;
+ try
+ {
+ await followUp;
+ }
+ catch (Exception error)
+ {
+ failure = error;
+ }
+
+ Assert.IsNotNull(failure, "A request issued during Disconnect() must not succeed - there is no connection to serve it.");
+ Assert.IsNotInstanceOfType(
+ failure,
+ "The follow-up request waited out RequestTimeout instead of being refused.");
+ Assert.IsTrue(
+ failure is Xrpl.Client.Exceptions.NotConnectedException
+ or Xrpl.Client.Exceptions.DisconnectedException
+ or OperationCanceledException,
+ $"A request issued during Disconnect() must be refused as not connected, not {failure.GetType().Name}: {failure.Message}");
+ }
+ }
+}
diff --git a/Xrpl/Client/connection.cs b/Xrpl/Client/connection.cs
index 0d88f35e..5c70fd8a 100644
--- a/Xrpl/Client/connection.cs
+++ b/Xrpl/Client/connection.cs
@@ -334,7 +334,9 @@ private static WebSocketClient CreateWebSocket(string url, ConnectionOptions con
public string url { get; private set; }
- public WebSocketClient ws;
+ // Volatile: this is the connectivity gate. It is cleared under _disconnectLock by the
+ // retirement paths and read lock-free by ShouldBeConnected/State/CheckIfNotConnected.
+ public volatile WebSocketClient ws;
private int? reconnectTimeoutID = null;
@@ -492,8 +494,8 @@ internal long? ActiveSessionId
/// timer stopped taking the processor down with it.
///
/// Volatile.Read rather than a plain read or _messageProcessorLock: the field is
- /// written under that lock, and holding it here would mean waiting out
- /// StopMessageProcessorInternal, which blocks up to two seconds on the reader task.
+ /// written under that lock, and a read that never contends with a stop in progress is all
+ /// this needs.
///
///
internal bool IsMessageProcessorRunning => Volatile.Read(ref _streamMessageChannel) != null;
@@ -518,7 +520,7 @@ internal long? ActiveSessionId
///
/// Completes the stream channel's writer without clearing the channel, reproducing the state
- /// StopMessageProcessorInternal leaves behind for anyone who read
+ /// DetachMessageProcessor leaves behind for anyone who read
/// _streamMessageChannel just before it was cleared.
///
///
@@ -766,17 +768,14 @@ public async Task ChangeServer(
// processor is explicit since StopPingTimerSync no longer does it as a side effect - this
// session's queue goes with the session.
StopPingTimerSync();
- StopMessageProcessor();
-
- // 3. Reject all pending requests BEFORE waiting for ping
- // This allows the ping handler to receive OperationCanceledException and exit quickly
- requestManager.RejectAllWithCancellation();
- connectionManager.RejectAllAwaitingWithCancellation();
-
- // 4. Now wait for ping to finish (should be very fast since requests were rejected)
- await WaitForPingToFinishAsync();
- // 5. Mark old session as retiring (callbacks will be ignored)
+ // 3. Mark old session as retiring (callbacks will be ignored) and clear ws - BEFORE
+ // rejecting the pending requests, on purpose.
+ // The rejection sweep resumes consumer continuations - inline on this thread when there
+ // is no synchronization context - and a consumer that issues its next request from there
+ // must already see no usable connection. Cleared afterwards, that request passes the
+ // connectivity check on the old socket and is written into it after the sweep that would
+ // have rejected it; nothing completes it before RequestTimeout (issue #177).
ConnectionSession? oldSession;
lock (_sessionLock)
{
@@ -784,7 +783,6 @@ public async Task ChangeServer(
oldSession?.MarkAsRetiring();
}
- // 6. Capture old socket and clear ws reference
WebSocketClient? oldSocket;
lock (_disconnectLock)
{
@@ -792,6 +790,19 @@ public async Task ChangeServer(
ws = null;
}
+ // 4. Reject all pending requests BEFORE waiting for ping
+ // This allows the ping handler to receive OperationCanceledException and exit quickly
+ requestManager.RejectAllWithCancellation();
+ connectionManager.RejectAllAwaitingWithCancellation();
+
+ // 5. The message processor goes with the session, and its reader is let go of after the
+ // sweep, not before: consumers are released first, and this is the first yield of the
+ // switch - what a single-threaded host runs their continuations on.
+ await StopMessageProcessorAsync();
+
+ // 6. Now wait for ping to finish (should be very fast since requests were rejected)
+ await WaitForPingToFinishAsync();
+
// 7. Mark old socket for intentional disconnect (per-socket tracking only)
// CRITICAL: Do NOT set global _isIntentionalDisconnect = true here - same rule as the ping/network
// recovery path. The global flag was only reset in OnceOpen, so if the NEW server never came up it
@@ -915,17 +926,14 @@ private async Task RetireCurrentSessionAndReconnectAsync(string reason)
// 4. Stop ping timer and the message processor (but don't wait yet) - the queue belongs
// to the session being retired.
StopPingTimerSync();
- StopMessageProcessor();
-
- // 5. Reject all pending requests BEFORE waiting for ping
- // This allows the ping handler to receive OperationCanceledException and exit quickly
- requestManager.RejectAllWithCancellation();
- connectionManager.RejectAllAwaitingWithCancellation();
-
- // 6. Now wait for ping to finish (should be very fast since requests were rejected)
- await WaitForPingToFinishAsync().ConfigureAwait(false);
- // 7. Mark old session as retiring (callbacks will be ignored)
+ // 5. Mark old session as retiring (callbacks will be ignored) and clear ws - BEFORE
+ // rejecting the pending requests, on purpose.
+ // The rejection sweep resumes consumer continuations - inline on this thread when there
+ // is no synchronization context - and a consumer that issues its next request from there
+ // must already see no usable connection. Cleared afterwards, that request passes the
+ // connectivity check on the old socket and is written into it after the sweep that would
+ // have rejected it; nothing completes it before RequestTimeout (issue #177).
ConnectionSession? oldSession;
lock (_sessionLock)
{
@@ -933,7 +941,6 @@ private async Task RetireCurrentSessionAndReconnectAsync(string reason)
oldSession?.MarkAsRetiring();
}
- // 8. Capture old socket and clear ws reference
WebSocketClient? oldSocket;
lock (_disconnectLock)
{
@@ -941,6 +948,18 @@ private async Task RetireCurrentSessionAndReconnectAsync(string reason)
ws = null;
}
+ // 6. Reject all pending requests BEFORE waiting for ping
+ // This allows the ping handler to receive OperationCanceledException and exit quickly
+ requestManager.RejectAllWithCancellation();
+ connectionManager.RejectAllAwaitingWithCancellation();
+
+ // 7. The message processor goes with the session (see ChangeServer for the ordering).
+ await StopMessageProcessorAsync().ConfigureAwait(false);
+
+ // 8. Now wait for the ping to finish. This method runs inside the ping check itself, so
+ // the wait returns at once - see WaitForPingToFinishAsync.
+ await WaitForPingToFinishAsync().ConfigureAwait(false);
+
// 9. Mark old socket for intentional disconnect (per-socket tracking only)
// CRITICAL: Do NOT set global _isIntentionalDisconnect = true for ping/network recoveries!
// The global flag would block OnConnectionFailed from processing new connection failures.
@@ -1063,6 +1082,19 @@ await NotifySessionEndedAsync(oldSession, SessionEndReason.ConnectionLost, reaso
return;
}
+ // The wait above ends in cancellation when its own source is retired, and on success the
+ // path that retires it is OnceOpen: the attempt above failed at the socket, the failure
+ // callback started the reconnect loop on this same source, and that loop connected
+ // first. A client that is connected has nothing to reconnect. Treating this as a failure
+ // started a second loop, whose first attempt retired the live socket and opened another -
+ // one reconnect became two, with a RestoringConnection reported on a healthy client.
+ if (IsConnected())
+ {
+ _isFastReconnectActive = false;
+ Debug.WriteLine($"{DateTime.Now}Fast reconnect settled by the reconnect loop: {ex.Message}");
+ return;
+ }
+
// Start the loop BEFORE notifying: SetConnectionState calls into consumer code, and an
// exception from a handler must not cost us the reconnect loop. Ordering matters more
// than the message here - without the loop the client never comes back.
@@ -1375,42 +1407,41 @@ public async Task Disconnect()
_isIntentionalDisconnect = true;
_permanentlyDisconnected = true;
- var currentSocket = ws;
- if (currentSocket != null)
+ // Capture the socket and clear ws BEFORE rejecting the pending requests - see ChangeServer
+ // for why: the sweep runs consumer continuations, and a request issued from one of them
+ // must not find the socket being closed still installed as the connection (issue #177).
+ WebSocketClient? socketToClose;
+ lock (_disconnectLock)
{
- MarkSocketAsUserInitiated(currentSocket);
- currentSocket.SetIntentionalDisconnect();
+ socketToClose = ws;
+ ws = null;
+
+ if (socketToClose != null)
+ {
+ MarkSocketAsUserInitiated(socketToClose);
+ socketToClose.SetIntentionalDisconnect();
+
+ if (_disconnectTcs == null || _disconnectTcs.Task.IsCompleted)
+ {
+ _disconnectTcs = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
+ }
+ }
}
ClearReconnectState(); // Clear all reconnect state on user disconnect
StopPingTimerSync();
- StopMessageProcessor();
-
+
// Reject pending requests so ping handler can exit quickly
requestManager.RejectAllWithCancellation();
connectionManager.RejectAllAwaitingWithCancellation();
-
+
+ await StopMessageProcessorAsync();
await WaitForPingToFinishAsync();
- WebSocketClient? socketToClose;
- lock (_disconnectLock)
+ if (socketToClose == null)
{
- socketToClose = ws;
- ws = null;
-
- if (socketToClose == null)
- {
- SetConnectionState(XrpConnectionState.Disconnected, message: "Already disconnected.");
- return 0;
- }
-
- MarkSocketAsUserInitiated(socketToClose);
- socketToClose.SetIntentionalDisconnect();
-
- if (_disconnectTcs == null || _disconnectTcs.Task.IsCompleted)
- {
- _disconnectTcs = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
- }
+ SetConnectionState(XrpConnectionState.Disconnected, message: "Already disconnected.");
+ return 0;
}
Interlocked.Exchange(ref _userInitiatedSocket, socketToClose);
@@ -1431,46 +1462,57 @@ public async Task DisconnectAndWaitAsync(TimeSpan timeout, CancellationToken can
_isIntentionalDisconnect = true;
_permanentlyDisconnected = true;
- var currentSocket = ws;
- if (currentSocket != null)
+ // Same ordering as Disconnect(): the socket leaves ws before the sweep runs (issue #177).
+ TaskCompletionSource? tcs = null;
+ WebSocketClient? socketToClose;
+ lock (_disconnectLock)
{
- MarkSocketAsUserInitiated(currentSocket);
- currentSocket.SetIntentionalDisconnect();
+ socketToClose = ws;
+ ws = null;
+
+ if (socketToClose != null)
+ {
+ MarkSocketAsUserInitiated(socketToClose);
+ socketToClose.SetIntentionalDisconnect();
+
+ if (_disconnectTcs == null || _disconnectTcs.Task.IsCompleted)
+ {
+ _disconnectTcs = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
+ }
+
+ tcs = _disconnectTcs;
+ }
}
ClearReconnectState(); // Clear all reconnect state on user disconnect
StopPingTimerSync();
- StopMessageProcessor();
-
+
// Reject pending requests so ping handler can exit quickly
requestManager.RejectAllWithCancellation();
connectionManager.RejectAllAwaitingWithCancellation();
-
- await WaitForPingToFinishAsync();
- TaskCompletionSource tcs;
- WebSocketClient? socketToClose;
+ await StopMessageProcessorAsync();
+ await WaitForPingToFinishAsync();
- lock (_disconnectLock)
+ if (socketToClose == null || tcs == null)
{
- socketToClose = ws;
- ws = null;
-
- if (socketToClose == null)
+ // Nothing here to close - but another DisconnectAndWaitAsync may be mid-way, having
+ // taken the socket already. This call promised to return once the socket is gone, so
+ // it waits on that one's completion source rather than reporting a disconnect that
+ // has not finished.
+ TaskCompletionSource? inProgress;
+ lock (_disconnectLock)
{
- SetConnectionState(XrpConnectionState.Disconnected, message: "Already disconnected.");
- return;
+ inProgress = _disconnectTcs;
}
- MarkSocketAsUserInitiated(socketToClose);
- socketToClose.SetIntentionalDisconnect();
-
- if (_disconnectTcs == null || _disconnectTcs.Task.IsCompleted)
+ if (inProgress is { Task.IsCompleted: false })
{
- _disconnectTcs = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
+ await Task.WhenAny(inProgress.Task, Task.Delay(timeout, cancellationToken));
}
- tcs = _disconnectTcs;
+ SetConnectionState(XrpConnectionState.Disconnected, message: "Already disconnected.");
+ return;
}
Interlocked.Exchange(ref _userInitiatedSocket, socketToClose);
@@ -2021,8 +2063,11 @@ public async Task>> Request(
{
WebsocketSendAsync(ws, _request.Message);
}
- catch (EncodingFormatException error)
+ catch (Exception error) when (error is EncodingFormatException or DisconnectedException)
{
+ // The connection can be retired between the check above and this send. The request
+ // never left, so it must not stay pending until RequestTimeout: the rejection is what
+ // the await below surfaces to the caller.
requestManager.Reject(_request.Id, error);
}
@@ -2043,8 +2088,11 @@ public async Task> GRequest(
{
WebsocketSendAsync(ws, _request.Message);
}
- catch (EncodingFormatException error)
+ catch (Exception error) when (error is EncodingFormatException or DisconnectedException)
{
+ // The connection can be retired between the check above and this send. The request
+ // never left, so it must not stay pending until RequestTimeout: the rejection is what
+ // the await below surfaces to the caller.
requestManager.Reject(_request.Id, error);
}
@@ -2080,6 +2128,17 @@ private async Task OnceOpen(WebSocketClient connectedSocket, long sessionId)
return;
}
+ // A user Disconnect() that landed while this socket was still connecting has marked it and
+ // is closing it. Installing it here would undo the disconnect: ws restored, the
+ // intentional-disconnect tracking cleared below, and the close that follows read as a
+ // network drop that starts a reconnect loop. The session check above does not cover this -
+ // Disconnect() does not retire the session, because OnceClose is what announces a user
+ // disconnect - so the socket's own marks are the signal.
+ if (_permanentlyDisconnected || IsSocketUserInitiated(connectedSocket))
+ {
+ return;
+ }
+
// Verify the connected socket matches current ws, or update ws if it was cleared
if (ws == null)
{
@@ -2171,10 +2230,6 @@ private async Task OnceOpen(WebSocketClient connectedSocket, long sessionId)
/// The exception thrown by the handler.
private async Task OnConnectHandlerFailedAsync(WebSocketClient failedSocket, Exception error)
{
- int failures = Interlocked.Increment(ref _connectHandlerFailures);
-
- Debug.WriteLine($"{DateTime.Now}OnConnected handler failed ({failures}): {error.Message}");
-
var errorHandler = OnError;
if (errorHandler is not null)
{
@@ -2190,6 +2245,23 @@ await errorHandler
}
}
+ // Ownership first, before anything below counts or tears down. WebSocketClient.Connect invokes
+ // its OnConnect callback without awaiting it, so this can run after a newer socket has replaced
+ // the one whose handler failed. That socket's failure is not a failure of the current
+ // connection: it must not count towards giving up, and the give-up branch - RejectAll and
+ // Disconnect() - would take the live connection down for a callback that belongs to a dead one.
+ // Read-only here; the clear under the same lock happens below, once this path owns the teardown.
+ if (!IsCurrentSocket(failedSocket))
+ {
+ failedSocket.Cancel();
+ failedSocket.Disconnect();
+ return;
+ }
+
+ int failures = Interlocked.Increment(ref _connectHandlerFailures);
+
+ Debug.WriteLine($"{DateTime.Now}OnConnected handler failed ({failures}): {error.Message}");
+
bool giveUp = config.StopAfterMaxAttempts && failures >= config.MaxReconnectAttempts;
if (giveUp)
{
@@ -2205,6 +2277,16 @@ await errorHandler
$"OnConnected handler failed {failures} time(s) in a row: {error.Message}. Giving up after {config.MaxReconnectAttempts} attempts. Call Connect() to retry.",
ConnectionCloseSeverity.Error);
+ // The notification above ran consumer code. A handler that answered "gave up" with a
+ // ChangeServer has already taken this socket out of ws and is opening another; the
+ // teardown below would then reject that connection's requests and close its socket.
+ if (!IsCurrentSocket(failedSocket))
+ {
+ failedSocket.Cancel();
+ failedSocket.Disconnect();
+ return;
+ }
+
// Rejected here, before Disconnect(), and with the reason that is actually true. The
// requests in flight are being stopped because this client gave up connecting, not
// because anyone cancelled them - and Disconnect() rejects with cancellation, which is
@@ -2227,14 +2309,13 @@ await errorHandler
ConnectionCloseSeverity.Warning,
reconnect: BuildReconnectInfo(failures));
- StopPingTimerSync();
- StopMessageProcessor();
- requestManager.RejectAllWithCancellation();
- await WaitForPingToFinishAsync();
-
// Always tear down the socket the handler actually ran for. WebSocketClient.Connect invokes its
// OnConnect callback without awaiting it, so the connect lock can be released while this method is
// still running: by now `ws` may already point at a newer socket that must not be touched.
+ // Cleared before the sweep below, for the reason given in ChangeServer (issue #177): the
+ // socket is open, and a request issued from a rejected continuation would otherwise go into it.
+ // Taken after the notification above on purpose: the check that comes with the clear is the
+ // one that sees what the consumer's handler did.
bool wasCurrentSocket;
lock (_disconnectLock)
{
@@ -2245,17 +2326,28 @@ await errorHandler
}
}
- // The socket is deliberately NOT marked as user-initiated: OnceClose must treat this as a real
- // close so the standard reconnect path runs instead of the "closed permanently" branch.
- failedSocket.Cancel();
- failedSocket.Disconnect();
-
if (!wasCurrentSocket)
{
- // A newer connection already replaced this socket - it owns the reconnect state now.
+ // Replaced since the ownership check at the top - the OnError notification, the give-up
+ // branch and the RestoringConnection notification in between all hand control to consumer
+ // code, and a ChangeServer from any of them retires this socket itself. A newer connection
+ // owns the ping timer, the pending requests, the message processor and the reconnect state
+ // now; this callback closes the socket its handler ran for and steps aside.
+ failedSocket.Cancel();
+ failedSocket.Disconnect();
return;
}
+ StopPingTimerSync();
+ requestManager.RejectAllWithCancellation();
+ await StopMessageProcessorAsync();
+ await WaitForPingToFinishAsync();
+
+ // The socket is deliberately NOT marked as user-initiated: OnceClose must treat this as a real
+ // close so the standard reconnect path runs instead of the "closed permanently" branch.
+ failedSocket.Cancel();
+ failedSocket.Disconnect();
+
// Take ownership of the reconnect state instead of asking "is a loop already running?".
// This method can run inside the reconnect loop's own attempt: that loop breaks as soon as the
// socket reports Open, which happens before the handler has even finished failing. Both this check
@@ -2271,6 +2363,18 @@ await errorHandler
RestartReconnectLoop(initialAttempts: failures);
}
+ ///
+ /// Whether is the one installed as the connection right now. Read
+ /// under _disconnectLock, the lock every retirement path clears ws under.
+ ///
+ private bool IsCurrentSocket(WebSocketClient socket)
+ {
+ lock (_disconnectLock)
+ {
+ return ReferenceEquals(ws, socket);
+ }
+ }
+
///
/// Retires the current reconnect session and installs a fresh one in a single transaction,
/// seeding the attempt counter with .
@@ -2369,7 +2473,7 @@ private async Task OnceClose(int? code, string? description, WebSocketClient clo
// Only for the current socket - and the message processor goes with it, this connection
// is over.
StopPingTimerSync();
- StopMessageProcessor();
+ await StopMessageProcessorAsync();
// Check if this is a network drop (FailureReason set by WebSocketClient)
var isNetworkDrop = closingSocket.FailureReason == SocketFailureReason.NetworkDrop;
@@ -2657,6 +2761,24 @@ private async Task ReconnectLoopAsync(CancellationTokenSource ownCts)
// =====================================================
// SESSION ISOLATION (same as ChangeServer)
// =====================================================
+ // Somebody else connected while this loop was waiting out its delay - another
+ // attempt on the same source, or the fast-reconnect path. Whatever is in ws is live,
+ // and retiring it below would trade a healthy connection for a reconnect nobody
+ // needed. Every path that starts this loop does so because the connection is gone,
+ // so an open socket here always belongs to someone who got there first.
+ if (IsConnected())
+ {
+ lock (_reconnectStateLock)
+ {
+ if (ReferenceEquals(_reconnectCts, ownCts))
+ {
+ _reconnectAttempts = 0;
+ }
+ }
+
+ break;
+ }
+
// Mark old session as retiring before creating new connection
// so late callbacks from old socket are properly ignored.
ConnectionSession? oldSession;
@@ -2760,6 +2882,18 @@ private async Task ReconnectLoopAsync(CancellationTokenSource ownCts)
private volatile int _pingRunning = 0;
+ ///
+ /// True inside this connection's ping check and everything it awaits. The fast-reconnect path
+ /// is awaited from there, and it must be able to tell that the ping it would wait for is the
+ /// one it is running in.
+ ///
+ ///
+ /// Per instance, not static: the value follows the execution context, so a consumer's
+ /// OnPing handler that awaits another connection would carry a static flag into that
+ /// connection and let it skip waiting for its own ping.
+ ///
+ private readonly AsyncLocal _insidePingCheck = new AsyncLocal();
+
private Task? _pingLoopTask = null;
private System.Threading.Timer? _wasmPingTimer;
@@ -2807,6 +2941,8 @@ private async Task ExecutePingCheckAndReleaseAsync(CancellationTokenSource cts,
private async Task ExecutePingCheckAsync(CancellationTokenSource cts)
{
+ _insidePingCheck.Value = true;
+
try
{
if (cts.IsCancellationRequested)
@@ -3066,7 +3202,17 @@ private async Task WaitForPingToFinishAsync()
{
// Clear the task reference
Interlocked.Exchange(ref _currentPingTask, value: null);
-
+
+ // Called from inside the ping check - RetireCurrentSessionAndReconnectAsync is awaited from
+ // there, and only from there. The flag polled below is this very check's, and it cannot
+ // clear before the check returns, which is after this method: the wait could only ever run
+ // out its timeout, and did, on every ping-triggered reconnect. The ping is not running
+ // alongside the retirement here; it is the retirement.
+ if (_insidePingCheck.Value)
+ {
+ return;
+ }
+
// Wait for _pingRunning to become 0 (ping task's finally block will reset it)
// Since we already rejected pending requests, the ping should exit very quickly
var startTime = DateTime.UtcNow;
@@ -3241,8 +3387,11 @@ private void StartMessageProcessor()
{
lock (_messageProcessorLock)
{
- // Stop any existing processor first
- StopMessageProcessorInternal();
+ // Detach any leftover without waiting for it: its channel is completed and its source
+ // cancelled, so it exits on its own, and its frames belong to a session that is already
+ // retired. Waiting here used to block OnceOpen on a single-threaded host.
+ (Task? leftoverTask, CancellationTokenSource? leftoverCts) = DetachMessageProcessor();
+ _ = AwaitMessageProcessorExitAsync(leftoverTask, leftoverCts);
// Create new session-bound channel and CTS
// Using bounded channel to prevent memory issues under high load
@@ -3301,47 +3450,66 @@ private void StartMessageProcessor()
///
/// Stops the background message processor and disposes resources.
///
- private void StopMessageProcessor()
+ private Task StopMessageProcessorAsync()
{
+ (Task? task, CancellationTokenSource? cts) detached;
lock (_messageProcessorLock)
{
- StopMessageProcessorInternal();
+ detached = DetachMessageProcessor();
}
+
+ return AwaitMessageProcessorExitAsync(detached.task, detached.cts);
}
///
- /// Internal stop logic - must be called with _messageProcessorLock held.
- /// Completes the channel, cancels the CTS, and awaits task completion.
+ /// Takes the processor's channel, source and task out of their fields, completes the channel
+ /// and cancels the source, so the reader exits on its own. Must be called with
+ /// _messageProcessorLock held. Does not wait for the reader: that is
+ /// , outside the lock.
///
- private void StopMessageProcessorInternal()
+ private (Task? task, CancellationTokenSource? cts) DetachMessageProcessor()
{
var channel = _streamMessageChannel;
var cts = _messageProcessorCts;
var task = _messageProcessorTask;
-
+
_streamMessageChannel = null;
_messageProcessorCts = null;
_messageProcessorTask = null;
-
+
// Complete the channel first to unblock WaitToReadAsync
if (channel != null)
{
try { channel.Writer.Complete(); } catch { }
}
-
+
// Then cancel the CTS
if (cts != null)
{
try { cts.Cancel(); } catch { }
}
-
- // Wait for task to complete (with timeout to prevent deadlock)
+
+ return (task, cts);
+ }
+
+ ///
+ /// Waits up to two seconds for a detached reader to exit, then disposes its source.
+ ///
+ ///
+ /// This used to be a blocking Task.Wait with the same cap. On a single-threaded host
+ /// (Blazor WebAssembly) the reader's continuation needs the very thread that was blocked in
+ /// order to observe the completed channel, so the wait never returned early: every
+ /// ChangeServer, fast reconnect and Disconnect stalled the UI for the full two
+ /// seconds. Awaited, the reader runs and is gone in milliseconds. The cap is kept for a reader
+ /// stuck inside a consumer handler, and the source is disposed regardless, as before.
+ ///
+ private static async Task AwaitMessageProcessorExitAsync(Task? task, CancellationTokenSource? cts)
+ {
if (task != null)
{
- try { task.Wait(TimeSpan.FromSeconds(2)); } catch { }
+ await Task.WhenAny(task, Task.Delay(TimeSpan.FromSeconds(2))).ConfigureAwait(false);
}
-
- // Dispose resources
+
cts?.Dispose();
}
@@ -3805,7 +3973,7 @@ private void EnqueueStreamMessage(byte[] frame, long? sessionId = null)
// The channel is bounded with DropOldest, so a full queue is not a refusal: TryWrite
// evicts the oldest frame, counts it through itemDropped and reports success. It
// refuses only a completed writer - and that happens on the ordinary path, not just in
- // some corner: StopMessageProcessorInternal completes the writer after clearing
+ // some corner: DetachMessageProcessor completes the writer after clearing
// _streamMessageChannel, so a reader that got the reference an instant earlier writes
// into a channel that is already closed. StartPingTimer tears the processor down and
// StartMessageProcessor builds it again on every connect, so the window recurs.
diff --git a/Xrpl/Xrpl.csproj b/Xrpl/Xrpl.csproj
index 7e3c7065..41014b18 100644
--- a/Xrpl/Xrpl.csproj
+++ b/Xrpl/Xrpl.csproj
@@ -14,7 +14,7 @@
Apache-2.0
https://github.com/StaticBit-io/XrplCSharp
XrplCSharp
- 11.3.1.0
+ 11.3.2.0