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)