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 f90f9dfef address review comments
f90f9dfef is described below

commit f90f9dfefeed0ce6349d220ea9d64bcbfcae74ff
Author: Grzegorz Koszyk <[email protected]>
AuthorDate: Wed Aug 26 22:20:34 2026 +0200

    address review comments
---
 core/sdk/src/leader_aware.rs                       |  92 +++-
 core/sdk/src/tcp/tcp_client.rs                     | 581 ++++++++++++++-------
 .../IggyClient/Implementations/TcpMessageStream.cs |  28 +-
 3 files changed, 481 insertions(+), 220 deletions(-)

diff --git a/core/sdk/src/leader_aware.rs b/core/sdk/src/leader_aware.rs
index 66cdad35f..b5c038f34 100644
--- a/core/sdk/src/leader_aware.rs
+++ b/core/sdk/src/leader_aware.rs
@@ -132,11 +132,32 @@ pub async fn check_and_redirect_to_leader<C: 
ClusterClient>(
     }
 }
 
+/// Every endpoint the roster names for `transport`, empty when the read did
+/// not answer.
+///
+/// No leader verdict and no waiting for an election: the caller is not moving
+/// anywhere, it only wants somewhere to dial once the node it is on dies.
+pub(crate) async fn read_transport_endpoints<C: ClusterClient>(
+    client: &C,
+    transport: TransportProtocol,
+) -> Vec<String> {
+    match client.get_cluster_metadata().await {
+        Ok(metadata) => transport_endpoints(&metadata, transport),
+        Err(error) => {
+            debug!("Failed to read the cluster roster: {error}");
+            Vec::new()
+        }
+    }
+}
+
 /// How long to wait for a transiently leaderless cluster to elect before
 /// proceeding on the current node anyway.
 const LEADERLESS_WAIT_BUDGET: std::time::Duration = 
std::time::Duration::from_secs(5);
 const LEADERLESS_POLL_INTERVAL: std::time::Duration = 
std::time::Duration::from_millis(250);
 
+/// Bound on one name lookup made to compare two addresses.
+const RESOLVE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(2);
+
 /// One leader-check verdict from a cluster-metadata snapshot.
 enum Outcome {
     /// A healthy leader exists elsewhere; reconnect to it.
@@ -156,11 +177,23 @@ fn transport_endpoints(metadata: &ClusterMetadata, 
transport: TransportProtocol)
         .iter()
         .filter_map(|node| {
             let port = transport_port(node, transport);
-            (port != 0).then(|| format!("{}:{port}", node.ip))
+            (port != 0).then(|| node_address(node, port))
         })
         .collect()
 }
 
+/// One node's `host:port`, bracketing a literal IPv6 address. Appending a
+/// port to a bare `::1` yields a spelling no dial can parse, so an IPv6
+/// cluster would hand out a roster of undialable entries that still count as
+/// endpoints to fail over to.
+fn node_address(node: &ClusterNode, port: u16) -> String {
+    if node.ip.contains(':') && !node.ip.starts_with('[') {
+        format!("[{}]:{port}", node.ip)
+    } else {
+        format!("{}:{port}", node.ip)
+    }
+}
+
 fn transport_port(node: &ClusterNode, transport: TransportProtocol) -> u16 {
     match transport {
         TransportProtocol::Tcp => node.endpoints.tcp,
@@ -193,7 +226,7 @@ async fn process_cluster_metadata(
     match leader {
         Some(leader_node) => {
             let leader_port = transport_port(leader_node, transport);
-            let leader_address = format!("{}:{}", leader_node.ip, leader_port);
+            let leader_address = node_address(leader_node, leader_port);
 
             info!(
                 "Found leader node: {} at {} (using {} transport)",
@@ -272,10 +305,17 @@ where
 }
 
 /// 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).
+/// resolver does not know the name or does not answer in time (which then
+/// compares unequal, at worst costing one extra dial).
 async fn resolve_all(addr: String) -> Option<Vec<SocketAddr>> {
-    let resolved: Vec<SocketAddr> = 
tokio::net::lookup_host(addr).await.ok()?.collect();
+    // A resolver that never answers must not own the request budget: this
+    // comparison runs on the connect and redirect paths, the redirect one
+    // inside the caller's request deadline, and `lookup_host` has no deadline
+    // of its own.
+    let lookup = tokio::time::timeout(RESOLVE_TIMEOUT, 
tokio::net::lookup_host(addr))
+        .await
+        .ok()?;
+    let resolved: Vec<SocketAddr> = lookup.ok()?.collect();
     (!resolved.is_empty()).then_some(resolved)
 }
 
@@ -385,6 +425,48 @@ mod tests {
         assert!(transport_endpoints(&metadata, 
TransportProtocol::Quic).is_empty());
     }
 
+    // A port appended to a bare IPv6 address parses as neither, so the roster
+    // of an IPv6 cluster would name endpoints no dial can use while still
+    // counting as somewhere to fail over to.
+    #[test]
+    fn an_ipv6_node_is_named_as_a_bracketed_address() {
+        let metadata = ClusterMetadata {
+            name: "iggy".to_string(),
+            nodes: vec![
+                node("iggy-1", "::1", 8090, ClusterNodeRole::Leader),
+                node("iggy-2", "[fd00::2]", 8090, ClusterNodeRole::Follower),
+            ],
+        };
+
+        let endpoints = transport_endpoints(&metadata, TransportProtocol::Tcp);
+        assert_eq!(endpoints, vec!["[::1]:8090", "[fd00::2]:8090"]);
+        for endpoint in endpoints {
+            assert!(
+                SocketAddr::from_str(&endpoint).is_ok(),
+                "the roster named an endpoint no dial can parse: {endpoint}"
+            );
+        }
+    }
+
+    // The address a redirect hands to the next dial comes from the same
+    // roster entry, so it has to be spelled the same way.
+    #[tokio::test]
+    async fn a_redirect_to_an_ipv6_leader_names_a_dialable_address() {
+        let metadata = ClusterMetadata {
+            name: "iggy".to_string(),
+            nodes: vec![
+                node("iggy-1", "fd00::1", 8090, ClusterNodeRole::Leader),
+                node("iggy-2", "fd00::2", 8090, ClusterNodeRole::Follower),
+            ],
+        };
+
+        match process_cluster_metadata(&metadata, "[fd00::2]:8090", 
TransportProtocol::Tcp).await {
+            Outcome::Redirect(leader) => assert_eq!(leader, "[fd00::1]:8090"),
+            Outcome::LeaderIsCurrent => panic!("the follower was taken for the 
leader"),
+            Outcome::NoLeader => panic!("the roster named a healthy leader"),
+        }
+    }
+
     #[tokio::test]
     async fn test_is_same_address() {
         assert!(is_same_address("127.0.0.1:8090", "127.0.0.1:8090").await);
diff --git a/core/sdk/src/tcp/tcp_client.rs b/core/sdk/src/tcp/tcp_client.rs
index 5648fffaf..a271cef8c 100644
--- a/core/sdk/src/tcp/tcp_client.rs
+++ b/core/sdk/src/tcp/tcp_client.rs
@@ -16,8 +16,8 @@
 // under the License.
 
 use crate::leader_aware::{
-    LeaderRedirectionState, check_and_redirect_to_leader, is_same_address,
-    is_unauthenticated_metadata_probe,
+    LeaderRedirectionState, check_and_redirect_to_leader, is_same_spelling,
+    is_unauthenticated_metadata_probe, read_transport_endpoints,
 };
 use crate::prelude::Client;
 use crate::prelude::TcpClientConfig;
@@ -46,6 +46,7 @@ use std::net::SocketAddr;
 use std::str::FromStr;
 use std::sync::Arc;
 use std::sync::Mutex as StdMutex;
+use std::sync::atomic::{AtomicBool, Ordering};
 #[cfg(test)]
 use tokio::net::TcpListener;
 use tokio::net::TcpStream;
@@ -76,6 +77,11 @@ const NOT_READY_RETRY_INTERVAL: std::time::Duration = 
std::time::Duration::from_
 /// overall.
 const TRANSIENT_FAILOVER_CHECK_INTERVAL: std::time::Duration = 
std::time::Duration::from_secs(2);
 
+/// Bound on the roster read that follows a sign-in the caller ran itself. The
+/// read is a convenience for a failover that may never happen, so a cluster
+/// that answers it slowly must not hold up the sign-in.
+const ROSTER_READ_TIMEOUT: std::time::Duration = 
std::time::Duration::from_secs(5);
+
 /// Bound on one dial while the client has other endpoints to try. A host
 /// that drops the SYN -- powered off, or partitioned away -- takes the OS
 /// connect timeout to fail, which is minutes, and every other endpoint waits
@@ -100,6 +106,9 @@ pub struct TcpClient {
     /// unreachable exactly when it is needed, so the client has to have
     /// remembered it while the connection was still healthy.
     roster_endpoints: Mutex<Vec<String>>,
+    /// Set once a sign-in on this client has gone looking for the roster, so
+    /// that read happens once (see [`TcpClient::learn_roster_once`]).
+    roster_learned: AtomicBool,
     /// Credentials a sign-in on this client succeeded with, so a reconnect --
     /// onto this node or, after a failover, another one -- can re-establish
     /// the session instead of surfacing `Unauthenticated`. Cleared on logout.
@@ -134,6 +143,15 @@ struct EstablishedConnection {
     remote_address: SocketAddr,
 }
 
+/// A sign-in that did not complete, and whether the connection it ran on went
+/// with it. A connection that is gone leaves the endpoints the sweep has not
+/// reached yet worth dialing; one that stands means only the session is
+/// missing, which no other endpoint would answer differently.
+struct SignInFailure {
+    error: IggyError,
+    connection_lost: bool,
+}
+
 impl Default for TcpClient {
     fn default() -> Self {
         TcpClient::create(Arc::new(TcpClientConfig::default())).unwrap()
@@ -282,7 +300,9 @@ impl BinaryTransport for TcpClient {
 /// reached the log may be re-sent.
 ///
 /// - the errors raised before the frame was written, and the server's own
-///   refusals, which precede execution;
+///   refusals, which precede execution. A `StaleClient` eviction is neither:
+///   it arrives out of band and is consumed in place of the pending reply, so
+///   the request it interrupted may already have committed;
 /// - operations that never enter the log: a non-replicated read, and a logout,
 ///   which ends whatever session the connection carried -- the reconnect
 ///   brought a new one, and refusing the replay would strand
@@ -298,7 +318,6 @@ fn replay_is_safe(code: u32, error: &IggyError) -> bool {
             IggyError::NotConnected
                 | IggyError::CannotEstablishConnection
                 | IggyError::Unauthenticated
-                | IggyError::StaleClient
         )
         || matches!(
             operation_for_code(code),
@@ -306,28 +325,6 @@ fn replay_is_safe(code: u32, error: &IggyError) -> bool {
         )
 }
 
-/// Why a TLS handshake failed, as far as retrying is concerned.
-///
-/// A certificate this client will never accept -- the wrong CA, a name it does
-/// not cover, a peer that answers a ClientHello with something else -- says 
the
-/// same thing on every attempt. Reported as a configuration fault it ends the
-/// connect after one sweep; reported as a lost connection it would be redialed
-/// every interval forever under `max_retries = None`, which is how a wrong CA
-/// looks like a flaky network.
-fn classify_handshake_failure(error: &std::io::Error) -> IggyError {
-    match error
-        .get_ref()
-        .and_then(|inner| inner.downcast_ref::<rustls::Error>())
-    {
-        Some(
-            rustls::Error::InvalidCertificate(_)
-            | rustls::Error::NoCertificatesPresented
-            | rustls::Error::InvalidMessage(_),
-        ) => IggyError::InvalidTlsCertificate,
-        _ => IggyError::CannotEstablishConnection,
-    }
-}
-
 impl iggy_common::VsrSessionSealed for TcpClient {}
 
 #[async_trait::async_trait]
@@ -365,6 +362,7 @@ impl iggy_common::VsrSessionControl for TcpClient {
                 credentials,
                 user_id,
             });
+        self.learn_roster_once().await;
     }
 
     async fn forget_session_credentials(&self) {
@@ -495,6 +493,7 @@ impl TcpClient {
             leader_redirection_state: 
Mutex::new(LeaderRedirectionState::new()),
             current_server_address: Mutex::new(server_address),
             roster_endpoints: Mutex::new(Vec::new()),
+            roster_learned: AtomicBool::new(false),
             configured_password: Mutex::new(None),
             session_credentials: Mutex::new(None),
             consensus_session: 
Arc::new(StdMutex::new(ConsensusSession::new())),
@@ -505,26 +504,33 @@ impl TcpClient {
 
     async fn connect(&self) -> Result<(), IggyError> {
         loop {
-            match self.get_state().await {
-                ClientState::Shutdown => {
-                    trace!("Cannot connect. Client is shutdown.");
-                    return Err(IggyError::ClientShutdown);
-                }
-                ClientState::Connected
-                | ClientState::Authenticating
-                | ClientState::Authenticated => {
-                    let client_address = self.get_client_address_value().await;
-                    trace!("Client: {client_address} is already connected.");
-                    return Ok(());
-                }
-                ClientState::Connecting => {
-                    trace!("Client is already connecting.");
-                    return Ok(());
+            // Read and claimed under one lock acquisition. Apart, two callers
+            // both find `Disconnected` and both sweep: the loser's
+            // `reset_vsr_session` re-mints the client id under the identity 
the
+            // winner is binding, and its `replace` below drops the live
+            // authenticated stream.
+            {
+                let mut state = self.state.lock().await;
+                match *state {
+                    ClientState::Shutdown => {
+                        trace!("Cannot connect. Client is shutdown.");
+                        return Err(IggyError::ClientShutdown);
+                    }
+                    ClientState::Connected
+                    | ClientState::Authenticating
+                    | ClientState::Authenticated => {
+                        let client_address = 
self.get_client_address_value().await;
+                        trace!("Client: {client_address} is already 
connected.");
+                        return Ok(());
+                    }
+                    ClientState::Connecting => {
+                        trace!("Client is already connecting.");
+                        return Ok(());
+                    }
+                    _ => *state = ClientState::Connecting,
                 }
-                _ => {}
             }
 
-            self.set_state(ClientState::Connecting).await;
             let mut candidates = self.dial_candidates().await;
             // `reestablish_after` paces reconnects to the endpoint this client
             // was last on, and to that one only: the other endpoints owe it no
@@ -538,17 +544,25 @@ impl TcpClient {
                 candidates.rotate_left(1);
             }
 
+            let skip_auto_login = {
+                let mut guard = self.skip_auto_login_once.lock().await;
+                std::mem::take(&mut *guard)
+            };
+
             let mut retry_count = 0;
-            let connection_stream: ConnectionStreamKind;
-            let remote_address;
-            let client_address;
             let mut candidate = 0;
             // A fault no retry can fix, remembered rather than returned at
-            // once: it belongs to the endpoint that raised it (a certificate
-            // that names another host, a domain that will not parse), and the
-            // endpoints behind that one may be perfectly usable.
+            // once: it belongs to the endpoint that raised it (an unreadable 
CA
+            // file, a domain that will not parse), and the endpoints behind
+            // that one may be perfectly usable.
             let mut config_fault: Option<IggyError> = None;
-            loop {
+            // A sign-in that failed together with the connection it ran on.
+            // The sweep carries on -- a node that answers the dial and then
+            // goes quiet must not own the client, and it is also the endpoint
+            // the next connect would lead with -- and this is the reason the
+            // caller gets if nothing behind it works out either.
+            let mut sign_in_failure: Option<IggyError> = None;
+            let should_redirect = loop {
                 let server_address = candidates[candidate].clone();
                 if server_address == paced_endpoint
                     && let Some(remaining) = self.reestablish_wait().await
@@ -560,6 +574,7 @@ impl TcpClient {
                 info!("{NAME} client is connecting to server: 
{server_address}...");
                 match self.establish_bounded(&server_address, 
&candidates).await {
                     Ok(connection) => {
+                        let dialed = server_address.clone();
                         // The endpoint that answered is where this client now
                         // lives: the leader check compares against it, and the
                         // next reconnect starts from it. Recorded only once 
the
@@ -567,11 +582,41 @@ impl TcpClient {
                         // fails the TLS handshake does not become sticky and
                         // shadow the endpoints behind it.
                         *self.current_server_address.lock().await = 
server_address;
-                        client_address = connection.client_address;
-                        remote_address = connection.remote_address;
+                        let client_address = connection.client_address;
                         
self.client_address.lock().await.replace(client_address);
-                        connection_stream = connection.stream;
-                        break;
+                        let now = IggyTimestamp::now();
+                        info!(
+                            "{NAME} client: {client_address} has connected to 
server: {} at: {now}",
+                            connection.remote_address,
+                        );
+                        self.stream.lock().await.replace(connection.stream);
+                        self.set_state(ClientState::Connected).await;
+                        self.connected_at.lock().await.replace(now);
+                        self.publish_event(DiagnosticEvent::Connected).await;
+
+                        match self
+                            .establish_session(client_address, skip_auto_login)
+                            .await
+                        {
+                            Ok(should_redirect) => break should_redirect,
+                            Err(failure) if failure.connection_lost => {
+                                warn!(
+                                    "The sign-in on the server: {dialed} did 
not complete: {}",
+                                    failure.error,
+                                );
+                                sign_in_failure = Some(failure.error);
+                                // The sweep owns the state again: the sign-in
+                                // took the connection down with it, and left
+                                // `Disconnected` another caller would start a
+                                // second sweep alongside this one.
+                                self.set_state(ClientState::Connecting).await;
+                            }
+                            // The connection stands and only the session is
+                            // missing: rejected credentials say the same thing
+                            // on every node, and no endpoint behind this one
+                            // would answer differently.
+                            Err(failure) => return Err(failure.error),
+                        }
                     }
                     Err(IggyError::CannotEstablishConnection) => {}
                     Err(error) => config_fault = Some(error),
@@ -586,12 +631,11 @@ impl TcpClient {
                 }
                 candidate = 0;
 
-                // An unreadable CA file, a certificate that names another
-                // host, a domain that will not parse: no endpoint answered and
-                // at least one said why in a way that a retry cannot change,
-                // so the caller gets that reason instead of a retry loop that
-                // buries it (`max_retries = None` would otherwise redial it
-                // every interval forever).
+                // An unreadable CA file, a domain that will not parse: no
+                // endpoint answered and at least one said why in a way that a
+                // retry cannot change, so the caller gets that reason instead
+                // of a retry loop that buries it (`max_retries = None` would
+                // otherwise redial it every interval forever).
                 if let Some(error) = config_fault {
                     self.fail_connect().await;
                     return Err(error);
@@ -604,7 +648,7 @@ impl TcpClient {
                 if !self.config.reconnection.enabled {
                     warn!("Automatic reconnection is disabled.");
                     self.fail_connect().await;
-                    return Err(IggyError::CannotEstablishConnection);
+                    return 
Err(sign_in_failure.unwrap_or(IggyError::CannotEstablishConnection));
                 }
 
                 let unlimited_retries = 
self.config.reconnection.max_retries.is_none();
@@ -629,109 +673,7 @@ impl TcpClient {
                 }
 
                 self.fail_connect().await;
-                return Err(IggyError::CannotEstablishConnection);
-            }
-
-            let now = IggyTimestamp::now();
-            info!(
-                "{NAME} client: {client_address} has connected to server: 
{remote_address} at: {now}",
-            );
-            self.stream.lock().await.replace(connection_stream);
-            self.set_state(ClientState::Connected).await;
-            self.connected_at.lock().await.replace(now);
-            self.publish_event(DiagnosticEvent::Connected).await;
-            let skip_auto_login = {
-                let mut guard = self.skip_auto_login_once.lock().await;
-                std::mem::take(&mut *guard)
-            };
-
-            // Handle auto-login
-            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
-                }
-                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;
-                        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}.")
-                            }
-                            Err(error) => {
-                                // With the transport up and only the session
-                                // missing, 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.
-                                //
-                                // A sign-in can also fail because the socket
-                                // died under it. Whatever is left of that
-                                // connection cannot carry a request, so it 
goes
-                                // rather than being kept behind a `Connected`
-                                // that makes the next `connect()` a no-op and
-                                // leaves every gated operation failing until
-                                // someone calls `disconnect()` by hand.
-                                if matches!(
-                                    error,
-                                    IggyError::Disconnected
-                                        | IggyError::EmptyResponse
-                                        | IggyError::NotConnected
-                                        | IggyError::CannotEstablishConnection
-                                        | IggyError::TcpError
-                                        | IggyError::StaleClient
-                                ) {
-                                    self.disconnect_transport().await?;
-                                } else if self.get_state().await == 
ClientState::Authenticating {
-                                    
self.set_state(ClientState::Connected).await;
-                                }
-                                // A rejected credential does not become valid 
on
-                                // the next reconnect, and replaying it costs 
an
-                                // argon2 on the server every time. Configured
-                                // credentials stay as configured -- they are 
the
-                                // caller's to fix -- so only the remembered
-                                // sign-in is dropped.
-                                if matches!(
-                                    error,
-                                    IggyError::InvalidCredentials
-                                        | IggyError::InvalidUsername
-                                        | IggyError::InvalidPassword
-                                        | IggyError::Unauthenticated
-                                ) {
-                                    self.forget_session_credentials().await;
-                                }
-                                return Err(error);
-                            }
-                        }
-
-                        // The sole leader settlement, and it runs
-                        // authenticated. Any node completes a login now -- a
-                        // backup forwards the register to the primary -- so
-                        // this decides where later ops land, not whether
-                        // sign-in works.
-                        self.handle_leader_redirection().await?
-                    }
-                }
+                return 
Err(sign_in_failure.unwrap_or(IggyError::CannotEstablishConnection));
             };
 
             if should_redirect {
@@ -742,6 +684,105 @@ impl TcpClient {
         }
     }
 
+    /// Re-establish the session on a connection that just came up and settle 
it
+    /// on the leader. Reports whether the leader check asks for a redirect.
+    async fn establish_session(
+        &self,
+        client_address: SocketAddr,
+        skip_auto_login: bool,
+    ) -> Result<bool, SignInFailure> {
+        let Some(credentials) = self.sign_in_credentials().await else {
+            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.
+            return Ok(false);
+        };
+
+        if skip_auto_login {
+            info!("Skipping automatic sign-in for a retried login/register 
request.");
+            return Ok(false);
+        }
+
+        info!("{NAME} client: {client_address} is signing in...");
+        self.set_state(ClientState::Authenticating).await;
+        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}."),
+            Err(error) => return Err(self.fail_sign_in(error).await),
+        }
+
+        // The sole leader settlement, and it runs authenticated. Any node
+        // completes a login now -- a backup forwards the register to the
+        // primary -- so this decides where later ops land, not whether sign-in
+        // works.
+        self.handle_leader_redirection()
+            .await
+            .map_err(|error| SignInFailure {
+                error,
+                connection_lost: false,
+            })
+    }
+
+    /// Put the client back into a state that describes what a failed sign-in
+    /// left behind, and report whether the connection survived it.
+    async fn fail_sign_in(&self, error: IggyError) -> SignInFailure {
+        // A sign-in can fail because the socket died under it. Whatever is 
left
+        // of that connection cannot carry a request, so it goes rather than
+        // being kept behind a `Connected` that makes the next `connect()` a
+        // no-op and leaves every gated operation failing until someone calls
+        // `disconnect()` by hand.
+        let connection_lost = matches!(
+            error,
+            IggyError::Disconnected
+                | IggyError::EmptyResponse
+                | IggyError::NotConnected
+                | IggyError::CannotEstablishConnection
+                | IggyError::TcpError
+                | IggyError::StaleClient
+        );
+        if connection_lost {
+            if let Err(teardown_error) = self.disconnect_transport().await {
+                warn!("Failed to drop the connection of a failed sign-in: 
{teardown_error}");
+            }
+        } else if self.get_state().await == ClientState::Authenticating {
+            // With the transport up and only the session missing, 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;
+        }
+
+        // A rejected credential does not become valid on the next reconnect,
+        // and replaying it costs an argon2 on the server every time. 
Configured
+        // credentials stay as configured -- they are the caller's to fix -- so
+        // only the remembered sign-in is dropped.
+        if matches!(
+            error,
+            IggyError::InvalidCredentials
+                | IggyError::InvalidUsername
+                | IggyError::InvalidPassword
+                | IggyError::Unauthenticated
+        ) {
+            self.forget_session_credentials().await;
+        }
+
+        SignInFailure {
+            error,
+            connection_lost,
+        }
+    }
+
     /// Checks cluster metadata and handles leader redirection if needed.
     /// Returns true if redirection occurred and reconnection is needed.
     pub(crate) async fn handle_leader_redirection(&self) -> Result<bool, 
IggyError> {
@@ -803,11 +844,12 @@ impl TcpClient {
     /// 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> {
-        // The sign-in that last succeeded, whoever ran it. One rule in every
-        // SDK: a client is whoever it last signed in as, so the same failure
-        // restores the same session everywhere. A configured `AutoLogin` signs
-        // in through this very path, so for a client that never signed in by
-        // hand the remembered credentials *are* the configured ones.
+        // The sign-in that last succeeded, whoever ran it: a client is whoever
+        // it last signed in as, so a reconnect restores the session the caller
+        // last asked for rather than one it had moved off. A configured
+        // `AutoLogin` signs in through this very path, so for a client that
+        // never signed in by hand the remembered credentials *are* the
+        // configured ones.
         if let Some(remembered) = 
self.session_credentials.lock().await.as_ref() {
             return Some(remembered.credentials.clone());
         }
@@ -831,6 +873,48 @@ impl TcpClient {
         }
     }
 
+    /// Read the cluster roster once, on the first sign-in that succeeds on a
+    /// client whose caller signs it in by hand.
+    ///
+    /// `connect()` follows the sign-in it runs itself with a leader check, and
+    /// that check is what refreshes the roster. A client with no configured
+    /// `AutoLogin` is signed in by its caller instead, and only `IggyClient`
+    /// follows that with a leader check, so a raw transport would know exactly
+    /// one endpoint -- the one it was configured with -- and redial the node
+    /// that died for as long as it lived.
+    ///
+    /// Once per client, which is also what keeps the read from nesting: it 
goes
+    /// through the reconnect path, whose sign-in calls straight back into 
here.
+    /// Bounded for the same reason: the read is a convenience for a failover
+    /// that may never happen, so it must not hold up the sign-in that 
triggered
+    /// it -- unbounded retries would do exactly that.
+    async fn learn_roster_once(&self) {
+        // Only a live session can read the roster, and only the caller's own
+        // sign-in leaves one behind here: a connect that signs in follows it
+        // with a leader check of its own.
+        if self.auto_login_configured()
+            || self.get_state().await != ClientState::Authenticated
+            || self.roster_learned.swap(true, Ordering::SeqCst)
+        {
+            return;
+        }
+
+        let read = read_transport_endpoints(self, TransportProtocol::Tcp);
+        let Ok(endpoints) = tokio::time::timeout(ROSTER_READ_TIMEOUT, 
read).await else {
+            warn!("Reading the cluster roster took longer than 
{ROSTER_READ_TIMEOUT:?}");
+            return;
+        };
+        if endpoints.is_empty() {
+            return;
+        }
+
+        info!(
+            "{NAME} client learned {} endpoint(s) to fail over to.",
+            endpoints.len()
+        );
+        *self.roster_endpoints.lock().await = endpoints;
+    }
+
     /// Endpoints to dial for one connect, likeliest first: where the client
     /// currently is, the address it was configured with, then the roster it
     /// learned while connected.
@@ -843,14 +927,13 @@ impl TcpClient {
         let roster = self.roster_endpoints.lock().await.clone();
         let configured = std::iter::once(&self.config.server_address);
         for endpoint in configured.chain(roster.iter()) {
-            let mut known = false;
-            for candidate in &candidates {
-                if is_same_address(candidate, endpoint).await {
-                    known = true;
-                    break;
-                }
-            }
-            if !known {
+            // Spellings only, no name resolution: one duplicate endpoint costs
+            // a dial that fails on its own, while a resolver that does not
+            // answer would stall the failover before it dialed anything.
+            if !candidates
+                .iter()
+                .any(|candidate| is_same_spelling(candidate, endpoint))
+            {
                 candidates.push(endpoint.clone());
             }
         }
@@ -943,8 +1026,14 @@ impl TcpClient {
             IggyError::InvalidTlsDomain
         })?;
         let stream = connector.connect(domain, stream).await.map_err(|error| {
+            // The verdict describes the peer, not this client: a certificate
+            // that names another host, a peer that answers a ClientHello with
+            // something else. The endpoints behind it may be fine, and with
+            // one roster entry per node the SNI is a bare address that no
+            // certificate has to cover, so this ends the dial rather than the
+            // connect.
             error!("Failed to establish a TLS connection to the server: 
{error}");
-            classify_handshake_failure(&error)
+            IggyError::CannotEstablishConnection
         })?;
 
         Ok(EstablishedConnection {
@@ -1016,8 +1105,18 @@ impl TcpClient {
     /// failover unauthenticated. The public [`Client::disconnect`] wraps this
     /// and forgets them first.
     async fn disconnect_transport(&self) -> Result<(), IggyError> {
-        if self.get_state().await == ClientState::Disconnected {
-            return Ok(());
+        match self.get_state().await {
+            ClientState::Disconnected => return Ok(()),
+            // A connect is already sweeping, and every caller here is tearing
+            // the connection down in order to reconnect -- which is what that
+            // sweep is doing. Tearing it down under the sweep would re-mint 
the
+            // client id the sign-in in flight is binding and take the stream 
it
+            // just installed.
+            ClientState::Connecting => {
+                trace!("Not disconnecting; a connect is already in flight.");
+                return Ok(());
+            }
+            _ => {}
         }
 
         let client_address = self.get_client_address_value().await;
@@ -1313,6 +1412,7 @@ const fn is_login_register_code(code: u32) -> bool {
 mod tests {
     use super::*;
     use iggy_binary_protocol::codes::{GET_ME_CODE, LOGOUT_USER_CODE, 
SEND_MESSAGES_CODE};
+    use std::sync::atomic::AtomicUsize;
     use tokio::io::{AsyncReadExt, AsyncWriteExt};
 
     const SESSION_USER_ID: u32 = 7;
@@ -1350,15 +1450,23 @@ mod tests {
     }
 
     /// A peer that accepts TCP and hangs up without a byte: enough for the
-    /// dial, never enough for a TLS handshake.
-    async fn endpoint_that_hangs_up() -> String {
+    /// dial, never enough for a TLS handshake or a sign-in. The counter says
+    /// how many dials reached it.
+    async fn counted_endpoint_that_hangs_up() -> (String, Arc<AtomicUsize>) {
         let (listener, address) = live_endpoint().await;
+        let dials = Arc::new(AtomicUsize::new(0));
+        let accepted = dials.clone();
         tokio::spawn(async move {
             while let Ok((stream, _)) = listener.accept().await {
+                accepted.fetch_add(1, Ordering::SeqCst);
                 drop(stream);
             }
         });
-        address
+        (address, dials)
+    }
+
+    async fn endpoint_that_hangs_up() -> String {
+        counted_endpoint_that_hangs_up().await.0
     }
 
     // With reconnection off there are no retries, but the endpoints the roster
@@ -1469,11 +1577,15 @@ mod tests {
     }
 
     /// A peer that accepts TCP and then answers a ClientHello with something
-    /// else: the handshake fails for a reason no retry changes.
-    async fn endpoint_that_speaks_no_tls() -> String {
+    /// else, so the handshake fails on the peer's own answer. The counter says
+    /// how many dials reached it.
+    async fn counted_endpoint_that_speaks_no_tls() -> (String, 
Arc<AtomicUsize>) {
         let (listener, address) = live_endpoint().await;
+        let dials = Arc::new(AtomicUsize::new(0));
+        let accepted = dials.clone();
         tokio::spawn(async move {
             while let Ok((mut stream, _)) = listener.accept().await {
+                accepted.fetch_add(1, Ordering::SeqCst);
                 tokio::spawn(async move {
                     let _ = stream.write_all(b"this is not a TLS 
record\n").await;
                     // Held open, so the failure is the handshake's verdict
@@ -1483,18 +1595,58 @@ mod tests {
                 });
             }
         });
-        address
+        (address, dials)
     }
 
-    // A certificate this client will never accept says the same thing on every
-    // attempt, so it has to reach the caller instead of being redialed every
-    // interval forever -- which is what `max_retries = None` did with it.
+    // A handshake verdict describes the peer -- a certificate that names
+    // another host, an answer that is not TLS at all -- and not this client's
+    // configuration, so it ends the dial rather than the connect: the 
endpoints
+    // behind it are untried, and a redial can find a repaired node.
     #[tokio::test]
-    async fn a_handshake_no_retry_can_fix_ends_the_connect() {
+    async fn a_handshake_the_peer_failed_is_dialed_again() {
+        let (plaintext, dials) = counted_endpoint_that_speaks_no_tls().await;
         let client = TcpClient::create(Arc::new(TcpClientConfig {
-            server_address: endpoint_that_speaks_no_tls().await,
+            server_address: plaintext,
             tls_enabled: true,
             tls_validate_certificate: false,
+            reconnection: TcpClientReconnectionConfig {
+                // One retry, so the pass runs twice and the connect still ends
+                // on its own.
+                max_retries: Some(1),
+                interval: 
NonZeroIggyDuration::from_str("100ms").expect("duration"),
+                ..TcpClientReconnectionConfig::default()
+            },
+            ..TcpClientConfig::default()
+        }))
+        .expect("create the client");
+
+        let connect = tokio::time::timeout(
+            std::time::Duration::from_secs(10),
+            std::pin::pin!(TcpClient::connect(&client)),
+        )
+        .await
+        .expect("the connect has to end on its own");
+        assert!(matches!(connect, Err(IggyError::CannotEstablishConnection)));
+        assert_eq!(
+            dials.load(Ordering::SeqCst),
+            2,
+            "a handshake the peer failed ended the connect instead of the dial"
+        );
+        assert_eq!(client.get_state().await, ClientState::Disconnected);
+    }
+
+    // A CA file that cannot be read is this client's own configuration, and it
+    // says the same thing on every attempt: reported as a lost connection it
+    // would be redialed every interval forever under `max_retries = None`,
+    // which is how a wrong CA path looks like a flaky network.
+    #[tokio::test]
+    async fn a_ca_file_that_cannot_be_read_ends_the_connect() {
+        let (_listener, endpoint) = live_endpoint().await;
+        let client = TcpClient::create(Arc::new(TcpClientConfig {
+            server_address: endpoint,
+            tls_enabled: true,
+            tls_validate_certificate: true,
+            tls_ca_file: Some("no-such-ca-file.pem".to_string()),
             reconnection: TcpClientReconnectionConfig {
                 // Unlimited retries, so a transient classification never
                 // returns and this test times out instead of failing.
@@ -1512,7 +1664,39 @@ mod tests {
         )
         .await
         .expect("the connect has to end on its own");
-        assert!(matches!(connect, Err(IggyError::InvalidTlsCertificate)));
+        assert!(matches!(connect, Err(IggyError::InvalidTlsCertificatePath)));
+        assert_eq!(client.get_state().await, ClientState::Disconnected);
+    }
+
+    // A node that answers the dial and then cannot carry the sign-in has to
+    // hand the sweep on. Ending it there leaves that node the one the client 
is
+    // recorded on, so every later connect leads with it and the endpoints
+    // behind it are never reached.
+    #[tokio::test]
+    async fn a_sign_in_that_failed_hands_the_sweep_on_to_the_next_endpoint() {
+        let (dialed_first, _) = counted_endpoint_that_hangs_up().await;
+        let (survivor, survivor_dials) = 
counted_endpoint_that_hangs_up().await;
+        let client = TcpClient::create(Arc::new(TcpClientConfig {
+            server_address: dialed_first,
+            auto_login: AutoLogin::Enabled(Credentials::UsernamePassword(
+                "iggy".to_string(),
+                "iggy".into(),
+            )),
+            reconnection: TcpClientReconnectionConfig {
+                enabled: false,
+                ..TcpClientReconnectionConfig::default()
+            },
+            ..TcpClientConfig::default()
+        }))
+        .expect("create the client");
+        *client.roster_endpoints.lock().await = vec![survivor];
+
+        assert!(TcpClient::connect(&client).await.is_err());
+        assert_eq!(
+            survivor_dials.load(Ordering::SeqCst),
+            1,
+            "the endpoint behind the one whose sign-in failed was never dialed"
+        );
         assert_eq!(client.get_state().await, ClientState::Disconnected);
     }
 
@@ -1563,9 +1747,11 @@ mod tests {
             SEND_MESSAGES_CODE,
             &IggyError::CannotEstablishConnection
         ));
-        assert!(replay_is_safe(SEND_MESSAGES_CODE, &IggyError::StaleClient));
         // Written, and its outcome unknown: a replicated write must not be
-        // re-sent under a session the fence cannot match it against.
+        // re-sent under a session the fence cannot match it against. An
+        // eviction is consumed in place of the reply, so it says nothing about
+        // whether the write committed.
+        assert!(!replay_is_safe(SEND_MESSAGES_CODE, &IggyError::StaleClient));
         assert!(!replay_is_safe(
             SEND_MESSAGES_CODE,
             &IggyError::Disconnected
@@ -1956,10 +2142,9 @@ mod tests {
         }
     }
 
-    // One rule in every SDK: a client is whoever it last signed in as, so the
-    // same failure restores the same session in each of them. The connection
-    // re-authenticates from the login it captured, and a redial that replayed
-    // somebody else would make the outcome depend on which got there first.
+    // A client is whoever it last signed in as: the connection 
re-authenticates
+    // from the login it captured, and a redial that replayed somebody else
+    // would make the outcome depend on which of the two got there first.
     #[tokio::test]
     async fn the_last_sign_in_outranks_the_configured_credentials() {
         let client = TcpClient::create(Arc::new(TcpClientConfig {
diff --git 
a/foreign/csharp/Iggy_SDK/IggyClient/Implementations/TcpMessageStream.cs 
b/foreign/csharp/Iggy_SDK/IggyClient/Implementations/TcpMessageStream.cs
index 38bbfebb2..245299d09 100644
--- a/foreign/csharp/Iggy_SDK/IggyClient/Implementations/TcpMessageStream.cs
+++ b/foreign/csharp/Iggy_SDK/IggyClient/Implementations/TcpMessageStream.cs
@@ -21,7 +21,6 @@ using System.Net;
 using System.Net.Security;
 using System.Net.Sockets;
 using System.Runtime.InteropServices;
-using System.Security.Authentication;
 using System.Security.Cryptography.X509Certificates;
 using Apache.Iggy.Configuration;
 using Apache.Iggy.Contracts;
@@ -1157,8 +1156,8 @@ public sealed partial class TcpMessageStream : IIggyClient
                 if (IsTlsConfigurationFault(e))
                 {
                     // A fault no retry can fix, kept aside rather than thrown 
at once: it belongs to the
-                    // endpoint that raised it - a certificate that names 
another host - and the endpoints
-                    // behind that one may be perfectly usable.
+                    // endpoint that raised it - a CA file that cannot be read 
- and the endpoints behind that
+                    // one may be perfectly usable.
                     configurationFault = e;
                 }
 
@@ -1174,8 +1173,8 @@ public sealed partial class TcpMessageStream : IIggyClient
                 _currentAddress = candidates[0];
 
                 // No endpoint answered and at least one said why in a way no 
retry changes: an unreadable CA
-                // file, a certificate this client will never accept. The 
caller gets that reason instead of a
-                // retry loop that buries it - unlimited retries would 
otherwise redial it forever.
+                // file. The caller gets that reason instead of a retry loop 
that buries it - unlimited retries
+                // would otherwise redial it forever.
                 if (configurationFault is not null)
                 {
                     SetConnectionState(ConnectionState.Disconnected);
@@ -1344,21 +1343,16 @@ public sealed partial class TcpMessageStream : 
IIggyClient
 
     /// <summary>
     ///     Whether bringing an endpoint up failed for a reason that says this 
client's own TLS configuration is
-    ///     wrong: a CA file that cannot be read, or a certificate it will 
never accept. Neither changes on a
-    ///     retry, so the sweep reports it instead of redialing forever.
+    ///     wrong: a CA file that cannot be read. It does not change on a 
retry, so the sweep reports it instead
+    ///     of redialing forever.
     /// </summary>
     private static bool IsTlsConfigurationFault(Exception e)
     {
-        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);
+        // A handshake verdict describes the peer, not this client: a 
certificate that names another host, a
+        // peer that answers a ClientHello with something else. The endpoints 
behind it may be fine, and with
+        // one roster entry per node the target host is a bare address no 
certificate has to cover, so a failed
+        // handshake ends the dial rather than the connect.
+        return e is InvalidCertificatePathException;
     }
 
     private async Task SendAckAsync(int code, ReadOnlyMemory<byte> body, 
CancellationToken token)

Reply via email to