hubcio commented on code in PR #3944:
URL: https://github.com/apache/iggy/pull/3944#discussion_r3866006416
##########
core/sdk/src/tcp/tcp_client.rs:
##########
@@ -486,36 +646,81 @@ impl TcpClient {
};
// Handle auto-login
- let should_redirect = match &self.config.auto_login {
- AutoLogin::Disabled => {
- info!("Automatic sign-in is disabled.");
+ let should_redirect = match self.sign_in_credentials().await {
+ None => {
+ info!("No credentials to sign in with.");
// Only `IggyClient` redirects after a manual sign-in, so
// a raw transport can stay on a backup: its first
// replicated write gets `TransientNotAccepted`, the
// redirect drops the session, and the retry fails
// `Unauthenticated` until the caller signs in again.
false
}
- AutoLogin::Enabled(credentials) => {
+ Some(credentials) => {
if skip_auto_login {
info!("Skipping automatic sign-in for a retried
login/register request.");
false
} else {
info!("{NAME} client: {client_address} is signing
in...");
self.set_state(ClientState::Authenticating).await;
- match credentials {
- Credentials::UsernamePassword(username, password)
=> {
- self.login_user(username,
password.expose_secret()).await?;
- info!(
- "{NAME} client: {client_address} has
signed in with the user credentials, username: {username}",
- );
+ let signed_in = match &credentials {
+ Credentials::UsernamePassword(username, password)
=> self
+ .login_user(username, password.expose_secret())
+ .await
+ .map(|_| format!("the user credentials,
username: {username}")),
+ Credentials::PersonalAccessToken(token) => self
+
.login_with_personal_access_token(token.expose_secret())
+ .await
+ .map(|_| "a personal access token".to_owned()),
+ };
+ match signed_in {
+ Ok(how) => {
+ info!("{NAME} client: {client_address} has
signed in with {how}.")
}
- Credentials::PersonalAccessToken(token) => {
-
self.login_with_personal_access_token(token.expose_secret())
- .await?;
- info!(
- "{NAME} client: {client_address} has
signed in with a personal access token.",
- );
+ Err(error) => {
+ // 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);
Review Comment:
a sign-in failure kills the whole sweep - `candidates[1..]` never get
dialed, and `:569` already pinned that node as `current_server_address`.
`reestablish_wait()` is `None` after the 30s stall so the rotate at `:537`
can't move it out either, so a node that accepts tcp then goes quiet owns the
client forever. java continues the sweep on any non-rejection failure
(`AsyncIggyTcpClient.java:756`); go, C# and node stop like this does.
##########
core/sdk/src/tcp/tcp_client.rs:
##########
@@ -202,6 +275,59 @@ impl BinaryTransport for TcpClient {
}
}
+/// Whether replaying `code` over a fresh connection cannot apply it twice.
+///
+/// The reconnect registers a new client identity, so the server's dedup fence
+/// no longer covers the original request: only requests that provably never
+/// reached the log may be re-sent.
+///
+/// - the errors raised before the frame was written, and the server's own
+/// refusals, which precede execution;
+/// - 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
+/// `logout_before_relogin`, whose failure aborts the sign-in that was about
+/// to replace the session;
+/// - login and register, where the replay is the protocol: the server stays
+/// deliberately silent on a transient register failure and relies on the
+/// client resending.
+fn replay_is_safe(code: u32, error: &IggyError) -> bool {
+ is_login_register_code(code)
+ || matches!(
+ error,
+ IggyError::NotConnected
+ | IggyError::CannotEstablishConnection
+ | IggyError::Unauthenticated
+ | IggyError::StaleClient
+ )
+ || matches!(
+ operation_for_code(code),
+ Operation::NonReplicated | Operation::Logout
+ )
+}
+
+/// 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(_),
Review Comment:
`InvalidMessage` and `NoCertificatesPresented` describe the peer, not this
client's config, but they become `config_fault` and end the connect at `:595`
with no retry. `InvalidCertificate` too - with `tls_domain` empty the SNI is
the candidate (`:931`), so a roster IP with no matching SAN is terminal. C# has
the same problem. cut by fault owner: only a bad CA path or unparsable domain.
##########
core/sdk/src/tcp/tcp_client.rs:
##########
@@ -486,36 +646,81 @@ impl TcpClient {
};
// Handle auto-login
- let should_redirect = match &self.config.auto_login {
- AutoLogin::Disabled => {
- info!("Automatic sign-in is disabled.");
+ let should_redirect = match self.sign_in_credentials().await {
+ None => {
+ info!("No credentials to sign in with.");
Review Comment:
a bare `TcpClient` with `AutoLogin::Disabled` never learns a roster - this
arm returns without reading one, and `roster_endpoints` is only written from
`handle_leader_redirection`, called from the `Some` arm below.
`dial_candidates()` stays one address and the reconnect redials the dead node
forever. every other SDK reads it from its own sign-in path.
##########
core/sdk/src/tcp/tcp_client.rs:
##########
@@ -202,6 +275,59 @@ impl BinaryTransport for TcpClient {
}
}
+/// Whether replaying `code` over a fresh connection cannot apply it twice.
+///
+/// The reconnect registers a new client identity, so the server's dedup fence
+/// no longer covers the original request: only requests that provably never
+/// reached the log may be re-sent.
+///
+/// - the errors raised before the frame was written, and the server's own
+/// refusals, which precede execution;
+/// - 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
+/// `logout_before_relogin`, whose failure aborts the sign-in that was about
+/// to replace the session;
+/// - login and register, where the replay is the protocol: the server stays
+/// deliberately silent on a transient register failure and relies on the
+/// client resending.
+fn replay_is_safe(code: u32, error: &IggyError) -> bool {
+ is_login_register_code(code)
Review Comment:
good that this replaced an unconditional replay, but `StaleClient` shouldn't
be replay-safe for replicated codes. the eviction arrives out-of-band and is
consumed in place of the pending reply (`vsr.rs:216`), so the write may already
be committed, and the replay goes out under a fresh `client_id` the dedup fence
can't match. the real predicate is whether the frame was written - C# splits
exactly that (`VsrConnection.cs:313`).
##########
core/sdk/src/leader_aware.rs:
##########
@@ -110,8 +147,31 @@ enum Outcome {
NoLeader,
}
+/// Every node's address for `transport`, in roster order. A node that does
+/// not expose the transport reports port 0 and is skipped: dialing it would
+/// burn a failover attempt on an endpoint that cannot answer.
+fn transport_endpoints(metadata: &ClusterMetadata, transport:
TransportProtocol) -> Vec<String> {
+ metadata
+ .nodes
+ .iter()
+ .filter_map(|node| {
+ let port = transport_port(node, transport);
+ (port != 0).then(|| format!("{}:{port}", node.ip))
Review Comment:
doesn't bracket ipv6, so `::1` becomes `::1:8090`, which
`SocketAddr::from_str` rejects. an ipv6 cluster gets a roster of undialable
entries that still push `candidates.len()` past 1. go uses `net.JoinHostPort`,
C# `ServerAddress.HostPort`, node and java keep host and port separate.
##########
core/common/src/types/configuration/tcp_config/tcp_client_reconnection_config.rs:
##########
@@ -20,10 +20,30 @@ use std::str::FromStr;
#[derive(Debug, Clone)]
pub struct TcpClientReconnectionConfig {
+ /// Whether a lost connection is redialed at all. With this off the
+ /// endpoints the client knows still get one pass, since they were
+ /// configured to be tried, but nothing is retried after it.
pub enabled: bool,
+ /// How many passes over the known endpoints *after the first*, or `None`
+ /// for unlimited. `Some(0)` still makes that one pass, since the endpoints
Review Comment:
rust and C# match this doc; go, java and node do N passes with no free first
one. `0` differs four ways - rust 1 pass, C# and go unlimited, java one pass
only with more than one candidate, node zero. worth converging, or documenting
the table.
##########
core/sdk/src/tcp/tcp_client.rs:
##########
@@ -321,155 +525,111 @@ impl TcpClient {
}
self.set_state(ClientState::Connecting).await;
Review Comment:
`get_state()` above and this are separate lock acquisitions, so two callers
can both read `Disconnected` and both sweep. the second's
`disconnect_transport` doesn't early-return on `Connecting`, so
`reset_vsr_session` re-mints `client_id` under the first's in-flight
`bind_vsr_session`, and the loser's `replace` at `:639` drops a live
authenticated stream. pre-existing, but the window is a whole sweep now; go
added both guards in this PR (`tcp_core.go:942`, `:1075`).
##########
core/sdk/src/tcp/tcp_client.rs:
##########
@@ -574,7 +787,235 @@ impl TcpClient {
}
}
- async fn disconnect(&self) -> Result<(), IggyError> {
+ /// Whether an `AutoLogin` is configured on this client, which makes the
+ /// session after any connect the configured user's rather than whoever
+ /// signed in by hand.
+ pub(crate) fn auto_login_configured(&self) -> bool {
+ matches!(self.config.auto_login, AutoLogin::Enabled(_))
+ }
+
+ /// Credentials to sign in with after connecting: the configured ones, or
+ /// else the ones a manual sign-in on this client succeeded with. A manual
+ /// sign-in is otherwise less reconnectable than a configured one, which
+ /// is a surprising difference between two ways of doing the same thing.
+ ///
+ /// A password change this client committed for the configured user is
+ /// applied on top: the configured password will never work again, and
every
+ /// later reconnect would otherwise fail `InvalidCredentials`.
+ async fn sign_in_credentials(&self) -> Option<Credentials> {
+ // 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.
+ if let Some(remembered) =
self.session_credentials.lock().await.as_ref() {
Review Comment:
remembered-first here and in java (`AsyncIggyTcpClient.java:809`),
configured-first in go (`tcp_core.go:1271`) and C#
(`TcpMessageStream.cs:1310`), nothing remembered in node. `AutoLogin(A)` +
`login_user(B)` comes back as B in two SDKs and A in three. the comment above
claims one rule everywhere and `:1964` pins it.
##########
core/sdk/src/leader_aware.rs:
##########
@@ -160,15 +215,70 @@ fn process_cluster_metadata(
}
}
-/// Check if two addresses refer to the same endpoint
-/// Handles various formats like 127.0.0.1:8090 vs localhost:8090
-fn is_same_address(addr1: &str, addr2: &str) -> bool {
+/// Whether two addresses are written the same way, up to canonicalization
+/// (`localhost` and `[::]` spellings, and a literal address compared as an
+/// address rather than as text).
+///
+/// Cheap and non-blocking, which is the whole point: the resolving comparison
+/// below is a `getaddrinfo`, and every caller reaches this first.
+pub(crate) fn is_same_spelling(addr1: &str, addr2: &str) -> bool {
match (parse_address(addr1), parse_address(addr2)) {
(Some(sock1), Some(sock2)) => sock1.ip() == sock2.ip() && sock1.port()
== sock2.port(),
_ => normalize_address(addr1) == normalize_address(addr2),
}
}
+/// Check if two addresses refer to the same endpoint
+/// Handles various formats like 127.0.0.1:8090 vs localhost:8090
+///
+/// A host name and the address it resolves to are one endpoint too: a client
+/// configured as `iggy-server:8090` whose roster advertises `10.0.0.5:8090`
+/// would otherwise dial that node twice per failover sweep, and a single-node
+/// deployment would be treated as a cluster.
+///
+/// Resolution is the last resort, only when the spellings differ and at least
+/// one side is not a literal address. It runs through the runtime's resolver
+/// rather than `ToSocketAddrs`: name lookup is a blocking `getaddrinfo`, and
+/// this is called from the connect and redirect paths, where stalling a
+/// runtime worker on a slow resolver would stall every task sharing it.
+pub(crate) async fn is_same_address(addr1: &str, addr2: &str) -> bool {
+ is_same_address_with(addr1, addr2, resolve_all).await
+}
+
+/// [`is_same_address`] against a caller-provided resolver, so the fallback can
+/// be exercised without depending on what the machine's resolver answers.
+async fn is_same_address_with<R, F>(addr1: &str, addr2: &str, resolve: R) ->
bool
+where
+ R: Fn(String) -> F,
+ F: Future<Output = Option<Vec<SocketAddr>>>,
+{
+ if is_same_spelling(addr1, addr2) {
+ return true;
+ }
+
+ // Two literal addresses that did not compare equal are different
+ // endpoints; resolving them would only hand back what they already say.
+ if parse_address(addr1).is_some() && parse_address(addr2).is_some() {
+ return false;
+ }
+
+ let (Some(first), Some(second)) = (
+ resolve(addr1.to_owned()).await,
+ resolve(addr2.to_owned()).await,
+ ) else {
+ return false;
+ };
+ first.iter().any(|resolved| second.contains(resolved))
+}
+
+/// Every socket address a host:port spelling resolves to, `None` when the
+/// resolver does not know the name (which then compares unequal, at worst
+/// costing one extra dial).
+async fn resolve_all(addr: String) -> Option<Vec<SocketAddr>> {
+ let resolved: Vec<SocketAddr> =
tokio::net::lookup_host(addr).await.ok()?.collect();
Review Comment:
`lookup_host` with no deadline on the failover path - per candidate pair
from `dial_candidates`, and once per leader check from
`process_cluster_metadata`, which sits on `send_raw`'s replay path inside the
30s budget. was a sync compare before this PR. spelling-only in
`dial_candidates`, and wrap this one in a `timeout` - the resolution here is
load-bearing.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]