hubcio commented on code in PR #3944:
URL: https://github.com/apache/iggy/pull/3944#discussion_r3853696586


##########
core/sdk/src/tcp/tcp_client.rs:
##########
@@ -485,36 +546,61 @@ impl TcpClient {
             };
 
             // Handle auto-login
-            let should_redirect = match &self.config.auto_login {
-                AutoLogin::Disabled => {
-                    info!("Automatic sign-in is disabled.");
+            let should_redirect = match self.sign_in_credentials().await {
+                None => {
+                    info!("No credentials to sign in with.");
                     // Only `IggyClient` redirects after a manual sign-in, so
                     // a raw transport can stay on a backup: its first
                     // replicated write gets `TransientNotAccepted`, the
                     // redirect drops the session, and the retry fails
                     // `Unauthenticated` until the caller signs in again.
                     false
                 }
-                AutoLogin::Enabled(credentials) => {
+                Some(credentials) => {
                     if skip_auto_login {
                         info!("Skipping automatic sign-in for a retried 
login/register request.");
                         false
                     } else {
                         info!("{NAME} client: {client_address} is signing 
in...");
                         self.set_state(ClientState::Authenticating).await;
-                        match credentials {
-                            Credentials::UsernamePassword(username, password) 
=> {
-                                self.login_user(username, 
password.expose_secret()).await?;
-                                info!(
-                                    "{NAME} client: {client_address} has 
signed in with the user credentials, username: {username}",
-                                );
+                        let signed_in = match &credentials {
+                            Credentials::UsernamePassword(username, password) 
=> self
+                                .login_user(username, password.expose_secret())
+                                .await
+                                .map(|_| format!("the user credentials, 
username: {username}")),
+                            Credentials::PersonalAccessToken(token) => self
+                                
.login_with_personal_access_token(token.expose_secret())
+                                .await
+                                .map(|_| "a personal access token".to_owned()),
+                        };
+                        match signed_in {
+                            Ok(how) => {
+                                info!("{NAME} client: {client_address} has 
signed in with {how}.")
                             }
-                            Credentials::PersonalAccessToken(token) => {
-                                
self.login_with_personal_access_token(token.expose_secret())
-                                    .await?;
-                                info!(
-                                    "{NAME} client: {client_address} has 
signed in with a personal access token.",
-                                );
+                            Err(error) => {
+                                // The transport is up and only the session is
+                                // not, so the state has to say so: left at
+                                // `Authenticating` every gated operation fails
+                                // client-side with `Disconnected`, `connect()`
+                                // returns ok without dialing, and nothing 
short
+                                // of an explicit `login_user` recovers.
+                                self.set_state(ClientState::Connected).await;

Review Comment:
   this also runs when the sign-in failed because the socket died: state goes 
to `Connected` with no stream, `connect()` is a no-op and every gated op fails 
until an explicit `disconnect()`. only reset to `Connected` if the stream 
survived.



##########
core/sdk/src/tcp/tcp_client.rs:
##########
@@ -573,7 +667,205 @@ impl TcpClient {
         }
     }
 
-    async fn disconnect(&self) -> Result<(), IggyError> {
+    /// Whether an `AutoLogin` is configured on this client, which makes the
+    /// session after any connect the configured user's rather than whoever
+    /// signed in by hand.
+    pub(crate) fn auto_login_configured(&self) -> bool {
+        matches!(self.config.auto_login, AutoLogin::Enabled(_))
+    }
+
+    /// Credentials to sign in with after connecting: the configured ones, or
+    /// 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.
+    async fn sign_in_credentials(&self) -> Option<Credentials> {
+        match &self.config.auto_login {
+            AutoLogin::Enabled(credentials) => Some(credentials.clone()),
+            AutoLogin::Disabled => self
+                .session_credentials
+                .lock()
+                .await
+                .as_ref()
+                .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.
+    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()) {
+            if !candidates
+                .iter()
+                .any(|candidate| is_same_address(candidate, endpoint))
+            {
+                candidates.push(endpoint.clone());
+            }
+        }
+        candidates
+    }
+
+    /// Bring one endpoint all the way up: TCP connect, socket options, and the
+    /// TLS handshake when it is configured. Nothing about the connection is
+    /// recorded until this succeeds, so a half-usable endpoint leaves no trace
+    /// for the next connect to lead with.
+    ///
+    /// `CannotEstablishConnection` means this endpoint failed and the next one
+    /// is worth trying; any other error is a configuration fault that no
+    /// endpoint can satisfy.
+    async fn establish(&self, server_address: &str) -> 
Result<EstablishedConnection, IggyError> {
+        let stream = TcpStream::connect(server_address).await.map_err(|error| {
+            error!("Failed to connect to server: {server_address}. Error: 
{error}");
+            IggyError::CannotEstablishConnection
+        })?;
+        let client_address = stream.local_addr().map_err(|error| {
+            error!("Failed to get the local address of the client: {error}");
+            IggyError::CannotEstablishConnection
+        })?;
+        let remote_address = stream.peer_addr().map_err(|error| {
+            error!("Failed to get the remote address of the server: {error}");
+            IggyError::CannotEstablishConnection
+        })?;
+
+        if let Err(error) = stream.set_nodelay(self.config.nodelay) {
+            error!("Failed to set the nodelay option on the client: {error}, 
continuing...");
+        }
+
+        if !self.config.tls_enabled {
+            return Ok(EstablishedConnection {
+                stream: 
ConnectionStreamKind::Tcp(TcpConnectionStream::new(client_address, stream)),
+                client_address,
+                remote_address,
+            });
+        }
+
+        let _ = 
rustls::crypto::aws_lc_rs::default_provider().install_default();
+        let config = if self.config.tls_validate_certificate {
+            let mut root_cert_store = rustls::RootCertStore::empty();
+            if let Some(certificate_path) = &self.config.tls_ca_file {
+                for cert in 
CertificateDer::pem_file_iter(certificate_path).map_err(|error| {
+                    error!("Failed to read the CA file: {certificate_path}. 
{error}");
+                    IggyError::InvalidTlsCertificatePath
+                })? {
+                    let certificate = cert.map_err(|error| {
+                        error!(
+                            "Failed to read a certificate from the CA file: 
{certificate_path}. {error}",
+                        );
+                        IggyError::InvalidTlsCertificate
+                    })?;
+                    root_cert_store.add(certificate).map_err(|error| {
+                        error!(
+                            "Failed to add a certificate to the root 
certificate store. {error}"
+                        );
+                        IggyError::InvalidTlsCertificate
+                    })?;
+                }
+            } else {
+                
root_cert_store.extend(webpki_roots::TLS_SERVER_ROOTS.iter().cloned());
+            }
+
+            rustls::ClientConfig::builder()
+                .with_root_certificates(root_cert_store)
+                .with_no_client_auth()
+        } else {
+            use crate::tcp::tcp_tls_verifier::NoServerVerification;
+            rustls::ClientConfig::builder()
+                .dangerous()
+                
.with_custom_certificate_verifier(Arc::new(NoServerVerification))
+                .with_no_client_auth()
+        };
+
+        let connector = TlsConnector::from(Arc::new(config));
+        let tls_domain = if self.config.tls_domain.is_empty() {
+            // Extract hostname/IP from server_address when tls_domain is not 
specified
+            server_address
+                .split(':')
+                .next()
+                .unwrap_or(server_address)
+                .to_string()
+        } else {
+            self.config.tls_domain.to_owned()
+        };
+        let domain = ServerName::try_from(tls_domain).map_err(|error| {
+            error!("Failed to create a server name from the domain. {error}");
+            IggyError::InvalidTlsDomain
+        })?;
+        let stream = connector.connect(domain, stream).await.map_err(|error| {
+            error!("Failed to establish a TLS connection to the server: 
{error}");
+            IggyError::CannotEstablishConnection

Review Comment:
   with a single endpoint and `max_retries = None` a wrong ca, a name mismatch 
or a plaintext peer now retries forever every second - before it returned at 
once. map rustls `InvalidCertificate` to a config fault so it aborts like the 
path/domain errors.



##########
core/sdk/src/tcp/tcp_client.rs:
##########
@@ -159,16 +204,38 @@ impl BinaryTransport for TcpClient {
             return Err(IggyError::Disconnected);
         }
 
-        if matches!(self.config.auto_login, AutoLogin::Disabled) && 
!is_login_register_code(code) {
-            // Without auto-login a reconnect cannot re-establish the session,
-            // so non-login requests fail fast. Login/register itself is the
+        if !is_login_register_code(code) && 
self.sign_in_credentials().await.is_none() {
+            // With no credentials -- neither configured nor remembered from a
+            // sign-in -- a reconnect cannot re-establish the session, so
+            // non-login requests fail fast. Login/register itself is the
             // exception: the server stays deliberately silent on transient
             // register failures (the server `surface_login_failure`) and
             // relies on the client timing out and replaying the request.
             return Err(error);
         }
 
-        self.disconnect().await?;
+        // Reconnecting heals the transport, but replaying the request over the
+        // new connection is a second attempt under a new session:
+        // `reset_vsr_session` drops the client id the server's dedup fence is
+        // keyed on, so a replicated write that committed before its reply was
+        // lost would apply a second time. Replay only what provably never
+        // reached the log -- the errors raised before the request was written,
+        // and the operations that never enter it.
+        //
+        // Login and register are the exception: the server stays deliberately
+        // silent on a transient register failure and relies on the client
+        // replaying, so that replay is the protocol rather than a retry.
+        let replay_after_reconnect = is_login_register_code(code)
+            || matches!(
+                error,
+                IggyError::NotConnected
+                    | IggyError::CannotEstablishConnection

Review Comment:
   `Operation::Logout` is not in this set, so a logout over a dropped socket is 
not replayed: `logout_before_relogin` fails `login_user(B)` on the first call 
and the reconnect re-signs the old user. a logout under a fresh session can't 
double-apply anything - treat it as replayable, like go does.



##########
core/sdk/src/leader_aware.rs:
##########
@@ -162,13 +217,41 @@ fn process_cluster_metadata(
 
 /// Check if two addresses refer to the same endpoint
 /// Handles various formats like 127.0.0.1:8090 vs localhost:8090
-fn is_same_address(addr1: &str, addr2: &str) -> bool {
+///
+/// A host name and the address it resolves to are one endpoint too: a client
+/// configured as `iggy-server:8090` whose roster advertises `10.0.0.5:8090`
+/// would otherwise dial that node twice per failover sweep, and a single-node
+/// deployment would be treated as a cluster. Resolution is the last resort,
+/// only when the spellings differ and at least one side is not a literal
+/// address, and it is a blocking lookup: this runs on the connect and redirect
+/// paths, which are rare and already wait on the network.
+pub(crate) fn is_same_address(addr1: &str, addr2: &str) -> bool {
     match (parse_address(addr1), parse_address(addr2)) {
         (Some(sock1), Some(sock2)) => sock1.ip() == sock2.ip() && sock1.port() 
== sock2.port(),
-        _ => normalize_address(addr1) == normalize_address(addr2),
+        (parsed1, parsed2) => {
+            if normalize_address(addr1) == normalize_address(addr2) {
+                return true;
+            }
+            if parsed1.is_some() && parsed2.is_some() {
+                return false;
+            }
+            resolve_all(addr1)
+                .zip(resolve_all(addr2))
+                .is_some_and(|(first, second)| {
+                    first.iter().any(|resolved| second.contains(resolved))
+                })
+        }
     }
 }
 
+/// Every socket address a host:port spelling resolves to, `None` when the
+/// resolver does not know the name (which then compares unequal, at worst
+/// costing one extra dial).
+fn resolve_all(addr: &str) -> Option<Vec<SocketAddr>> {

Review Comment:
   `to_socket_addrs` is a blocking getaddrinfo on the async path 
(`dial_candidates`, `process_cluster_metadata`, the recheck on 58). a slow 
resolver stalls the runtime worker. spawn_blocking it or resolve once and cache.



##########
core/sdk/src/leader_aware.rs:
##########
@@ -253,6 +366,20 @@ mod tests {
         assert!(!is_same_address("192.168.1.1:8090", "127.0.0.1:8090"));
     }
 
+    // A host name and the address it resolves to name one endpoint; a name the
+    // resolver does not know compares unequal rather than erroring.
+    #[test]
+    fn a_host_name_matches_the_address_it_resolves_to() {

Review Comment:
   this passes with the resolver fallback replaced by `false`: `localhost` is 
rewritten by `normalize_address` before resolution runs. use a name that only 
resolves.



##########
core/integration/tests/sdk/disconnect_relogin.rs:
##########
@@ -0,0 +1,66 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! An explicit disconnect ends the session for good: the credentials a manual
+//! sign-in remembered (for reconnecting across involuntary drops and
+//! failovers) must not resurrect it. Pins at the Rust layer the contract the
+//! C++ e2e suite asserts through the FFI 
(`DisconnectThenReconnectWithoutRelogin`,
+//! `GetStatsBeforeLoginThrows`), so a regression fails here first instead of
+//! three suites downstream.
+
+use iggy::prelude::*;
+use integration::iggy_harness;
+
+#[iggy_harness]
+async fn 
given_a_logged_in_client_when_explicitly_disconnected_should_require_a_fresh_login(
+    harness: &TestHarness,
+) {
+    let client = harness.new_client().await.unwrap();
+    client
+        .login_user(DEFAULT_ROOT_USERNAME, DEFAULT_ROOT_PASSWORD)
+        .await
+        .unwrap();
+    client.get_me().await.expect("authenticated get_me works");
+
+    client.disconnect().await.unwrap();
+    client.connect().await.unwrap();
+    assert!(
+        matches!(client.get_me().await, Err(IggyError::Unauthenticated)),
+        "an explicit disconnect is caller intent, like a logout: the sign-in 
it ended \
+         must not be silently replayed by the reconnect, so the server sees an 
\

Review Comment:
   the request never reaches the server - `Unauthenticated` comes from the 
client-side gate. fix the message.



##########
core/common/src/traits/binary_impls/users.rs:
##########
@@ -174,6 +174,7 @@ impl<B: BinaryClient> UserClient for B {
             .to_bytes(),
         )
         .await?;
+        self.refresh_session_password(user_id, new_password).await;

Review Comment:
   this only refreshes the remembered login. with `AutoLogin::Enabled` the 
configured password stays stale after a change, so every later drop re-logs 
with the old one and fails `InvalidCredentials`. refresh it too or say so on 
`refresh_session_password`.



##########
foreign/go/client/tcp/tcp_session_management.go:
##########
@@ -165,14 +187,39 @@ func (c *IggyTcpClient) settleOnLeader(ctx 
context.Context, code uint32, body []
 
 // endBoundSession logs out a live session before a re-login, so the server
 // drops its client-table entry instead of leaving it to be fenced.
+//
+// The logout runs connect-scoped, and a failure it could recover from is
+// swallowed. Both because this call holds registerMtx: a logout that entered
+// the reconnect path would reconnect, sign in with the remembered credentials,
+// and deadlock on that lock. There is nothing to salvage either way -- a
+// session whose logout cannot be delivered died with its socket, and the
+// server fences what it left behind -- and the sign-in that follows replays
+// through its own reconnect.
 func (c *IggyTcpClient) endBoundSession(ctx context.Context) error {
        c.mtx.Lock()
        bound := c.session.Bound()
        c.mtx.Unlock()
        if !bound {
                return nil
        }
-       return c.LogoutUser(ctx)
+
+       err := c.LogoutUser(context.WithValue(ctx, connectScoped{}, struct{}{}))
+       if err == nil {
+               return nil
+       }
+       if !isReconnectable(err) {
+               return err
+       }
+
+       c.logger.Debug("The bound session's logout was not delivered; its 
socket ended it.",
+               slog.Any("error", err))
+       c.mtx.Lock()

Review Comment:
   this branch keeps the old remembered login: alice signed in, logout dropped, 
`LoginUser("bob")` rejected, next dropped request registers alice again. a 
delivered logout ends signed out. call `forgetLogin()` here.



##########
foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java:
##########
@@ -457,10 +494,40 @@ private AsyncTcpConnection openConnection(ConnectionInfo 
target) {
                 heartbeatInterval,
                 maxVsrFrameSize,
                 this::retryTransientOnLeader,
-                routingState::clearAssignments,
+                this::onSessionReset,
                 this::onConnectionFailure);
     }
 
+    /**
+     * 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: the
+     * remembered sign-in ends with it, so only credentials configured on the
+     * builder may bring the session back.
+     */
+    private void onSessionReset(int errorCode) {
+        routingState.clearAssignments();
+        if (errorCode != IggyErrorCode.STALE_CLIENT.getCode()) {
+            return;
+        }
+        // Credentials configured on the builder are what every connect of this
+        // client signs in as, so an eviction does not revoke them and the
+        // client recovers on its own. A sign-in a caller ran is different: the
+        // server ended that session deliberately, and reviving it behind the
+        // caller's back is what an explicit logout must not be able to do
+        // either. Dropped in both places, because the connection replays its
+        // own captured login to bring up a replacement channel.
+        if (username.isPresent() && password.isPresent()) {

Review Comment:
   with builder credentials A and a hand-run `login(B)`, an eviction revives B 
(captured payload) while a connection loss replays A. two different users 
depending on the failure kind.



##########
foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java:
##########
@@ -585,42 +654,112 @@ private CompletableFuture<Void> redialAttempt(int 
attempt, RetryPolicy policy) {
             log.error("Redial gave up after {} attempts, next request will 
fail fast", policy.getMaxRetries());
             return CompletableFuture.completedFuture(null);
         }
-        ConnectionInfo target = ReconnectPlan.target(connectionInfo, 
seedConnectionInfo, attempt);
-        Duration delay = ReconnectPlan.delay(policy, attempt);
+        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
+        // node just lost may be gone for good.
+        Duration delay = attempt == 1 && candidates.size() > 1 ? Duration.ZERO 
: ReconnectPlan.delay(policy, attempt);
         Executor delayedExecutor = 
CompletableFuture.delayedExecutor(delay.toMillis(), TimeUnit.MILLISECONDS);
-        return CompletableFuture.supplyAsync(() -> null, 
delayedExecutor).thenCompose(ignored -> {
-            if (closed) {
-                return CompletableFuture.completedFuture(null);
-            }
-            log.info("Redial attempt {}/{} to {}", attempt, 
policy.getMaxRetries(), target.serverAddress());
-            return retarget(target)
-                    .thenCompose(retargeted -> replayLogin())
-                    .handle((ok, error) -> {
-                        if (error == null) {
-                            log.info("Reconnected to {}", 
target.serverAddress());
-                            return 
CompletableFuture.<Void>completedFuture(null);
-                        }
+        return CompletableFuture.supplyAsync(() -> null, delayedExecutor)
+                .thenCompose(ignored -> sweepCandidates(candidates, 0, 
attempt, policy));
+    }
+
+    /**
+     * Dials one endpoint of a rotation and, if it does not come up, the next
+     * one. Every endpoint gets its turn inside one attempt, so a full pass 
over
+     * the cluster costs one retry rather than one per endpoint: with the
+     * default policy, rotating one endpoint per attempt would first dial a
+     * two-node survivor two delays in.
+     */
+    private CompletableFuture<Void> sweepCandidates(
+            List<ConnectionInfo> candidates, int index, int attempt, 
RetryPolicy policy) {
+        if (closed) {
+            return CompletableFuture.completedFuture(null);
+        }
+        if (index >= candidates.size()) {
+            return redialAttempt(attempt + 1, policy);
+        }
+        ConnectionInfo target = candidates.get(index);
+        log.info(
+                "Redial attempt {}/{} to {} ({}/{})",
+                attempt,
+                policy.getMaxRetries(),
+                target.serverAddress(),
+                index + 1,
+                candidates.size());
+        return retarget(target)
+                .handle((retargeted, dialError) -> {
+                    if (dialError != null) {
+                        log.warn("Redial to {} failed: {}", 
target.serverAddress(), dialError.getMessage());
+                        return sweepCandidates(candidates, index + 1, attempt, 
policy);
+                    }
+                    return replaySignInOn(target, candidates, index, attempt, 
policy);
+                })
+                .thenCompose(Function.identity());
+    }
+
+    /**
+     * Re-establishes the session on an endpoint that just came up.
+     *
+     * A sign-in the server rejected -- a rotated password, an expired token --
+     * ends the redial: the connection is up, no other endpoint would answer
+     * differently, and retrying would tear the working connection down on the
+     * next rotation and leave the client connected but unauthenticated anyway.
+     * The rejected credentials are dropped so nothing replays them.
+     */
+    private CompletableFuture<Void> replaySignInOn(
+            ConnectionInfo target, List<ConnectionInfo> candidates, int index, 
int attempt, RetryPolicy policy) {
+        return replayLogin()
+                .handle((ok, loginError) -> {
+                    if (loginError == null) {
+                        log.info("Reconnected to {}", target.serverAddress());
+                        return CompletableFuture.<Void>completedFuture(null);
+                    }
+                    if (isConnectionLoss(unwrap(loginError))) {
                         log.warn(
-                                "Redial attempt {} to {} failed: {}",
-                                attempt,
+                                "The sign-in on {} was lost with the 
connection: {}",

Review Comment:
   `isConnectionLoss` misses `IggyConnectionException` (channel closed before 
the reply), so a survivor that dies mid sign-in counts as a rejected login: 
remembered login dropped, client published on a dead node, every later call 
fails "not authenticated". reproduced. only a server rejection with a 
non-transient code is a rejection; everything else should move on.



##########
foreign/csharp/Iggy_SDK/IggyClient/Implementations/TcpMessageStream.cs:
##########
@@ -1026,7 +1063,17 @@ private async Task TryEstablishConnectionAsync(bool 
autoLogin, CancellationToken
                 // trailing segment of a large request until the previous one 
is acked.
                 socket.NoDelay = true;
 
-                await socket.ConnectAsync(host, port, token);
+                // Neither ConnectAsync nor the TLS handshake has a deadline 
of its own, and nothing up the
+                // stack adds one: a node whose syns are dropped - or one that 
accepts TCP and then never
+                // answers the ClientHello - would hold the sweep, and the 
connection semaphore with it, for
+                // the whole kernel connect timeout while a survivor goes 
untried.
+                using var dialCancellation = candidates.Length > 1
+                    ? CancellationTokenSource.CreateLinkedTokenSource(token)
+                    : null;
+                dialCancellation?.CancelAfter(FailoverDialTimeout);
+                var dialToken = dialCancellation?.Token ?? token;
+
+                await socket.ConnectAsync(host, port, dialToken);
                 dialed = true;

Review Comment:
   set before the handshake, so a handshake the 2s bound cuts short lands in 
the fatal filter at 1129: sweep aborted, the caller gets 
`OperationCanceledException` with a live token, survivor never dialed. set it 
after the handshake or use a separate `established` flag.



##########
foreign/node/src/client/client.connection.ts:
##########
@@ -227,6 +245,38 @@ export class IggyConnection extends EventEmitter {
     return connectPromise;
   }
 
+  /**
+   * Waits for one dial, bounded while other endpoints are queued behind it.
+   *
+   * A socket has no connect deadline of its own and there is no 'timeout'
+   * listener on it, so a node whose syns are dropped holds the pass for the
+   * whole OS connect timeout -- and it leads every pass, because the current
+   * endpoint only moves on success. The bound matches the Rust SDK's.
+   */
+  private async _dialWithin(socket: Socket, bounded: boolean): Promise<this> {

Review Comment:
   the bound covers the `'connect'` event only. on a tls socket that fires 
before `'secureConnect'`, so the handshake is still unbounded - rust, go and c# 
bound it too.



##########
foreign/node/src/client/client.connection.ts:
##########
@@ -67,6 +67,16 @@ const getTransport = (config: ClientConfig): Socket => {
   }
 };
 
+/** One node of the cluster, as a redial candidate. */
+export type Endpoint = { host: string, port: number };

Review Comment:
   one inline `{ host: string, port: number }` left at client.socket.ts:551.



##########
core/sdk/src/leader_aware.rs:
##########
@@ -162,13 +217,41 @@ fn process_cluster_metadata(
 
 /// Check if two addresses refer to the same endpoint
 /// Handles various formats like 127.0.0.1:8090 vs localhost:8090
-fn is_same_address(addr1: &str, addr2: &str) -> bool {
+///
+/// A host name and the address it resolves to are one endpoint too: a client
+/// configured as `iggy-server:8090` whose roster advertises `10.0.0.5:8090`
+/// would otherwise dial that node twice per failover sweep, and a single-node
+/// deployment would be treated as a cluster. Resolution is the last resort,
+/// only when the spellings differ and at least one side is not a literal
+/// address, and it is a blocking lookup: this runs on the connect and redirect
+/// paths, which are rare and already wait on the network.
+pub(crate) fn is_same_address(addr1: &str, addr2: &str) -> bool {
     match (parse_address(addr1), parse_address(addr2)) {
         (Some(sock1), Some(sock2)) => sock1.ip() == sock2.ip() && sock1.port() 
== sock2.port(),
-        _ => normalize_address(addr1) == normalize_address(addr2),
+        (parsed1, parsed2) => {
+            if normalize_address(addr1) == normalize_address(addr2) {
+                return true;
+            }
+            if parsed1.is_some() && parsed2.is_some() {

Review Comment:
   unreachable - the `(Some, Some)` arm above already returns.



##########
foreign/go/client/tcp/tcp_core.go:
##########
@@ -994,6 +1009,77 @@ func (c *IggyTcpClient) Connect(ctx context.Context) 
error {
        return nil
 }
 
+// awaitReestablish waits out what is left of the reestablishAfter window since
+// the last successful connection, if any.
+func (c *IggyTcpClient) awaitReestablish(connectedAt time.Time) {
+       if connectedAt.IsZero() {
+               return
+       }
+
+       elapsed := time.Since(connectedAt)
+       c.logger.Debug("Elapsed time since last connection", 
slog.Duration("elapsed", elapsed))
+       if remaining := c.config.reconnection.reestablishAfter - elapsed; 
remaining > 0 {
+               c.logger.Info("Trying to connect to the server", 
slog.Duration("remaining", remaining))
+               time.Sleep(remaining)
+       }
+}
+
+// dialCandidate brings one endpoint all the way up, wrapping it in TLS when
+// configured, and records the endpoint that answered: the leader check
+// compares against it and the next reconnect starts from it.
+//
+// bounded caps the whole attempt at failoverDialTimeout, for when other
+// endpoints are queued behind this one. Neither the dial nor the handshake has
+// a deadline of its own, and a node whose syns are dropped -- or one that
+// accepts TCP and then never answers the ClientHello -- would hold the sweep
+// for minutes while a survivor goes untried.
+func (c *IggyTcpClient) dialCandidate(ctx context.Context, address string, 
bounded bool) (net.Conn, error) {
+       c.logger.Info("Iggy client is connecting to server...", 
slog.String("server_address", address))
+       if bounded {
+               var cancel context.CancelFunc
+               ctx, cancel = context.WithTimeout(ctx, failoverDialTimeout)
+               defer cancel()
+       }
+
+       connection, err := (&net.Dialer{}).DialContext(ctx, "tcp", address)
+       if err != nil {
+               c.logger.Error("Failed to establish TCP connection to the 
server", slog.Any("error", err))
+               return nil, ierror.ErrCannotEstablishConnection
+       }
+
+       tc := connection.(*net.TCPConn)
+       if err := tc.SetNoDelay(c.config.noDelay); err != nil {
+               c.logger.Error("Failed to set the nodelay option on the client, 
continuing...", slog.Any("error", err))
+       }
+
+       established := connection
+       if c.config.tlsEnabled {
+               tlsConfig, err := c.createTLSConfig()

Review Comment:
   `createTLSConfig` takes the sni from `c.currentServerAddress`, which during 
a sweep is the endpoint just lost, not the candidate being dialed. with 
validation on and no `WithTLSDomain` a failover to a host with another name/ip 
fails `x509: certificate is valid for 127.0.0.1, not 127.0.0.2`. pass `address` 
in. the new tls tests run with validation off so they don't see it. also: a bad 
ca path/cert/domain from here is retried forever under the default `maxRetries` 
0; rust aborts on those.



##########
foreign/go/client/tcp/tcp_core.go:
##########
@@ -994,6 +1009,77 @@ func (c *IggyTcpClient) Connect(ctx context.Context) 
error {
        return nil
 }
 
+// awaitReestablish waits out what is left of the reestablishAfter window since
+// the last successful connection, if any.
+func (c *IggyTcpClient) awaitReestablish(connectedAt time.Time) {
+       if connectedAt.IsZero() {
+               return
+       }
+
+       elapsed := time.Since(connectedAt)
+       c.logger.Debug("Elapsed time since last connection", 
slog.Duration("elapsed", elapsed))
+       if remaining := c.config.reconnection.reestablishAfter - elapsed; 
remaining > 0 {
+               c.logger.Info("Trying to connect to the server", 
slog.Duration("remaining", remaining))
+               time.Sleep(remaining)

Review Comment:
   `time.Sleep` ignores `ctx` and `Close`. with the rotation a sweep whose 
other candidates are refused reaches this endpoint in ms and then sleeps the 
rest of `reestablishAfter` in `Connecting` - 60s under a 5s ctx. select on 
`ctx.Done()`.



##########
foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java:
##########
@@ -457,10 +494,40 @@ private AsyncTcpConnection openConnection(ConnectionInfo 
target) {
                 heartbeatInterval,
                 maxVsrFrameSize,
                 this::retryTransientOnLeader,
-                routingState::clearAssignments,
+                this::onSessionReset,
                 this::onConnectionFailure);
     }
 
+    /**
+     * 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: the
+     * remembered sign-in ends with it, so only credentials configured on the
+     * builder may bring the session back.
+     */
+    private void onSessionReset(int errorCode) {
+        routingState.clearAssignments();
+        if (errorCode != IggyErrorCode.STALE_CLIENT.getCode()) {
+            return;
+        }
+        // Credentials configured on the builder are what every connect of this
+        // client signs in as, so an eviction does not revoke them and the
+        // client recovers on its own. A sign-in a caller ran is different: the
+        // server ended that session deliberately, and reviving it behind the
+        // caller's back is what an explicit logout must not be able to do
+        // either. Dropped in both places, because the connection replays its
+        // own captured login to bring up a replacement channel.
+        if (username.isPresent() && password.isPresent()) {
+            return;
+        }
+        log.warn("The server evicted this session as stale; the sign-in it ran 
will not be replayed");
+        rememberedLogin = null;

Review Comment:
   not pinned: removing this line keeps both failover suites green.



##########
foreign/go/client/tcp/tcp_session_management.go:
##########
@@ -165,14 +187,39 @@ func (c *IggyTcpClient) settleOnLeader(ctx 
context.Context, code uint32, body []
 
 // endBoundSession logs out a live session before a re-login, so the server
 // drops its client-table entry instead of leaving it to be fenced.
+//
+// The logout runs connect-scoped, and a failure it could recover from is
+// swallowed. Both because this call holds registerMtx: a logout that entered
+// the reconnect path would reconnect, sign in with the remembered credentials,
+// and deadlock on that lock. There is nothing to salvage either way -- a
+// session whose logout cannot be delivered died with its socket, and the
+// server fences what it left behind -- and the sign-in that follows replays
+// through its own reconnect.
 func (c *IggyTcpClient) endBoundSession(ctx context.Context) error {
        c.mtx.Lock()
        bound := c.session.Bound()
        c.mtx.Unlock()
        if !bound {
                return nil
        }
-       return c.LogoutUser(ctx)
+
+       err := c.LogoutUser(context.WithValue(ctx, connectScoped{}, struct{}{}))

Review Comment:
   still deadlocks when the server answers the logout with 
`TransientNotAccepted` (it does that when it's not primary): `sendFrame` 
follows the redirect with `Connect` (tcp_core.go:646) without 
`skipAutoLoginOnce`, so `establishSession` -> `LoginUser` -> `register` waits 
on the `registerMtx` this goroutine holds. reproduced with a two-node fake. set 
the skip before that `Connect`.



##########
foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java:
##########
@@ -732,7 +874,35 @@ CompletableFuture<Optional<ConnectionInfo>> 
findLeaderElsewhere(ConnectionInfo c
         if (currentSystemClient == null) {
             return CompletableFuture.completedFuture(Optional.empty());
         }
-        return 
LeaderAwareness.findLeaderElsewhere(currentSystemClient::getClusterMetadata, 
currentTarget);
+        return 
LeaderAwareness.findLeaderElsewhere(currentSystemClient::getClusterMetadata, 
currentTarget)
+                .thenApply(lookup -> {
+                    // Replaced wholesale rather than merged: the roster is the
+                    // cluster's own answer about where its nodes are, so a 
node
+                    // it dropped stops being dialed. The configured seed is
+                    // kept separately and outlives it.
+                    if (!lookup.endpoints().isEmpty()) {

Review Comment:
   untested: assigning `lookup.endpoints()` unconditionally keeps everything 
green.



##########
foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java:
##########
@@ -638,6 +777,9 @@ CompletableFuture<IdentityInfo> 
loginOnLeader(Supplier<CompletableFuture<Identit
         CompletableFuture<IdentityInfo> callerFuture = new 
CompletableFuture<>();
         transaction.whenComplete((identity, error) -> {
             gate.complete(null);
+            if (error == null) {
+                rememberedLogin = loginAttempt;

Review Comment:
   a login still in flight when `close()` nulled this sets it again, so a later 
`connect()` + loss replays it.



##########
foreign/csharp/Iggy_SDK/IggyClient/Implementations/TcpMessageStream.cs:
##########
@@ -1095,6 +1146,20 @@ private async Task TryEstablishConnectionAsync(bool 
autoLogin, CancellationToken
 
                 _logger.LogError(e, "Failed to connect");
 
+                // Every other endpoint gets its turn before the retry delay: 
the node just lost may be gone for
+                // good, and pausing on it helps nothing.
+                if (++candidate < candidates.Length)
+                {
+                    _currentAddress = candidates[candidate];
+                    continue;
+                }
+
+                candidate = 0;
+                _currentAddress = candidates[0];
+
+                // The sweep is what the reconnection budget applies to, not a 
single dial: checked per dial,

Review Comment:
   no test for this: moving the checks back above the advance keeps the suite 
green.



##########
foreign/node/src/client/client.socket.ts:
##########
@@ -186,21 +203,24 @@ export class CommandResponseStream extends EventEmitter {
       if (!this.isAuthenticated && !this.isUnloggedCommand(command))
         await this.authenticate(this.options.credentials);
 
-      const response = await new Promise<CommandResponse>(
-        (resolve, reject) => {
-          const job = {
-            command,
-            payload,
-            handleResponse,
-            resolve,
-            reject
-          };
-          if (last)
-            this._execQueue.push(job);
-          else
-            this._execQueue.unshift(job);
-          this._processQueue();
-        });
+      // The roster read is itself a queued command and the queue is
+      // single-flighted, so the leader re-check cannot happen inside
+      // `_processVsr`. The refusal comes back out here instead, where the
+      // queue is free, and the command is re-issued on the node that now
+      // leads.
+      let response: CommandResponse;
+      for (let move = 0; ; move += 1) {
+        try {
+          response = await this._queueCommand(command, payload, handleResponse,
+            last);
+          break;
+        } catch (error) {
+          if (!(error instanceof LeaderMovedError))
+            throw error;
+          if (move >= MAX_LEADER_REDIRECTS || !await this._followLeaderMove())
+            throw responseError(command, error.refusal.errorCode);

Review Comment:
   when the roster still names this node or has no healthy leader the command 
now fails with 58 after one 2s window. before it was replayed for 30s, and rust 
still keeps replaying on the same connection and re-checking every 2s until the 
30s budget. keep that.



##########
foreign/node/src/client/client.socket.ts:
##########
@@ -186,21 +203,24 @@ export class CommandResponseStream extends EventEmitter {
       if (!this.isAuthenticated && !this.isUnloggedCommand(command))
         await this.authenticate(this.options.credentials);

Review Comment:
   after `_followLeaderMove()` the client has no session: `redirect()` emits 
`'disconnected'` which resets it, and this step sits above the loop. the 
retried replicated command fails `Unauthenticated` client-side, a 
non-replicated one goes out with session 0. authenticate again inside the loop 
after a move.



##########
foreign/csharp/Iggy_SDK/IggyClient/Implementations/TcpMessageStream.cs:
##########
@@ -1186,7 +1300,10 @@ private async Task<Stream> 
CreateSslStreamAndAuthenticate(Socket socket, TlsSett
         var stream = new NetworkStream(socket, true);
         var sslStream = new SslStream(stream, false, 
RemoteCertificateValidationCallback);
 
-        await sslStream.AuthenticateAsClientAsync(tlsSettings.Hostname);
+        // The token carries the dial bound when other endpoints are queued 
behind this one: a peer that

Review Comment:
   the `SslStream` is never disposed when the handshake fails.



##########
foreign/node/src/client/client.socket.ts:
##########
@@ -216,6 +236,57 @@ export class CommandResponseStream extends EventEmitter {
     }
   }
 
+  private _queueCommand(
+    command: number,
+    payload: Buffer,
+    handleResponse: boolean,
+    last: boolean
+  ): Promise<CommandResponse> {
+    return new Promise<CommandResponse>((resolve, reject) => {
+      const job = {
+        command,
+        payload,
+        handleResponse,
+        resolve,
+        reject
+      };
+      if (last)
+        this._execQueue.push(job);
+      else
+        this._execQueue.unshift(job);
+      this._processQueue();
+    });
+  }
+
+  /**
+   * Re-reads the roster and moves to the leader it names.
+   *
+   * @returns Whether the client moved, so the refused request is worth
+   * re-issuing
+   */
+  private async _followLeaderMove(): Promise<boolean> {

Review Comment:
   the roster read goes through `sendCommand` and hits the same 
`LeaderMovedError` loop with no depth guard. also, concurrent commands during a 
move: the first redirect's `'disconnected'` fails the others' roster reads, so 
they fail with 58 instead of being re-issued.



##########
foreign/node/src/client/client.connection.ts:
##########
@@ -300,36 +350,60 @@ export class IggyConnection extends EventEmitter {
   ): Promise<this> {
     let lastError = initialError;
     let expectedSocket = this.socket;
-    let attempt = 0;
+    let firstPass = true;
     while (enabled && this.reconnectCount < maxRetries) {
       this.connecting = true;
       this.reconnectCount += 1;
-      await waitForReconnect(interval);
+      const candidates = this._redialCandidates();
+      // The backoff paces retries against a single endpoint. With other
+      // 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)
+        await waitForReconnect(interval);
+      firstPass = false;
       if (this.ending)
         throw new Error('connection is closed', { cause: lastError });
       // A redirect may replace the socket at any point. Defer to the active
       // connection instead of dialing the superseded endpoint.
       if (this.connected || this.socket !== expectedSocket)
         return this.connect();
 
-      const options = this._reconnectTarget(attempt);
-      attempt += 1;
-      const socket = this._installSocket(
-        getTransport({ ...this.config, options })
-      );
-      this.socket = socket;
-      expectedSocket = socket;
-      try {
-        await this._waitForConnection(socket);
-        if (this.socket !== socket)
+      // Every endpoint gets its turn inside one attempt, so a full pass over
+      // the cluster costs one retry rather than one per endpoint: a pass that
+      // stopped at the first refusal would never reach the survivors of a
+      // client configured for a single retry.
+      for (const options of candidates) {
+        // Re-checked every iteration, not once above the loop: a destroy or a
+        // redirect mid-pass has to stop the pass. Left running, the next
+        // endpoint that answers would leave an open socket nobody closes, a
+        // 'connect' event after the destroy, and the process alive.
+        if (this.ending)
+          throw new Error('connection is closed', { cause: lastError });
+        if (this.socket !== expectedSocket)
           return this.connect();
-        this.config.options = options;
-        return this;
-      } catch (error) {
-        lastError = error instanceof Error
-          ? error
-          : new Error(String(error));
-        debug('reconnect attempt failed', lastError);
+
+        const socket = this._installSocket(
+          getTransport({ ...this.config, options })
+        );
+        this.socket = socket;
+        expectedSocket = socket;
+        try {
+          await this._dialWithin(socket, candidates.length > 1);

Review Comment:
   no test covers the bound - `_dialWithin(socket, false)` keeps the suite 
green. same for the first-pass skip at 357.



##########
foreign/node/src/client/client.connection.ts:
##########
@@ -67,6 +67,16 @@ const getTransport = (config: ClientConfig): Socket => {
   }
 };
 
+/** One node of the cluster, as a redial candidate. */
+export type Endpoint = { host: string, port: number };
+
+/**
+ * Bound on one dial while other endpoints are queued behind it. A socket has
+ * no connect deadline of its own, so a node whose syns are dropped would hold
+ * the whole pass. Matches the Rust SDK.
+ */
+const FAILOVER_DIAL_TIMEOUT_MS = 2_000;
+
 /**
  * Default reconnection settings.

Review Comment:
   the `DefaultReconnectOption` comment below and `maxRetries` in 
`client.type.ts:105` still say attempts; one retry is now a full pass over 
every candidate, and the first pass with a roster doesn't wait.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to