This is an automated email from the ASF dual-hosted git repository.
numinnex pushed a commit to branch multi_endpoint_failover
in repository https://gitbox.apache.org/repos/asf/iggy.git
The following commit(s) were added to refs/heads/multi_endpoint_failover by
this push:
new d0a95b1cb address review comments
d0a95b1cb is described below
commit d0a95b1cbb7966a18946ca16c26246696cf545a0
Author: Grzegorz Koszyk <[email protected]>
AuthorDate: Wed Aug 26 09:15:42 2026 +0200
address review comments
---
core/sdk/src/tcp/tcp_client.rs | 128 +++++++++++++--------
.../Iggy_SDK.Tests.Integration/HeartbeatTests.cs | 26 ++---
.../Implementations/TcpMessageStream.Vsr.cs | 8 +-
.../IggyClient/Implementations/TcpMessageStream.cs | 35 +++---
.../VsrTests/EndpointFailoverTests.cs | 114 ++++++++++--------
foreign/csharp/README.md | 27 +++--
foreign/go/client/tcp/tcp_connect_test.go | 13 ++-
foreign/go/client/tcp/tcp_core.go | 79 +++++++------
foreign/go/client/tcp/tcp_core_review_test.go | 13 ++-
foreign/go/client/tcp/tcp_failover_test.go | 27 +++--
foreign/go/client/tcp/tcp_session_management.go | 18 +--
.../iggy/client/async/tcp/AsyncIggyTcpClient.java | 84 ++++++--------
.../iggy/client/async/tcp/AsyncTcpConnection.java | 15 ---
.../AsyncIggyTcpClientTransientFailoverTest.java | 104 ++++-------------
foreign/node/src/client/client.connection.ts | 11 +-
foreign/node/src/client/client.socket.test.ts | 72 ++++++++++++
foreign/node/src/client/client.socket.ts | 84 +++++++++-----
foreign/node/src/client/client.type.ts | 8 +-
18 files changed, 485 insertions(+), 381 deletions(-)
diff --git a/core/sdk/src/tcp/tcp_client.rs b/core/sdk/src/tcp/tcp_client.rs
index ba3080bc5..59cb24b01 100644
--- a/core/sdk/src/tcp/tcp_client.rs
+++ b/core/sdk/src/tcp/tcp_client.rs
@@ -104,6 +104,11 @@ pub struct TcpClient {
/// onto this node or, after a failover, another one -- can re-establish
/// the session instead of surfacing `Unauthenticated`. Cleared on logout.
session_credentials: Mutex<Option<RememberedSignIn>>,
+ /// The password a committed change gave the user a configured `AutoLogin`
+ /// signs in as. The configured credentials cannot be rewritten, and the
+ /// password they carry is dead once the change commits, so every later
+ /// sign-in reads this instead.
+ configured_password: Mutex<Option<SecretString>>,
// `std::sync::Mutex` (not `tokio::sync::Mutex`): the critical section
// is `encode_request_header`, which is pure CPU and never awaits. The
// tokio variant would pay a waker alloc + internal semaphore on
@@ -120,11 +125,6 @@ pub struct TcpClient {
struct RememberedSignIn {
credentials: Credentials,
user_id: u32,
- /// Set by a committed password change for `user_id`. The configured
- /// `AutoLogin` credentials cannot be rewritten -- the config is shared and
- /// immutable -- so this marks the remembered copy as the newer one for
- /// that same user, and [`TcpClient::sign_in_credentials`] prefers it.
- password_refreshed: bool,
}
/// A connection that completed every step of coming up, TLS included.
@@ -364,7 +364,6 @@ impl iggy_common::VsrSessionControl for TcpClient {
.replace(RememberedSignIn {
credentials,
user_id,
- password_refreshed: false,
});
}
@@ -388,9 +387,27 @@ impl iggy_common::VsrSessionControl for TcpClient {
.get_cow_str_value()
.is_ok_and(|name| name.as_ref() == username),
};
- if targets_session_user {
- *password = SecretString::from(new_password.to_owned());
- sign_in.password_refreshed = true;
+ if !targets_session_user {
+ return;
+ }
+
+ *password = SecretString::from(new_password.to_owned());
+ // The configured credentials cannot be rewritten -- the config is
+ // shared and immutable -- and the password they carry will never work
+ // again, so the new one is kept here instead. Kept outside the
+ // remembered sign-in on purpose: that record is replaced wholesale by
+ // every later login, so a marker on it would survive exactly one
+ // reconnect and the one after that would replay the dead password.
+ let configured_user_changed = matches!(
+ &self.config.auto_login,
+ AutoLogin::Enabled(Credentials::UsernamePassword(configured, _))
+ if configured == username
+ );
+ if configured_user_changed {
+ self.configured_password
+ .lock()
+ .await
+ .replace(SecretString::from(new_password.to_owned()));
}
}
@@ -457,6 +474,7 @@ impl TcpClient {
leader_redirection_state:
Mutex::new(LeaderRedirectionState::new()),
current_server_address: Mutex::new(server_address),
roster_endpoints: Mutex::new(Vec::new()),
+ configured_password: Mutex::new(None),
session_credentials: Mutex::new(None),
consensus_session:
Arc::new(StdMutex::new(ConsensusSession::new())),
skip_auto_login_once: Mutex::new(false),
@@ -759,44 +777,44 @@ impl TcpClient {
/// else the ones a manual sign-in on this client succeeded with. A manual
/// sign-in is otherwise less reconnectable than a configured one, which
/// is a surprising difference between two ways of doing the same thing.
+ ///
+ /// A password change this client committed for the configured user is
+ /// applied on top: the configured password will never work again, and
every
+ /// later reconnect would otherwise fail `InvalidCredentials`.
async fn sign_in_credentials(&self) -> Option<Credentials> {
- let remembered =
- self.session_credentials
+ match &self.config.auto_login {
+ AutoLogin::Enabled(Credentials::UsernamePassword(username,
configured_password)) => {
+ let password = self
+ .configured_password
+ .lock()
+ .await
+ .clone()
+ .unwrap_or_else(|| configured_password.clone());
+ Some(Credentials::UsernamePassword(username.clone(), password))
+ }
+ AutoLogin::Enabled(credentials) => Some(credentials.clone()),
+ AutoLogin::Disabled => self
+ .session_credentials
.lock()
.await
.as_ref()
- .map(|remembered: &RememberedSignIn| {
- (
- remembered.credentials.clone(),
- remembered.password_refreshed,
- )
- });
-
- match (&self.config.auto_login, remembered) {
- // A committed password change for the configured user outranks the
- // configured password: the config cannot be rewritten, and signing
- // in with the password this client itself replaced would fail
- // `InvalidCredentials` on every later reconnect. Same user either
- // way -- `refresh_session_password` only marks a change that
- // targeted the signed-in one.
- (
-
AutoLogin::Enabled(Credentials::UsernamePassword(configured_username, _)),
- Some((Credentials::UsernamePassword(username, password),
true)),
- ) if configured_username == &username => {
- Some(Credentials::UsernamePassword(username, password))
- }
- (AutoLogin::Enabled(configured), _) => Some(configured.clone()),
- (AutoLogin::Disabled, remembered) => remembered.map(|(credentials,
_)| credentials),
+ .map(|remembered| remembered.credentials.clone()),
}
}
/// Endpoints to dial for one connect, likeliest first: where the client
- /// currently is, then the roster it learned while connected, then the
- /// configured seeds.
+ /// currently is, the addresses it was configured with, then the roster it
+ /// learned while connected.
+ ///
+ /// Configured before learned, as in the other SDKs: those are the
endpoints
+ /// the caller vouched for, while a roster read from a cluster that has
since
+ /// changed shape may name nodes that are gone.
async fn dial_candidates(&self) -> Vec<String> {
let mut candidates =
vec![self.current_server_address.lock().await.clone()];
let roster = self.roster_endpoints.lock().await.clone();
- for endpoint in
roster.iter().chain(self.config.failover_addresses.iter()) {
+ let configured = std::iter::once(&self.config.server_address)
+ .chain(self.config.failover_addresses.iter());
+ for endpoint in configured.chain(roster.iter()) {
let mut known = false;
for candidate in &candidates {
if is_same_address(candidate, endpoint).await {
@@ -1609,15 +1627,16 @@ mod tests {
"127.0.0.1:8092".to_string(),
];
- // The current endpoint leads, the roster follows, and neither the
- // roster's copy of the current endpoint nor a seed that only spells
- // the same endpoint differently earns a second dial.
+ // The current endpoint leads, the configured ones follow, and the
+ // roster comes last -- the same order as the other SDKs. Neither the
+ // roster's copy of an endpoint already named nor a seed that only
+ // spells one differently earns a second dial.
assert_eq!(
client.dial_candidates().await,
vec![
"127.0.0.1:8090".to_string(),
- "127.0.0.1:8091".to_string(),
"127.0.0.1:8092".to_string(),
+ "127.0.0.1:8091".to_string(),
]
);
}
@@ -1804,12 +1823,31 @@ mod tests {
)
.await;
- match client.sign_in_credentials().await {
- Some(Credentials::UsernamePassword(username, password)) => {
- assert_eq!(username, "iggy");
- assert_eq!(password.expose_secret(), "new");
+ // Every later reconnect, not just the next one: each of them signs in
+ // and remembers that sign-in afresh, and the configured password is
+ // dead for good once the change commits.
+ for reconnect in 0..3 {
+ match client.sign_in_credentials().await {
+ Some(Credentials::UsernamePassword(username, password)) => {
+ assert_eq!(username, "iggy");
+ assert_eq!(
+ password.expose_secret(),
+ "new",
+ "the configured password came back on reconnect
{reconnect}"
+ );
+ // What the reconnect's own login does with what it signed
+ // in with.
+ client
+ .remember_session_credentials(
+ Credentials::UsernamePassword(username, password),
+ SESSION_USER_ID,
+ )
+ .await;
+ }
+ other => {
+ panic!("expected the configured user with the new
password, got {other:?}")
+ }
}
- other => panic!("expected the configured user with the new
password, got {other:?}"),
}
}
diff --git a/foreign/csharp/Iggy_SDK.Tests.Integration/HeartbeatTests.cs
b/foreign/csharp/Iggy_SDK.Tests.Integration/HeartbeatTests.cs
index 7a9de7fb6..b0e71feee 100644
--- a/foreign/csharp/Iggy_SDK.Tests.Integration/HeartbeatTests.cs
+++ b/foreign/csharp/Iggy_SDK.Tests.Integration/HeartbeatTests.cs
@@ -104,28 +104,28 @@ public class HeartbeatTests
me.ConsumerGroupsCount.ShouldBe(0);
}
+ /// <summary>
+ /// An eviction is the server's heartbeat verifier reacting to
silence, not caller intent, so a client
+ /// that signed in by hand recovers from it exactly like one whose
credentials were configured: the
+ /// sign-in it remembered re-establishes the session. Only an explicit
sign-out or Dispose ends it.
+ /// </summary>
[Test]
- public async Task
EvictedClient_WithoutAutoLogin_Should_FailFast_And_NotReconnect()
+ public async Task
EvictedClient_WithoutAutoLogin_Should_ReestablishItsSession()
{
using var client = await CreateClient(TimeSpan.FromHours(1), false);
await client.LoginUserAsync("iggy", "iggy");
var (streamName, _) = await JoinFreshGroup(client);
- var reconnected = false;
- client.SubscribeConnectionEvents(args =>
- {
- reconnected |= args.CurrentState == ConnectionState.Connecting;
- return Task.CompletedTask;
- });
-
await Task.Delay(IdleFor);
- // A reconnect could not bring the session back, so the request
surfaces the loss instead of coming back
- // over an unauthenticated connection.
+ // The evicted request surfaces the loss; the one after it comes back
over a session the remembered
+ // sign-in re-established.
await Should.ThrowAsync<Exception>(() => client.GetMeAsync());
- reconnected.ShouldBeFalse();
- await Should.ThrowAsync<NotConnectedException>(() =>
- client.GetStreamByIdAsync(Identifier.String(streamName)));
+
+ var stream = await
client.GetStreamByIdAsync(Identifier.String(streamName));
+ stream.ShouldNotBeNull();
+ var me = await client.GetMeAsync();
+ me.ShouldNotBeNull();
}
private Task<IIggyClient> CreateClient(TimeSpan heartbeatInterval, bool
autoLogin = true)
diff --git
a/foreign/csharp/Iggy_SDK/IggyClient/Implementations/TcpMessageStream.Vsr.cs
b/foreign/csharp/Iggy_SDK/IggyClient/Implementations/TcpMessageStream.Vsr.cs
index ea41f1e0f..e100fcc54 100644
--- a/foreign/csharp/Iggy_SDK/IggyClient/Implementations/TcpMessageStream.Vsr.cs
+++ b/foreign/csharp/Iggy_SDK/IggyClient/Implementations/TcpMessageStream.Vsr.cs
@@ -569,11 +569,9 @@ public sealed partial class TcpMessageStream :
ISessionGenerationProvider
if (attempt.Error is VsrSessionEvictedException evicted)
{
- // Whatever the eviction interrupted, the session it
evicted is gone. Reported as an
- // outcome-unknown write it never reaches the
lost-connection path, so the remembered
- // sign-in would otherwise survive it and the next
reconnect would resurrect the session.
- ForgetSessionAfterEviction(evicted.Verdict);
-
+ // The session is gone, but the sign-in that established
it is not: a stale-client
+ // eviction comes off the server's heartbeat timer, so the
reconnect re-establishes it.
+ // Only an explicit sign-out or Dispose ends it.
if (attempt.RequestStarted &&
!VsrOperations.IsReplaySafeRead(code, isLoginRegister, body.Span))
{
throw new VsrRequestOutcomeUnknownException(evicted);
diff --git
a/foreign/csharp/Iggy_SDK/IggyClient/Implementations/TcpMessageStream.cs
b/foreign/csharp/Iggy_SDK/IggyClient/Implementations/TcpMessageStream.cs
index a1e267c3c..38bbfebb2 100644
--- a/foreign/csharp/Iggy_SDK/IggyClient/Implementations/TcpMessageStream.cs
+++ b/foreign/csharp/Iggy_SDK/IggyClient/Implementations/TcpMessageStream.cs
@@ -91,7 +91,8 @@ public sealed partial class TcpMessageStream : IIggyClient
// The credentials a sign-in succeeded with, so a reconnect - on this node
or, after a failover, another one -
// can re-establish the session instead of leaving every later request
unauthenticated. A caller that signs in
// by hand is otherwise less reconnectable than one that configures auto
login, which is a surprising
- // difference between two ways of doing the same thing. Cleared on
sign-out and on a server eviction.
+ // difference between two ways of doing the same thing. Cleared on
sign-out and on Dispose; a server-side
+ // eviction is not caller intent and leaves them in place.
private AutoLoginSettings? _rememberedLogin;
private bool IsConnecting => Volatile.Read(ref _isConnecting) != 0;
@@ -1348,7 +1349,16 @@ public sealed partial class TcpMessageStream :
IIggyClient
/// </summary>
private static bool IsTlsConfigurationFault(Exception e)
{
- return e is AuthenticationException or InvalidCertificatePathException;
+ if (e is InvalidCertificatePathException)
+ {
+ return true;
+ }
+
+ // AuthenticationException also carries a handshake that died on the
wire - a reset, a closed socket -
+ // and that says nothing about the configuration. Only a verdict
reached without transport trouble is
+ // one no retry can change.
+ return e is AuthenticationException
+ && e.InnerException is not (IOException or SocketException);
}
private async Task SendAckAsync(int code, ReadOnlyMemory<byte> body,
CancellationToken token)
@@ -1367,12 +1377,9 @@ public sealed partial class TcpMessageStream :
IIggyClient
{
_logger.LogWarning("Connection lost");
- // A server-side eviction is the server ending this session
authoritatively, like a logout: the
- // remembered sign-in ends with it, so only a configured auto
login may bring the session back.
- // Remembered credentials exist for transport loss, where the
session died with the socket rather
- // than by anyone's decision.
- ForgetSessionAfterEviction(e);
-
+ // A stale-client eviction is not caller intent: the server's
heartbeat verifier sends it after a
+ // gc pause or a laptop sleep, so the remembered sign-in survives
it and the reconnect below
+ // re-establishes the session. Only an explicit sign-out or
Dispose ends it. Same rule in every SDK.
if (!_configuration.ReconnectionSettings.Enabled)
{
_logger.LogWarning("Reconnection is disabled");
@@ -1402,18 +1409,6 @@ public sealed partial class TcpMessageStream :
IIggyClient
|| e is IggyInvalidStatusCodeException { StatusCode:
VsrError.STALE_CLIENT, FromServer: true };
}
- // An eviction ends this session authoritatively, like a logout: the
sign-in this client remembered goes
- // with it, so only a configured auto login may bring the session back.
Called from the consensus layer,
- // which sees the eviction whatever it interrupted - an eviction that
landed on a replicated write is
- // reported as VsrRequestOutcomeUnknownException and never reaches the
lost-connection path at all.
- private void ForgetSessionAfterEviction(Exception verdict)
- {
- if (verdict is IggyInvalidStatusCodeException { StatusCode:
VsrError.STALE_CLIENT, FromServer: true })
- {
- _rememberedLogin = null;
- }
- }
-
private async Task<IMemoryOwner<byte>> HandleReconnectionAsync(int code,
ReadOnlyMemory<byte> body,
bool autoLogin, CancellationToken token)
{
diff --git a/foreign/csharp/Iggy_SDK_Tests/VsrTests/EndpointFailoverTests.cs
b/foreign/csharp/Iggy_SDK_Tests/VsrTests/EndpointFailoverTests.cs
index 2d8880184..b835fbe56 100644
--- a/foreign/csharp/Iggy_SDK_Tests/VsrTests/EndpointFailoverTests.cs
+++ b/foreign/csharp/Iggy_SDK_Tests/VsrTests/EndpointFailoverTests.cs
@@ -105,12 +105,12 @@ public sealed class EndpointFailoverTests
/// <summary>
/// Mirrors the integration contract (HeartbeatTests
- /// EvictedClient_WithoutAutoLogin_Should_FailFast_And_NotReconnect):
a server-side eviction ends the
- /// session authoritatively, so the credentials a manual sign-in
remembered must not resurrect it - the
- /// evicted request surfaces the loss with no reconnect attempt.
+ /// EvictedClient_WithoutAutoLogin_Should_ReestablishItsSession): an
eviction comes off the server's
+ /// heartbeat timer rather than from the caller, so the sign-in this
client remembered survives it and
+ /// the reconnect re-establishes the session. Same rule in every SDK.
/// </summary>
[Fact]
- public async Task ServerEvictionForgetsTheRememberedSignIn()
+ public async Task ServerEvictionReplaysTheRememberedSignIn()
{
using var node = new MockNode();
var evict = false;
@@ -121,11 +121,15 @@ public sealed class EndpointFailoverTests
return Reply(OperationRegister, RegisterBody(session: 128));
}
- return evict
- ? EvictionFrame(EvictionStaleClient)
- : Reply(OperationNonReplicated, request.Code ==
GetClusterMetadataCode
- ? ClusterMetadata(node.Port, node.Port, node.Port)
- : []);
+ if (evict)
+ {
+ evict = false;
+ return EvictionFrame(EvictionStaleClient);
+ }
+
+ return Reply(OperationNonReplicated, request.Code ==
GetClusterMetadataCode
+ ? ClusterMetadata(node.Port, node.Port, node.Port)
+ : []);
});
var configuration = new IggyClientConfigurator
@@ -144,29 +148,25 @@ public sealed class EndpointFailoverTests
await client.ConnectAsync(TestContext.Current.CancellationToken);
await client.LoginUserAsync("iggy", "iggy",
TestContext.Current.CancellationToken);
await client.PingAsync(TestContext.Current.CancellationToken);
- var connectionsBeforeEviction = node.Connections;
+ var registrationsBeforeEviction = node.Registrations;
+ // A ping is replay-safe, so the eviction is absorbed: the reconnect
signs in again with the
+ // remembered credentials and the request completes over the session
it re-established.
evict = true;
- var evicted = await
Assert.ThrowsAsync<IggyInvalidStatusCodeException>(() =>
- client.PingAsync(TestContext.Current.CancellationToken));
- Assert.Equal(VsrError.STALE_CLIENT, evicted.StatusCode);
- Assert.True(evicted.FromServer);
- Assert.Equal(connectionsBeforeEviction, node.Connections);
-
- // The dropped connection leaves the next call transport-shaped, but
the eviction forgot the remembered
- // sign-in, so it must fail fast instead of reconnecting into a
resurrected session.
- await Assert.ThrowsAsync<NotConnectedException>(() =>
- client.PingAsync(TestContext.Current.CancellationToken));
- Assert.Equal(connectionsBeforeEviction, node.Connections);
+ await client.PingAsync(TestContext.Current.CancellationToken);
+
+ Assert.True(node.Registrations > registrationsBeforeEviction,
+ "the reconnect signed in again with the remembered credentials");
+ await client.PingAsync(TestContext.Current.CancellationToken);
}
/// <summary>
- /// The same contract when the eviction lands on a replicated write:
that request is reported as
- /// outcome-unknown rather than as a lost connection, so it never
passes through the lost-connection
- /// path, and the session it evicted still has to be forgotten.
+ /// The same rule when the eviction lands on a replicated write: that
request is reported as
+ /// outcome-unknown, because its own outcome is unknown, but the
session behind it is still
+ /// re-established for the requests that follow.
/// </summary>
[Fact]
- public async Task
ServerEvictionDuringAReplicatedWriteForgetsTheRememberedSignIn()
+ public async Task
ServerEvictionDuringAReplicatedWriteReplaysTheRememberedSignIn()
{
using var node = new MockNode();
var evict = false;
@@ -177,11 +177,15 @@ public sealed class EndpointFailoverTests
return Reply(OperationRegister, RegisterBody(session: 128));
}
- return evict
- ? EvictionFrame(EvictionStaleClient)
- : Reply(OperationNonReplicated, request.Code ==
GetClusterMetadataCode
- ? ClusterMetadata(node.Port, node.Port, node.Port)
- : []);
+ if (evict)
+ {
+ evict = false;
+ return EvictionFrame(EvictionStaleClient);
+ }
+
+ return Reply(OperationNonReplicated, request.Code ==
GetClusterMetadataCode
+ ? ClusterMetadata(node.Port, node.Port, node.Port)
+ : []);
});
var configuration = new IggyClientConfigurator
@@ -200,28 +204,27 @@ public sealed class EndpointFailoverTests
await client.ConnectAsync(TestContext.Current.CancellationToken);
await client.LoginUserAsync("iggy", "iggy",
TestContext.Current.CancellationToken);
await client.PingAsync(TestContext.Current.CancellationToken);
- var connectionsBeforeEviction = node.Connections;
+ var registrationsBeforeEviction = node.Registrations;
evict = true;
await Assert.ThrowsAsync<VsrRequestOutcomeUnknownException>(() =>
client.CreateStreamAsync("evicted-mid-write", token:
TestContext.Current.CancellationToken));
- Assert.Equal(connectionsBeforeEviction, node.Connections);
- await Assert.ThrowsAsync<NotConnectedException>(() =>
- client.PingAsync(TestContext.Current.CancellationToken));
- Assert.Equal(connectionsBeforeEviction, node.Connections);
+ await client.PingAsync(TestContext.Current.CancellationToken);
+ Assert.True(node.Registrations > registrationsBeforeEviction,
+ "the reconnect signed in again with the remembered credentials");
}
/// <summary>
- /// A survivor that only comes up after the first rotation still has
to be found. The reconnection
- /// budget counts rotations, not dials, so one retry is one full pass
over every endpoint the client
- /// knows rather than one dial of the endpoint it started from.
+ /// A survivor that is not listening yet when its node dies still has
to be found: the client keeps
+ /// rotating over every endpoint it knows, so one that comes up while
it is retrying is dialed on a
+ /// later pass rather than only on the first.
/// </summary>
[Fact]
- public async Task ResumesOnASurvivorThatComesUpAfterTheFirstRotation()
+ public async Task ResumesOnASurvivorThatComesUpWhileTheClientIsRetrying()
{
- // A port nothing listens on yet: the survivor comes up on it only
after the first rotation has
- // already failed on both endpoints.
+ // A port nothing listens on: the survivor binds it only after the
client has already failed on both
+ // endpoints, so the first pass cannot be the one that finds it.
var probe = new TcpListener(IPAddress.Loopback, 0);
probe.Start();
var survivorPort = (ushort)((IPEndPoint)probe.LocalEndpoint).Port;
@@ -240,10 +243,8 @@ public sealed class EndpointFailoverTests
ReconnectionSettings = new ReconnectionSettings
{
Enabled = true,
- // One retry: the rotation after the first failure is the last
one, which is where the
- // budget check used to cut the sweep short.
- MaxRetries = 1,
- InitialDelay = TimeSpan.FromMilliseconds(600)
+ MaxRetries = 4,
+ InitialDelay = TimeSpan.FromMilliseconds(200)
}
};
using var client = new TcpMessageStream(configuration,
NullLoggerFactory.Instance);
@@ -253,19 +254,32 @@ public sealed class EndpointFailoverTests
await client.PingAsync(TestContext.Current.CancellationToken);
primary.Kill();
- using var survivor = new MockNode(survivorPort);
+ MockNode? survivor = null;
+ // Constructed inside the delay, because the listener starts in the
constructor: built up front, the
+ // survivor would be answering from the very first dial and nothing
about the later passes would be
+ // exercised.
var comesUp = Task.Run(async () =>
{
- await Task.Delay(250, TestContext.Current.CancellationToken);
+ await Task.Delay(300, TestContext.Current.CancellationToken);
+ survivor = new MockNode(survivorPort);
survivor.Serve(request => request.Code == GetClusterMetadataCode
? Reply(OperationNonReplicated, ClusterMetadata(primary.Port,
survivorPort, survivorPort))
: Answer(request));
}, TestContext.Current.CancellationToken);
- await client.PingAsync(TestContext.Current.CancellationToken);
- await comesUp;
+ try
+ {
+ await client.PingAsync(TestContext.Current.CancellationToken);
+ await comesUp;
- Assert.True(survivor.Registrations >= 1, "the session was
re-established on the survivor");
+ Assert.NotNull(survivor);
+ Assert.True(survivor!.Registrations >= 1, "the session was
re-established on the survivor");
+ }
+ finally
+ {
+ await comesUp;
+ survivor?.Dispose();
+ }
}
private static byte[] EvictionFrame(byte reason)
diff --git a/foreign/csharp/README.md b/foreign/csharp/README.md
index 316ef3cb0..3fed133b1 100644
--- a/foreign/csharp/README.md
+++ b/foreign/csharp/README.md
@@ -108,8 +108,9 @@ var client = IggyClientFactory.CreateClient(new
IggyClientConfigurator
BackoffMultiplier = 2.0
},
- // Auto-login after connection. Reconnection needs it: without credentials
to replay a reconnect cannot
- // restore the session, so a lost connection fails the request instead
+ // Auto-login after connection. Optional for reconnection: a client that
signs in with
+ // LoginUserAsync has that sign-in replayed on a reconnect too. Without
either, a reconnect
+ // cannot restore the session and a lost connection fails the request
AutoLoginSettings = AutoLoginSettings.For("your_username",
"your_password"),
// or AutoLoginSettings.ForPersonalAccessToken("your_token")
@@ -196,14 +197,20 @@ The SDK replays a request whenever the server says it
never admitted it. Two cas
to `WithConnection`. Before, a builder-created client came back from a
reconnect unauthenticated; now the credentials are held for the lifetime of
the connection and replayed.
- The TCP client now pings the server every `HeartbeatInterval` (5 seconds,
always on) on its own, and
- reconnection is on by default (it was off before) with unlimited retries,
like the Rust client. Only a
- failed dial is retried; a rejected certificate, bad credentials or a missing
leader is thrown right away.
- A dropped connection fails every in-flight request at once; they share a
single reconnect and are replayed
- on the connection it establishes. With the default `MaxRetries = 0` an
unreachable server is retried
- forever, so a request that passes no `CancellationToken` waits for as long
as the server stays down - set
- `MaxRetries` or pass a token to bound it. Reconnection only replays a
request when `AutoLoginSettings` can
- restore the session; a client that logged in by hand fails fast on a lost
connection. Set
- `ReconnectionSettings.Enabled = false` to opt out of reconnection.
+ reconnection is on by default (it was off before) with unlimited retries,
like the Rust client. Bringing an
+ endpoint up is what gets another attempt, and one retry is a full pass over
every endpoint the client knows -
+ where it is, the configured address, and every node the roster named -
rather than one dial of the first.
+ Bad credentials or a missing leader is thrown right away, and so is a TLS
fault no retry can fix (an
+ unreadable CA file, a certificate this client will never accept), once the
pass has given the other
+ endpoints their turn. A dropped connection fails every in-flight request at
once; they share a single
+ reconnect and are replayed on the connection it establishes. With the
default `MaxRetries = 0` an
+ unreachable server is retried forever, so a request that passes no
`CancellationToken` waits for as long as
+ the server stays down - set `MaxRetries` or pass a token to bound it. A
reconnect restores the session from
+ `AutoLoginSettings` or from the sign-in a `LoginUserAsync` call succeeded
with, so a client that logged in by
+ hand reconnects too; without either there is nothing to restore and the
request fails. A server-side
+ eviction (the heartbeat verifier reacting to silence) is recovered from the
same way; only `LogoutUserAsync`
+ or `Dispose` ends a session for good. Set `ReconnectionSettings.Enabled =
false` to opt out of
+ reconnection.
- `AutoLoginSettings` properties are now `init`-only, as is
`IggyClientConfigurator.HeartbeatInterval`. Build
them with an object initializer or the `AutoLoginSettings.For` /
`AutoLoginSettings.ForPersonalAccessToken`
factories instead of assigning after construction.
diff --git a/foreign/go/client/tcp/tcp_connect_test.go
b/foreign/go/client/tcp/tcp_connect_test.go
index aa3ce3a12..12fcd615a 100644
--- a/foreign/go/client/tcp/tcp_connect_test.go
+++ b/foreign/go/client/tcp/tcp_connect_test.go
@@ -495,7 +495,18 @@ func TestExchange_DoesNotPreemptAReplayedSignIn(t
*testing.T) {
// The connect sign-in, the dropped explicit one, and its replay. An
// automatic sign-in on the new connection would make a fourth.
assert.Equal(t, 3, signIns)
- assert.False(t, client.skipAutoLoginOnce, "the suppression fires
exactly once")
+
+ // The suppression rides the replay's own context, so nothing about it
+ // outlives that call: the next Connect signs in again.
+ require.NoError(t, client.disconnect())
+ require.NoError(t, client.Connect(context.Background()))
+ signIns = 0
+ for _, recorded := range server.recorded() {
+ if recorded.operation() == vsr.OperationRegister {
+ signIns++
+ }
+ }
+ assert.Equal(t, 4, signIns, "the suppression leaked past the call that
meant it")
}
func TestExchange_FailsFastWhenAutoLoginIsOff(t *testing.T) {
diff --git a/foreign/go/client/tcp/tcp_core.go
b/foreign/go/client/tcp/tcp_core.go
index 0ebb91c31..9cad23f74 100644
--- a/foreign/go/client/tcp/tcp_core.go
+++ b/foreign/go/client/tcp/tcp_core.go
@@ -79,9 +79,6 @@ type IggyTcpClient struct {
// session carries the consensus client identity and request watermark;
// guarded by c.mtx.
session *vsr.Session
- // skipAutoLoginOnce suppresses the next automatic sign-in so a replayed
- // login is not preempted by one the reconnect issues; guarded by c.mtx.
- skipAutoLoginOnce bool
// loggedOut records an explicit sign-out, so a reconnect's automatic
// sign-in does not silently reverse it; guarded by c.mtx.
loggedOut bool
@@ -466,6 +463,19 @@ func appendCommandFrame(buf []byte, cmd command.Command)
([]byte, error) {
// deadlock on that lock or recurse Connect without a bound.
type connectScoped struct{}
+// skipAutoLogin marks the context of a Connect whose caller owns the sign-in:
+// a replayed login, or a redirect inside the sign-in transaction. Carried on
+// the context rather than on the client, so it cannot outlive the call that
+// meant it -- a client-wide flag leaks when Connect returns early on the
+// already-connected gate, and then suppresses somebody else's auto-login.
+type skipAutoLogin struct{}
+
+// suppressAutoLogin returns ctx marked so the Connect it drives does not sign
+// in by itself.
+func suppressAutoLogin(ctx context.Context) context.Context {
+ return context.WithValue(ctx, skipAutoLogin{}, struct{}{})
+}
+
// localPreconditionError marks a request that failed before its frame was
// written. The connection is healthy, so exchange must not tear it down and
// re-dial over what is purely local state.
@@ -482,18 +492,10 @@ func (c *IggyTcpClient) exchange(ctx context.Context,
code uint32, frame []byte)
return response, err
}
- // A stale-client eviction is the server ending this session
- // authoritatively, like a logout: the remembered sign-in ends with it,
so
- // only a configured auto-login may bring the session back. Remembered
- // credentials exist for transport loss, where the session died with the
- // socket rather than by anyone's decision. This runs before the gates
- // below because they all return: with reconnection disabled the
eviction
- // would otherwise never be forgotten, and the next manual Connect would
- // sign in with the evicted session's credentials.
- if errors.Is(err, ierror.ErrStaleClient) {
- c.forgetLogin()
- }
-
+ // A stale-client eviction is not caller intent: the heartbeat verifier
+ // sends it after a gc pause or a laptop sleep, so the remembered
sign-in
+ // survives it and the reconnect re-establishes the session. Only an
+ // explicit sign-out ends it. Same rule in every SDK.
var precondition *localPreconditionError
if errors.As(err, &precondition) {
return nil, err
@@ -532,22 +534,20 @@ func (c *IggyTcpClient) exchange(ctx context.Context,
code uint32, frame []byte)
if disconnectErr := c.disconnect(); disconnectErr != nil {
return nil, disconnectErr
}
+ reconnectCtx := ctx
if login {
- c.mtx.Lock()
- c.skipAutoLoginOnce = true
- c.mtx.Unlock()
+ // The caller replays the login itself, so the reconnect must
not.
+ reconnectCtx = suppressAutoLogin(ctx)
}
+ c.mtx.Lock()
+ serverAddress := c.currentServerAddress
+ c.mtx.Unlock()
c.logger.Info("Reconnecting to the server...",
- slog.String("server_address", c.currentServerAddress),
+ slog.String("server_address", serverAddress),
slog.Any("error", err))
- if reconnectErr := c.Connect(ctx); reconnectErr != nil {
- if login {
- c.mtx.Lock()
- c.skipAutoLoginOnce = false
- c.mtx.Unlock()
- }
+ if reconnectErr := c.Connect(reconnectCtx); reconnectErr != nil {
return nil, reconnectErr
}
return c.sendFrame(ctx, code, frame)
@@ -648,18 +648,15 @@ func (c *IggyTcpClient) sendFrame(ctx context.Context,
code uint32, frame []byte
// on the reconnect path would wait on that
lock forever. The
// transaction signs in itself on the node it
lands on, so the
// reconnect must not.
- connectScopedRequest :=
ctx.Value(connectScoped{}) != nil
- if connectScopedRequest {
- c.mtx.Lock()
- c.skipAutoLoginOnce = true
- c.mtx.Unlock()
+ redirectCtx := ctx
+ if ctx.Value(connectScoped{}) != nil {
+ // Issued from inside the sign-in
transaction, which holds
+ // registerMtx: the automatic sign-in
on the reconnect path
+ // would wait on that lock forever. The
transaction signs in
+ // itself on the node it lands on, so
the reconnect must not.
+ redirectCtx = suppressAutoLogin(ctx)
}
- if connectErr := c.Connect(ctx); connectErr !=
nil {
- if connectScopedRequest {
- c.mtx.Lock()
- c.skipAutoLoginOnce = false
- c.mtx.Unlock()
- }
+ if connectErr := c.Connect(redirectCtx);
connectErr != nil {
return nil, connectErr
}
stamped = false
@@ -1029,12 +1026,14 @@ func (c *IggyTcpClient) Connect(ctx context.Context)
error {
// The server fence does not survive the old socket, so the new
connection
// starts from a fresh client identity.
c.session.Reset()
- skipAutoLogin := c.skipAutoLoginOnce
- c.skipAutoLoginOnce = false
- c.logger.Info("Iggy client has connected to the Iggy server",
slog.String("client_address", c.clientAddress), slog.String("server_address",
c.currentServerAddress))
+ clientAddress := c.clientAddress
+ serverAddress := c.currentServerAddress
c.mtx.Unlock()
+ c.logger.Info("Iggy client has connected to the Iggy server",
+ slog.String("client_address", clientAddress),
+ slog.String("server_address", serverAddress))
- if err := c.establishSession(ctx, skipAutoLogin); err != nil {
+ if err := c.establishSession(ctx, ctx.Value(skipAutoLogin{}) != nil);
err != nil {
_ = c.disconnect()
return err
}
diff --git a/foreign/go/client/tcp/tcp_core_review_test.go
b/foreign/go/client/tcp/tcp_core_review_test.go
index f9deff4f7..4d0c2e5cd 100644
--- a/foreign/go/client/tcp/tcp_core_review_test.go
+++ b/foreign/go/client/tcp/tcp_core_review_test.go
@@ -136,14 +136,19 @@ func TestConnect_SuppressedSignInSendsNothing(t
*testing.T) {
client := newDialingClient(t, server.address(),
WithAutoLogin(NewUsernamePasswordCredentials("iggy", "iggy")))
- // The state a replayed login leaves behind before its reconnect.
- client.skipAutoLoginOnce = true
- require.NoError(t, client.Connect(context.Background()))
+ // The context a replayed login drives its reconnect with.
+ require.NoError(t,
client.Connect(suppressAutoLogin(context.Background())))
assert.Empty(t, server.recorded(),
"the replayed login owns the sign-in; Connect must not preempt
it")
- assert.False(t, client.skipAutoLoginOnce, "the suppression is consumed
exactly once")
+
+ // And it is that context, not the client, that carries the suppression:
+ // a Connect without it signs in.
+ require.NoError(t, client.disconnect())
+ require.NoError(t, client.Connect(context.Background()))
+ assert.NotEmpty(t, server.recorded(),
+ "the suppression outlived the call that meant it")
}
func TestClose_InterruptsAnInFlightReplayWait(t *testing.T) {
diff --git a/foreign/go/client/tcp/tcp_failover_test.go
b/foreign/go/client/tcp/tcp_failover_test.go
index 30badd8a3..020e8b2b8 100644
--- a/foreign/go/client/tcp/tcp_failover_test.go
+++ b/foreign/go/client/tcp/tcp_failover_test.go
@@ -124,17 +124,20 @@ func TestFailover_FailsFastWhenNothingEverSignedIn(t
*testing.T) {
"a client that never signed in cannot restore a session by
reconnecting")
}
-// A stale-client eviction is the server ending the session authoritatively,
-// like a logout: the remembered sign-in must not resurrect it, so the evicted
-// request surfaces the loss instead of reconnecting into a fresh session.
-func TestFailover_ServerEvictionForgetsTheRememberedSignIn(t *testing.T) {
+// A stale-client eviction is not caller intent: the heartbeat verifier sends
it
+// after a gc pause or a laptop sleep, and a client that signed in by hand has
+// to recover from it exactly like one with a configured auto-login. Same rule
+// in every SDK.
+func TestFailover_ServerEvictionReplaysTheRememberedSignIn(t *testing.T) {
var server *testListener
var evict atomic.Bool
+ var registers atomic.Int32
server = listenVSR(t, nil, func(_, _ int, read request) []byte {
if read.operation() == vsr.OperationRegister {
+ registers.Add(1)
return registerReplyFrame(7, 128)
}
- if evict.Load() {
+ if evict.CompareAndSwap(true, false) {
return evictionFrame(vsr.EvictionStaleClient, 0, 0)
}
if read.code() == uint32(command.GetClusterMetadataCode) {
@@ -148,15 +151,19 @@ func
TestFailover_ServerEvictionForgetsTheRememberedSignIn(t *testing.T) {
require.NoError(t, client.Connect(ctx))
_, err := client.LoginUser(ctx, "iggy", "iggy")
require.NoError(t, err)
- connectionsBefore := server.connections()
+ registersBefore := registers.Load()
evict.Store(true)
- require.Error(t, client.Ping(ctx), "the evicted request surfaces the
loss")
+ // The evicted request is answered by the eviction, and the reconnect it
+ // triggers signs in again with the credentials the sign-in remembered.
+ _ = client.Ping(ctx)
_, remembered := client.signInCredentials()
- assert.False(t, remembered, "the eviction forgot the remembered
sign-in")
- assert.Equal(t, connectionsBefore, server.connections(),
- "no reconnect dial resurrected the evicted session")
+ assert.True(t, remembered, "an eviction is not a sign-out; the
credentials stay")
+ require.NoError(t, client.Ping(ctx), "the session came back on its own")
+ assert.Greater(t, registers.Load(), registersBefore,
+ "the reconnect re-established the session")
+ assert.True(t, client.session.Bound())
}
// An explicit sign-out is caller intent: the reconnect must not sign back in
diff --git a/foreign/go/client/tcp/tcp_session_management.go
b/foreign/go/client/tcp/tcp_session_management.go
index d2e534116..4da1fb23b 100644
--- a/foreign/go/client/tcp/tcp_session_management.go
+++ b/foreign/go/client/tcp/tcp_session_management.go
@@ -76,7 +76,10 @@ func (c *IggyTcpClient) register(
c.registerMtx.Lock()
defer c.registerMtx.Unlock()
- c.logger.Info("Iggy client is signing in...",
slog.String("client_address", c.clientAddress))
+ c.mtx.Lock()
+ clientAddress := c.clientAddress
+ c.mtx.Unlock()
+ c.logger.Info("Iggy client is signing in...",
slog.String("client_address", clientAddress))
if err := c.endBoundSession(ctx); err != nil {
return nil, err
@@ -136,8 +139,11 @@ func (c *IggyTcpClient) signIn(ctx context.Context, code
uint32, body []byte) (*
return nil, err
}
+ c.mtx.Lock()
+ signedInAddress := c.clientAddress
+ c.mtx.Unlock()
c.logger.Info("Iggy client has signed in successfully.",
- slog.String("client_address", c.clientAddress),
+ slog.String("client_address", signedInAddress),
slog.String("server_version", registered.ServerVersion))
return &iggcon.IdentityInfo{UserId: registered.UserID}, nil
}
@@ -169,13 +175,7 @@ func (c *IggyTcpClient) settleOnLeader(ctx
context.Context, code uint32, body []
// The replayed sign-in below owns the session; the redirected
Connect
// must not sign in on its own, or the replay commits a second
Register.
- c.mtx.Lock()
- c.skipAutoLoginOnce = true
- c.mtx.Unlock()
- if err := c.Connect(ctx); err != nil {
- c.mtx.Lock()
- c.skipAutoLoginOnce = false
- c.mtx.Unlock()
+ if err := c.Connect(suppressAutoLogin(ctx)); err != nil {
return nil, err
}
settled, err = c.signIn(ctx, code, body)
diff --git
a/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java
b/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java
index ef2baa5ed..2bec95acb 100644
---
a/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java
+++
b/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java
@@ -500,42 +500,21 @@ public class AsyncIggyTcpClient {
/**
* A server-side eviction reached this client. The routing state it cached
- * belonged to the evicted session, and a stale-client eviction is the
- * server ending that session authoritatively, like a logout: nothing may
- * revive it behind the caller's back.
+ * belonged to the evicted session, so it goes.
*
- * The captured login always goes, whatever is configured. The connection
- * replays it to bring up a replacement channel, so keeping it would revive
- * exactly the session the server ended -- and on a client whose caller
- * signed in as somebody other than the configured user, it would revive a
- * different user than a redial replays, making the outcome depend on which
- * failure ran.
- *
- * What may re-establish the session is what every connect of this client
- * signs in as: credentials configured on the builder. A sign-in the caller
- * ran is theirs to repeat.
+ * The sign-in stays. A stale-client eviction is not caller intent: the
+ * server's heartbeat verifier sends it after a gc pause or a laptop sleep,
+ * and a client that signed in by hand has to recover from it exactly like
+ * one whose credentials were configured. The connection re-authenticates
+ * the replacement channel from the login it captured, which is the same
+ * sign-in a redial would replay. Same rule in every SDK; only an explicit
+ * sign-out or close ends a session for good.
*/
private void onSessionReset(int errorCode) {
routingState.clearAssignments();
- if (errorCode != IggyErrorCode.STALE_CLIENT.getCode()) {
- return;
- }
- AsyncTcpConnection currentConnection = connection.get();
- if (currentConnection != null) {
- currentConnection.forgetCapturedLogin();
+ if (errorCode == IggyErrorCode.STALE_CLIENT.getCode()) {
+ log.debug("The server evicted this session as stale; the next
request re-establishes it");
}
- if (username.isEmpty() || password.isEmpty()) {
- log.warn("The server evicted this session as stale; the sign-in it
ran will not be replayed");
- rememberedLogin = null;
- return;
- }
-
- log.info("The server evicted this session as stale; signing in again
with the configured credentials");
- replayLogin().whenComplete((ignored, error) -> {
- if (error != null) {
- log.warn("Signing in again after the eviction failed: {}",
error.getMessage());
- }
- });
}
/**
@@ -660,11 +639,18 @@ public class AsyncIggyTcpClient {
if (closed) {
return CompletableFuture.completedFuture(null);
}
- if (attempt > policy.getMaxRetries()) {
+ List<ConnectionInfo> candidates = redialCandidates();
+ // The retry budget bounds the rotations, not the endpoints. A policy
of
+ // zero retries with several endpoints known still gets one rotation:
+ // those endpoints - the address the client was configured with, the
+ // nodes the roster named - were made known in order to be tried, and
+ // the other SDKs sweep them once too. With one endpoint known, zero
+ // retries redials nothing, which is what it asked for.
+ boolean sweepOnce = attempt == 1 && candidates.size() > 1;
+ if (attempt > policy.getMaxRetries() && !sweepOnce) {
log.error("Redial gave up after {} attempts, next request will
fail fast", policy.getMaxRetries());
return CompletableFuture.completedFuture(null);
}
- List<ConnectionInfo> candidates = redialCandidates();
// The delay paces rotations, not dials. The first rotation runs at
once
// when there is somewhere else to go: pausing before dialing a
survivor
// only pushes the failover past the window the caller waits in, and
the
@@ -767,27 +753,33 @@ public class AsyncIggyTcpClient {
}
/**
- * Re-establishes the session on the freshly published connection: with the
- * credentials configured on the builder, or else the sign-in a caller ran
- * by hand. Configured credentials win, as in the Rust and Go SDKs -- they
- * are what every connect of this client is meant to sign in as, and a
- * remembered sign-in exists to make a hand-run login as reconnectable as
- * a configured one, not to override it.
+ * Re-establishes the session on the freshly published connection with the
+ * sign-in that last succeeded, falling back to the credentials configured
+ * on the builder when no login has run yet.
+ *
+ * The last sign-in rather than the configured one, which is where this
+ * differs from the Rust and Go SDKs: the connection re-authenticates a
+ * replacement channel from the login payload it captured, which is that
+ * same last sign-in. Replaying a different user here would make the same
+ * eviction land on a different session depending on whether the channel
+ * or the redial got there first. A client that only ever used its
+ * configured credentials remembers exactly those, so nothing changes for
+ * it.
*
* The login runs through the users client, so leader discovery retargets
* again before Register when the redialed node is not the leader.
*/
private CompletableFuture<Void> replayLogin() {
+ Supplier<CompletableFuture<IdentityInfo>> replay = rememberedLogin;
+ if (replay != null) {
+ // Runs through loginOnLeader, so a redial that landed on a backup
+ // still settles on the leader before the session is used.
+ return loginOnLeader(replay).thenApply(identity -> null);
+ }
if (username.isPresent() && password.isPresent() && usersClient !=
null) {
return usersClient.login(username.get(),
password.get()).thenApply(identity -> null);
}
- Supplier<CompletableFuture<IdentityInfo>> replay = rememberedLogin;
- if (replay == null) {
- return CompletableFuture.completedFuture(null);
- }
- // Runs through loginOnLeader, so a redial that landed on a backup
- // still settles on the leader before the session is used.
- return loginOnLeader(replay).thenApply(identity -> null);
+ return CompletableFuture.completedFuture(null);
}
/**
diff --git
a/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncTcpConnection.java
b/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncTcpConnection.java
index 94515314c..71a0bca73 100644
---
a/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncTcpConnection.java
+++
b/foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncTcpConnection.java
@@ -841,21 +841,6 @@ public class AsyncTcpConnection {
sessionResetListener.accept(errorCode);
}
- /**
- * Drops the captured sign-in, so the next channel comes up
- * unauthenticated instead of replaying it.
- *
- * The channel replays the payload it captured to re-authenticate a
- * replacement channel, which is right for a lost connection and wrong
- * after an eviction the server decided on: that would resurrect the very
- * session the server ended.
- */
- void forgetCapturedLogin() {
- authenticated = false;
- authGeneration.incrementAndGet();
- releaseLoginPayload();
- }
-
private void captureLoginPayloadIfNeeded(int commandCode, ByteBuf payload)
{
if (isLoginCode(commandCode)) {
updateLoginPayload(commandCode, payload);
diff --git
a/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClientTransientFailoverTest.java
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClientTransientFailoverTest.java
index d339cebf7..f67405969 100644
---
a/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClientTransientFailoverTest.java
+++
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClientTransientFailoverTest.java
@@ -213,13 +213,14 @@ class AsyncIggyTcpClientTransientFailoverTest {
}
/**
- * The user an eviction re-establishes must not depend on which failure
- * ran. Configured credentials are what every connect of this client signs
- * in as, so they are what comes back -- not whoever the caller signed in
- * as by hand, which is what the connection's captured login would replay.
+ * A stale-client eviction is not caller intent: the server's heartbeat
+ * verifier sends it after a gc pause or a laptop sleep. A client that
+ * signed in by hand recovers from it exactly like one whose credentials
+ * were configured (the test above), and the sign-in it recovers with is
the
+ * one that last succeeded. Same rule in every SDK.
*/
@Test
- void shouldSignInAsTheConfiguredUserAfterAnEviction() throws Exception {
+ void shouldReviveTheSignInAfterAStaleClientEviction() throws Exception {
InetAddress loopback = InetAddress.getLoopbackAddress();
try (ServerSocket serverSocket = new ServerSocket(0, 4, loopback)) {
AtomicInteger registrations = new AtomicInteger();
@@ -233,69 +234,13 @@ class AsyncIggyTcpClientTransientFailoverTest {
registeredLogins.add(request.bodyAsText());
return Response.success(OPERATION_REGISTER,
registerBody(registrations.incrementAndGet()));
}
- if (request.operation() == OPERATION_CREATE_STREAM &&
evict.compareAndSet(true, false)) {
- return Response.eviction(EVICTION_STALE_CLIENT);
- }
- return Response.success(request.operation(),
Unpooled.EMPTY_BUFFER);
- });
-
- AsyncIggyTcpClient client = AsyncIggyTcpClient.builder()
- .host(loopback.getHostAddress())
- .port(serverSocket.getLocalPort())
- .credentials("configured", "configured")
- .requestTimeout(Duration.ofSeconds(2))
- .build();
- try {
- client.connect().get(5, TimeUnit.SECONDS);
- client.login().get(5, TimeUnit.SECONDS);
- client.users().login("handrun", "handrun").get(5,
TimeUnit.SECONDS);
- int registrationsBeforeEviction = registrations.get();
-
- assertThatThrownBy(() ->
client.sendBinaryRequest(CREATE_STREAM_CODE, new byte[0])
- .get(5, TimeUnit.SECONDS))
- .hasCauseInstanceOf(IggyServerException.class);
-
- // The sign-in that follows the eviction runs on its own, so
- // give it a moment to land.
- long deadline = System.nanoTime() +
Duration.ofSeconds(5).toNanos();
- while (registrations.get() == registrationsBeforeEviction &&
System.nanoTime() < deadline) {
- Thread.sleep(25);
- }
- assertThat(registrations.get())
- .as("the eviction was not followed by a sign-in")
- .isGreaterThan(registrationsBeforeEviction);
- assertThat(registeredLogins.get(registeredLogins.size() - 1))
- .as("the eviction revived the hand-run sign-in instead
of the configured one")
- .contains("configured")
- .doesNotContain("handrun");
- } finally {
- client.close().get(5, TimeUnit.SECONDS);
- }
- server.completeExceptionally(new IllegalStateException("test
over"));
- }
- }
-
- /**
- * A stale-client eviction is the server ending the session
- * authoritatively, like a logout. A client whose credentials were
- * configured signs in again on every connect and recovers (the test
- * above); one whose session came from a caller's own sign-in must not have
- * that session revived behind the caller's back.
- */
- @Test
- void shouldNotReviveAHandRunSignInAfterAStaleClientEviction() throws
Exception {
- InetAddress loopback = InetAddress.getLoopbackAddress();
- try (ServerSocket serverSocket = new ServerSocket(0, 4, loopback)) {
- AtomicInteger registrations = new AtomicInteger();
- CompletableFuture<Void> server = serve(serverSocket, 4, request ->
{
- if (request.is(GET_CLUSTER_METADATA_CODE,
OPERATION_NON_REPLICATED)) {
- return Response.success(OPERATION_NON_REPLICATED,
singleNodeMetadata(serverSocket.getLocalPort()));
- }
- if (request.operation() == OPERATION_REGISTER) {
- return Response.success(OPERATION_REGISTER,
registerBody(registrations.incrementAndGet()));
- }
if (request.operation() == OPERATION_CREATE_STREAM) {
- return Response.eviction(EVICTION_STALE_CLIENT);
+ if (evict.compareAndSet(true, false)) {
+ return Response.eviction(EVICTION_STALE_CLIENT);
+ }
+ ByteBuf body = Unpooled.buffer(Integer.BYTES);
+ body.writeIntLE(0);
+ return Response.success(OPERATION_CREATE_STREAM, body);
}
return Response.success(request.operation(),
Unpooled.EMPTY_BUFFER);
});
@@ -308,27 +253,26 @@ class AsyncIggyTcpClientTransientFailoverTest {
.build();
try {
client.connect().get(5, TimeUnit.SECONDS);
- client.users().login("iggy", "iggy").get(5, TimeUnit.SECONDS);
+ client.users().login("handrun", "handrun").get(5,
TimeUnit.SECONDS);
int registrationsBeforeEviction = registrations.get();
assertThatThrownBy(() ->
client.sendBinaryRequest(CREATE_STREAM_CODE, new byte[0])
.get(5, TimeUnit.SECONDS))
.hasCauseInstanceOf(IggyServerException.class);
- // Whatever the caller does next, nothing may sign this session
- // back in on its own.
- assertThatThrownBy(() ->
client.sendBinaryRequest(CREATE_STREAM_CODE, new byte[0])
+ assertThat(client.hasRememberedLogin())
+ .as("an eviction is not a sign-out; the credentials
stay")
+ .isTrue();
+
+ // The next request brings the session back, under the sign-in
+ // that last succeeded.
+ assertThat(client.sendBinaryRequest(CREATE_STREAM_CODE, new
byte[0])
.get(5, TimeUnit.SECONDS))
- .isNotNull();
- assertThat(registrations)
- .as("the evicted session was signed back in")
- .hasValue(registrationsBeforeEviction);
- // The connection replays the login it captured to bring up a
- // replacement channel, so that copy has to go too: kept, the
- // next channel revives exactly the session the server ended.
- assertThat(client.currentConnection().authenticationSnapshot())
- .as("the captured login outlived the session it
established")
.isEmpty();
+ assertThat(registrations.get())
+ .as("the evicted session was not re-established")
+ .isGreaterThan(registrationsBeforeEviction);
+ assertThat(registeredLogins.get(registeredLogins.size() -
1)).contains("handrun");
} finally {
client.close().get(5, TimeUnit.SECONDS);
}
diff --git a/foreign/node/src/client/client.connection.ts
b/foreign/node/src/client/client.connection.ts
index ea5c7dcbd..0226147e1 100644
--- a/foreign/node/src/client/client.connection.ts
+++ b/foreign/node/src/client/client.connection.ts
@@ -381,7 +381,14 @@ export class IggyConnection extends EventEmitter {
let lastError = initialError;
let expectedSocket = this.socket;
let firstPass = true;
- while (enabled && this.reconnectCount < maxRetries) {
+ // Reconnection settings bound the retries, not the endpoints. With them
off
+ // and several endpoints known - the address the client was configured
with,
+ // the nodes the roster named - those endpoints were made known in order to
+ // be tried, so they get one pass and no backoff, as in the other SDKs. A
+ // client that knows one endpoint and turned reconnection off redials
+ // nothing, which is what it asked for.
+ const sweepOnce = !enabled && this._redialCandidates().length > 1;
+ while (this.reconnectCount < maxRetries || (sweepOnce && firstPass)) {
this.connecting = true;
this.reconnectCount += 1;
const candidates = this._redialCandidates();
@@ -389,7 +396,7 @@ export class IggyConnection extends EventEmitter {
// endpoints known there is somewhere else to go, and pausing first only
// pushes the failover past the interval a caller is willing to wait;
// later passes still back off.
- if (!firstPass || candidates.length === 1)
+ if (enabled && (!firstPass || candidates.length === 1))
await waitForReconnect(interval);
firstPass = false;
if (this.ending)
diff --git a/foreign/node/src/client/client.socket.test.ts
b/foreign/node/src/client/client.socket.test.ts
index e938c7603..2da50fdcd 100644
--- a/foreign/node/src/client/client.socket.test.ts
+++ b/foreign/node/src/client/client.socket.test.ts
@@ -209,6 +209,12 @@ const vsrConfig = (port: number): ClientConfig => ({
});
/** Shrinks the leaderless poll so a test observes it without waiting on it. */
+/** Drives the leader re-check a refused request runs. */
+const followLeaderMove = (client: CommandResponseStream): Promise<boolean> =>
+ (client as unknown as {
+ _followLeaderMove: () => Promise<boolean>
+ })._followLeaderMove();
+
const compressLeaderlessPoll = (
client: CommandResponseStream,
budget: number
@@ -969,6 +975,72 @@ describe('VSR client socket', () => {
}
});
+ it('re-issues every request refused by a demoted node, not just the first',
+ async () => {
+ // One demotion, several refused requests: each of them re-checking on
+ // its own would move the client once per request, and the first
+ // redirect's drop would fail the others' roster reads - reporting a
+ // refusal they never had to.
+ const leader = await startVsrServer((frame, socket) => {
+ singleNodeHandler(leader.port)(frame, socket);
+ });
+ // Leader at login, so the settlement leaves the client here, then
+ // demoted: the refusals below are what tells the client to look again.
+ let demotedYet = false;
+ const demoted = await startVsrServer((frame, socket) => {
+ const code = frame.readUInt32LE(REQUEST_OFFSET.reserved);
+ if (code === COMMAND_CODE.GetClusterMetadata) {
+ socket.write(replyFrame(
+ Operation.NonReplicated,
+ demotedYet
+ ? twoNodeMetadataBody(demoted.port, leader.port)
+ : twoNodeMetadataBody(leader.port, demoted.port)
+ ));
+ return;
+ }
+ if (code === 60_032) {
+ socket.write(replyFrame(Operation.NonReplicated, Buffer.alloc(0),
58));
+ return;
+ }
+ singleNodeHandler(demoted.port)(frame, socket);
+ });
+
+ const client = new CommandResponseStream(vsrConfig(demoted.port));
+ try {
+ await client.authenticate(vsrConfig(demoted.port).credentials);
+ demotedYet = true;
+ const rosterReadsBefore = demoted.frames.filter(
+ (frame) => frame.readUInt32LE(REQUEST_OFFSET.reserved) ===
+ COMMAND_CODE.GetClusterMetadata
+ ).length;
+
+ // Two refusals, one re-check: the second caller shares the move the
+ // first started rather than starting its own or being told 58.
+ const moves = await Promise.all([
+ followLeaderMove(client),
+ followLeaderMove(client)
+ ]);
+
+ assert.deepEqual(moves, [true, true]);
+ const rosterReads = demoted.frames.filter(
+ (frame) => frame.readUInt32LE(REQUEST_OFFSET.reserved) ===
+ COMMAND_CODE.GetClusterMetadata
+ ).length - rosterReadsBefore;
+ assert.equal(rosterReads, 1,
+ 'each refusal re-read the roster on its own'
+ );
+ const connection = (client as unknown as {
+ connection: { isConnectedTo: (host: string, port: number) => boolean
}
+ }).connection;
+ assert.equal(connection.isConnectedTo('127.0.0.1', leader.port), true);
+ } finally {
+ client.destroy();
+ await leader.close();
+ await demoted.close();
+ }
+ }
+ );
+
it('keeps re-issuing a not-admitted request while the roster still names
this node',
async () => {
const server = await startVsrServer((frame, socket) => {
diff --git a/foreign/node/src/client/client.socket.ts
b/foreign/node/src/client/client.socket.ts
index 43d378b84..225141875 100644
--- a/foreign/node/src/client/client.socket.ts
+++ b/foreign/node/src/client/client.socket.ts
@@ -119,11 +119,10 @@ export class CommandResponseStream extends EventEmitter {
/** Whether a login is already being moved to the leader */
private settlingLeader: boolean;
/**
- * Whether a refused request is already re-checking the leader. The roster
- * read that re-check runs can be refused the same way, and answering a
- * leader check with another leader check would recurse.
+ * The leader re-check a refused request started, shared with every other
+ * request refused by the same node so one demotion moves the client once.
*/
- private followingLeaderMove = false;
+ private leaderMoveInFlight?: Promise<boolean>;
/** How long a leaderless roster is polled before settling in place */
private leaderlessWaitBudget: number;
/** Delay between roster reads while the cluster elects */
@@ -202,7 +201,8 @@ export class CommandResponseStream extends EventEmitter {
try {
const {
handleResponse = true,
- last = true
+ last = true,
+ followsLeaderMoves = true
} = options;
if (!this.connection.connected)
@@ -222,6 +222,7 @@ export class CommandResponseStream extends EventEmitter {
// one window: the roster can still name this node -- an election in
// flight, a leader that has not moved yet -- and that is a wait, not a
// verdict.
+ //
// One budget for the whole request: the transient replays on a
// connection, the leader re-checks, and the re-issues after a move all
// spend it, so a request cannot outlive it by moving.
@@ -235,12 +236,24 @@ export class CommandResponseStream extends EventEmitter {
} catch (error) {
if (!(error instanceof LeaderMovedError))
throw error;
- // The roster read itself is refused: it runs through this same
- // path, and re-checking the leader to answer a leader check would
- // recurse. Its caller reads a failure as "stay where you are".
- if (this.followingLeaderMove || Date.now() >= deadline)
+ // The roster read that a re-check runs is itself a command that can
+ // be refused this way, and answering a leader check with another
+ // leader check would recurse. Its caller reads a failure as "stay
+ // where you are".
+ if (!followsLeaderMoves || Date.now() >= deadline)
throw responseError(command, error.refusal.errorCode);
- await this._followLeaderMove();
+
+ const moved = await this._followLeaderMove();
+ if (!moved) {
+ // Nowhere else to go yet: the roster still names this node, or it
+ // could not be read. Paced, because the in-connection replay
+ // window belongs to the request's budget and has already been
+ // spent -- re-issuing straight away would spin.
+ const remaining = deadline - Date.now();
+ if (remaining <= 0)
+ throw responseError(command, error.refusal.errorCode);
+ await delay(Math.min(VSR_FAILOVER_CHECK_MS, remaining));
+ }
// A move drops the session with the socket it was bound to, so the
// re-issue would otherwise go out under no session: a replicated
// command fails client-side, a non-replicated one goes out with
@@ -292,29 +305,40 @@ export class CommandResponseStream extends EventEmitter {
* Re-reads the roster and moves to the leader it names.
*
* Best effort: an unreadable roster, or one that still names this node,
- * leaves the client where it is and the refused request is re-issued here
- * anyway. Guarded against re-entry, since the roster read is itself a
- * command that can be refused the same way.
+ * leaves the client where it is and the refused request is re-issued anyway.
+ *
+ * Single-flighted, and concurrent callers share the outcome instead of
+ * failing: several commands are refused by the same demoted node, and each
+ * starting its own redirect would move the client once per command. The
+ * first redirect's `'disconnected'` also fails the others' roster reads, so
+ * a caller that raced one would report a refusal it never had to.
*
* @returns Whether the client moved
*/
- private async _followLeaderMove(): Promise<boolean> {
- if (this.followingLeaderMove)
- return false;
- this.followingLeaderMove = true;
- try {
- const leader = await this._readLeaderEndpoint();
- if (!leader || this.connection.isConnectedTo(leader.host, leader.port))
+ private _followLeaderMove(): Promise<boolean> {
+ const inFlight = this.leaderMoveInFlight;
+ if (inFlight)
+ return inFlight;
+
+ const move = (async () => {
+ try {
+ const leader = await this._readLeaderEndpoint();
+ if (!leader || this.connection.isConnectedTo(leader.host, leader.port))
+ return false;
+ debug(`the leader moved to ${leader.host}:${leader.port}, following
it`);
+ await this.connection.redirect(leader.host, leader.port);
+ return true;
+ } catch (error) {
+ debug('the leader could not be re-checked, staying on this node',
error);
return false;
- debug(`the leader moved to ${leader.host}:${leader.port}, following it`);
- await this.connection.redirect(leader.host, leader.port);
- return true;
- } catch (error) {
- debug('the leader could not be re-checked, staying on this node', error);
- return false;
- } finally {
- this.followingLeaderMove = false;
- }
+ }
+ })();
+ this.leaderMoveInFlight = move;
+ void move.finally(() => {
+ if (this.leaderMoveInFlight === move)
+ this.leaderMoveInFlight = undefined;
+ });
+ return move;
}
private _rememberRoster(response: CommandResponse): void {
@@ -617,7 +641,7 @@ export class CommandResponseStream extends EventEmitter {
const response = await this.sendCommand(
GET_CLUSTER_METADATA.code,
GET_CLUSTER_METADATA.serialize(),
- { last: false }
+ { last: false, followsLeaderMoves: false }
);
// The redial candidates are fed by `_processVsr` for every roster
// read, leaderless ones included: a roster with no leader still names
diff --git a/foreign/node/src/client/client.type.ts
b/foreign/node/src/client/client.type.ts
index 63703005f..bb9f89f0f 100644
--- a/foreign/node/src/client/client.type.ts
+++ b/foreign/node/src/client/client.type.ts
@@ -47,7 +47,13 @@ export type SendCommandOptions = {
/** Whether the response uses the standard command response decoder */
handleResponse?: boolean,
/** Whether to append rather than prepend the command to the queue */
- last?: boolean
+ last?: boolean,
+ /**
+ * Whether a not-admitted refusal re-checks the leader and re-issues the
+ * command. False for the roster read a re-check itself runs: answering a
+ * leader check with another leader check would recurse.
+ */
+ followsLeaderMoves?: boolean
};
/**