hubcio commented on code in PR #3944:
URL: https://github.com/apache/iggy/pull/3944#discussion_r3853696586
##########
core/sdk/src/tcp/tcp_client.rs:
##########
@@ -485,36 +546,61 @@ impl TcpClient {
};
// Handle auto-login
- let should_redirect = match &self.config.auto_login {
- AutoLogin::Disabled => {
- info!("Automatic sign-in is disabled.");
+ let should_redirect = match self.sign_in_credentials().await {
+ None => {
+ info!("No credentials to sign in with.");
// Only `IggyClient` redirects after a manual sign-in, so
// a raw transport can stay on a backup: its first
// replicated write gets `TransientNotAccepted`, the
// redirect drops the session, and the retry fails
// `Unauthenticated` until the caller signs in again.
false
}
- AutoLogin::Enabled(credentials) => {
+ Some(credentials) => {
if skip_auto_login {
info!("Skipping automatic sign-in for a retried
login/register request.");
false
} else {
info!("{NAME} client: {client_address} is signing
in...");
self.set_state(ClientState::Authenticating).await;
- match credentials {
- Credentials::UsernamePassword(username, password)
=> {
- self.login_user(username,
password.expose_secret()).await?;
- info!(
- "{NAME} client: {client_address} has
signed in with the user credentials, username: {username}",
- );
+ let signed_in = match &credentials {
+ Credentials::UsernamePassword(username, password)
=> self
+ .login_user(username, password.expose_secret())
+ .await
+ .map(|_| format!("the user credentials,
username: {username}")),
+ Credentials::PersonalAccessToken(token) => self
+
.login_with_personal_access_token(token.expose_secret())
+ .await
+ .map(|_| "a personal access token".to_owned()),
+ };
+ match signed_in {
+ Ok(how) => {
+ info!("{NAME} client: {client_address} has
signed in with {how}.")
}
- Credentials::PersonalAccessToken(token) => {
-
self.login_with_personal_access_token(token.expose_secret())
- .await?;
- info!(
- "{NAME} client: {client_address} has
signed in with a personal access token.",
- );
+ Err(error) => {
+ // The transport is up and only the session is
+ // not, so the state has to say so: left at
+ // `Authenticating` every gated operation fails
+ // client-side with `Disconnected`, `connect()`
+ // returns ok without dialing, and nothing
short
+ // of an explicit `login_user` recovers.
+ self.set_state(ClientState::Connected).await;
Review Comment:
this also runs when the sign-in failed because the socket died: state goes
to `Connected` with no stream, `connect()` is a no-op and every gated op fails
until an explicit `disconnect()`. only reset to `Connected` if the stream
survived.
##########
core/sdk/src/tcp/tcp_client.rs:
##########
@@ -573,7 +667,205 @@ impl TcpClient {
}
}
- async fn disconnect(&self) -> Result<(), IggyError> {
+ /// Whether an `AutoLogin` is configured on this client, which makes the
+ /// session after any connect the configured user's rather than whoever
+ /// signed in by hand.
+ pub(crate) fn auto_login_configured(&self) -> bool {
+ matches!(self.config.auto_login, AutoLogin::Enabled(_))
+ }
+
+ /// Credentials to sign in with after connecting: the configured ones, or
+ /// else the ones a manual sign-in on this client succeeded with. A manual
+ /// sign-in is otherwise less reconnectable than a configured one, which
+ /// is a surprising difference between two ways of doing the same thing.
+ async fn sign_in_credentials(&self) -> Option<Credentials> {
+ match &self.config.auto_login {
+ AutoLogin::Enabled(credentials) => Some(credentials.clone()),
+ AutoLogin::Disabled => self
+ .session_credentials
+ .lock()
+ .await
+ .as_ref()
+ .map(|remembered| remembered.credentials.clone()),
+ }
+ }
+
+ /// Endpoints to dial for one connect, likeliest first: where the client
+ /// currently is, then the roster it learned while connected, then the
+ /// configured seeds.
+ async fn dial_candidates(&self) -> Vec<String> {
+ let mut candidates =
vec![self.current_server_address.lock().await.clone()];
+ let roster = self.roster_endpoints.lock().await.clone();
+ for endpoint in
roster.iter().chain(self.config.failover_addresses.iter()) {
+ if !candidates
+ .iter()
+ .any(|candidate| is_same_address(candidate, endpoint))
+ {
+ candidates.push(endpoint.clone());
+ }
+ }
+ candidates
+ }
+
+ /// Bring one endpoint all the way up: TCP connect, socket options, and the
+ /// TLS handshake when it is configured. Nothing about the connection is
+ /// recorded until this succeeds, so a half-usable endpoint leaves no trace
+ /// for the next connect to lead with.
+ ///
+ /// `CannotEstablishConnection` means this endpoint failed and the next one
+ /// is worth trying; any other error is a configuration fault that no
+ /// endpoint can satisfy.
+ async fn establish(&self, server_address: &str) ->
Result<EstablishedConnection, IggyError> {
+ let stream = TcpStream::connect(server_address).await.map_err(|error| {
+ error!("Failed to connect to server: {server_address}. Error:
{error}");
+ IggyError::CannotEstablishConnection
+ })?;
+ let client_address = stream.local_addr().map_err(|error| {
+ error!("Failed to get the local address of the client: {error}");
+ IggyError::CannotEstablishConnection
+ })?;
+ let remote_address = stream.peer_addr().map_err(|error| {
+ error!("Failed to get the remote address of the server: {error}");
+ IggyError::CannotEstablishConnection
+ })?;
+
+ if let Err(error) = stream.set_nodelay(self.config.nodelay) {
+ error!("Failed to set the nodelay option on the client: {error},
continuing...");
+ }
+
+ if !self.config.tls_enabled {
+ return Ok(EstablishedConnection {
+ stream:
ConnectionStreamKind::Tcp(TcpConnectionStream::new(client_address, stream)),
+ client_address,
+ remote_address,
+ });
+ }
+
+ let _ =
rustls::crypto::aws_lc_rs::default_provider().install_default();
+ let config = if self.config.tls_validate_certificate {
+ let mut root_cert_store = rustls::RootCertStore::empty();
+ if let Some(certificate_path) = &self.config.tls_ca_file {
+ for cert in
CertificateDer::pem_file_iter(certificate_path).map_err(|error| {
+ error!("Failed to read the CA file: {certificate_path}.
{error}");
+ IggyError::InvalidTlsCertificatePath
+ })? {
+ let certificate = cert.map_err(|error| {
+ error!(
+ "Failed to read a certificate from the CA file:
{certificate_path}. {error}",
+ );
+ IggyError::InvalidTlsCertificate
+ })?;
+ root_cert_store.add(certificate).map_err(|error| {
+ error!(
+ "Failed to add a certificate to the root
certificate store. {error}"
+ );
+ IggyError::InvalidTlsCertificate
+ })?;
+ }
+ } else {
+
root_cert_store.extend(webpki_roots::TLS_SERVER_ROOTS.iter().cloned());
+ }
+
+ rustls::ClientConfig::builder()
+ .with_root_certificates(root_cert_store)
+ .with_no_client_auth()
+ } else {
+ use crate::tcp::tcp_tls_verifier::NoServerVerification;
+ rustls::ClientConfig::builder()
+ .dangerous()
+
.with_custom_certificate_verifier(Arc::new(NoServerVerification))
+ .with_no_client_auth()
+ };
+
+ let connector = TlsConnector::from(Arc::new(config));
+ let tls_domain = if self.config.tls_domain.is_empty() {
+ // Extract hostname/IP from server_address when tls_domain is not
specified
+ server_address
+ .split(':')
+ .next()
+ .unwrap_or(server_address)
+ .to_string()
+ } else {
+ self.config.tls_domain.to_owned()
+ };
+ let domain = ServerName::try_from(tls_domain).map_err(|error| {
+ error!("Failed to create a server name from the domain. {error}");
+ IggyError::InvalidTlsDomain
+ })?;
+ let stream = connector.connect(domain, stream).await.map_err(|error| {
+ error!("Failed to establish a TLS connection to the server:
{error}");
+ IggyError::CannotEstablishConnection
Review Comment:
with a single endpoint and `max_retries = None` a wrong ca, a name mismatch
or a plaintext peer now retries forever every second - before it returned at
once. map rustls `InvalidCertificate` to a config fault so it aborts like the
path/domain errors.
##########
core/sdk/src/tcp/tcp_client.rs:
##########
@@ -159,16 +204,38 @@ impl BinaryTransport for TcpClient {
return Err(IggyError::Disconnected);
}
- if matches!(self.config.auto_login, AutoLogin::Disabled) &&
!is_login_register_code(code) {
- // Without auto-login a reconnect cannot re-establish the session,
- // so non-login requests fail fast. Login/register itself is the
+ if !is_login_register_code(code) &&
self.sign_in_credentials().await.is_none() {
+ // With no credentials -- neither configured nor remembered from a
+ // sign-in -- a reconnect cannot re-establish the session, so
+ // non-login requests fail fast. Login/register itself is the
// exception: the server stays deliberately silent on transient
// register failures (the server `surface_login_failure`) and
// relies on the client timing out and replaying the request.
return Err(error);
}
- self.disconnect().await?;
+ // Reconnecting heals the transport, but replaying the request over the
+ // new connection is a second attempt under a new session:
+ // `reset_vsr_session` drops the client id the server's dedup fence is
+ // keyed on, so a replicated write that committed before its reply was
+ // lost would apply a second time. Replay only what provably never
+ // reached the log -- the errors raised before the request was written,
+ // and the operations that never enter it.
+ //
+ // Login and register are the exception: the server stays deliberately
+ // silent on a transient register failure and relies on the client
+ // replaying, so that replay is the protocol rather than a retry.
+ let replay_after_reconnect = is_login_register_code(code)
+ || matches!(
+ error,
+ IggyError::NotConnected
+ | IggyError::CannotEstablishConnection
Review Comment:
`Operation::Logout` is not in this set, so a logout over a dropped socket is
not replayed: `logout_before_relogin` fails `login_user(B)` on the first call
and the reconnect re-signs the old user. a logout under a fresh session can't
double-apply anything - treat it as replayable, like go does.
##########
core/sdk/src/leader_aware.rs:
##########
@@ -162,13 +217,41 @@ fn process_cluster_metadata(
/// Check if two addresses refer to the same endpoint
/// Handles various formats like 127.0.0.1:8090 vs localhost:8090
-fn is_same_address(addr1: &str, addr2: &str) -> bool {
+///
+/// A host name and the address it resolves to are one endpoint too: a client
+/// configured as `iggy-server:8090` whose roster advertises `10.0.0.5:8090`
+/// would otherwise dial that node twice per failover sweep, and a single-node
+/// deployment would be treated as a cluster. Resolution is the last resort,
+/// only when the spellings differ and at least one side is not a literal
+/// address, and it is a blocking lookup: this runs on the connect and redirect
+/// paths, which are rare and already wait on the network.
+pub(crate) fn is_same_address(addr1: &str, addr2: &str) -> bool {
match (parse_address(addr1), parse_address(addr2)) {
(Some(sock1), Some(sock2)) => sock1.ip() == sock2.ip() && sock1.port()
== sock2.port(),
- _ => normalize_address(addr1) == normalize_address(addr2),
+ (parsed1, parsed2) => {
+ if normalize_address(addr1) == normalize_address(addr2) {
+ return true;
+ }
+ if parsed1.is_some() && parsed2.is_some() {
+ return false;
+ }
+ resolve_all(addr1)
+ .zip(resolve_all(addr2))
+ .is_some_and(|(first, second)| {
+ first.iter().any(|resolved| second.contains(resolved))
+ })
+ }
}
}
+/// Every socket address a host:port spelling resolves to, `None` when the
+/// resolver does not know the name (which then compares unequal, at worst
+/// costing one extra dial).
+fn resolve_all(addr: &str) -> Option<Vec<SocketAddr>> {
Review Comment:
`to_socket_addrs` is a blocking getaddrinfo on the async path
(`dial_candidates`, `process_cluster_metadata`, the recheck on 58). a slow
resolver stalls the runtime worker. spawn_blocking it or resolve once and cache.
##########
core/sdk/src/leader_aware.rs:
##########
@@ -253,6 +366,20 @@ mod tests {
assert!(!is_same_address("192.168.1.1:8090", "127.0.0.1:8090"));
}
+ // A host name and the address it resolves to name one endpoint; a name the
+ // resolver does not know compares unequal rather than erroring.
+ #[test]
+ fn a_host_name_matches_the_address_it_resolves_to() {
Review Comment:
this passes with the resolver fallback replaced by `false`: `localhost` is
rewritten by `normalize_address` before resolution runs. use a name that only
resolves.
##########
core/integration/tests/sdk/disconnect_relogin.rs:
##########
@@ -0,0 +1,66 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! An explicit disconnect ends the session for good: the credentials a manual
+//! sign-in remembered (for reconnecting across involuntary drops and
+//! failovers) must not resurrect it. Pins at the Rust layer the contract the
+//! C++ e2e suite asserts through the FFI
(`DisconnectThenReconnectWithoutRelogin`,
+//! `GetStatsBeforeLoginThrows`), so a regression fails here first instead of
+//! three suites downstream.
+
+use iggy::prelude::*;
+use integration::iggy_harness;
+
+#[iggy_harness]
+async fn
given_a_logged_in_client_when_explicitly_disconnected_should_require_a_fresh_login(
+ harness: &TestHarness,
+) {
+ let client = harness.new_client().await.unwrap();
+ client
+ .login_user(DEFAULT_ROOT_USERNAME, DEFAULT_ROOT_PASSWORD)
+ .await
+ .unwrap();
+ client.get_me().await.expect("authenticated get_me works");
+
+ client.disconnect().await.unwrap();
+ client.connect().await.unwrap();
+ assert!(
+ matches!(client.get_me().await, Err(IggyError::Unauthenticated)),
+ "an explicit disconnect is caller intent, like a logout: the sign-in
it ended \
+ must not be silently replayed by the reconnect, so the server sees an
\
Review Comment:
the request never reaches the server - `Unauthenticated` comes from the
client-side gate. fix the message.
##########
core/common/src/traits/binary_impls/users.rs:
##########
@@ -174,6 +174,7 @@ impl<B: BinaryClient> UserClient for B {
.to_bytes(),
)
.await?;
+ self.refresh_session_password(user_id, new_password).await;
Review Comment:
this only refreshes the remembered login. with `AutoLogin::Enabled` the
configured password stays stale after a change, so every later drop re-logs
with the old one and fails `InvalidCredentials`. refresh it too or say so on
`refresh_session_password`.
##########
foreign/go/client/tcp/tcp_session_management.go:
##########
@@ -165,14 +187,39 @@ func (c *IggyTcpClient) settleOnLeader(ctx
context.Context, code uint32, body []
// endBoundSession logs out a live session before a re-login, so the server
// drops its client-table entry instead of leaving it to be fenced.
+//
+// The logout runs connect-scoped, and a failure it could recover from is
+// swallowed. Both because this call holds registerMtx: a logout that entered
+// the reconnect path would reconnect, sign in with the remembered credentials,
+// and deadlock on that lock. There is nothing to salvage either way -- a
+// session whose logout cannot be delivered died with its socket, and the
+// server fences what it left behind -- and the sign-in that follows replays
+// through its own reconnect.
func (c *IggyTcpClient) endBoundSession(ctx context.Context) error {
c.mtx.Lock()
bound := c.session.Bound()
c.mtx.Unlock()
if !bound {
return nil
}
- return c.LogoutUser(ctx)
+
+ err := c.LogoutUser(context.WithValue(ctx, connectScoped{}, struct{}{}))
+ if err == nil {
+ return nil
+ }
+ if !isReconnectable(err) {
+ return err
+ }
+
+ c.logger.Debug("The bound session's logout was not delivered; its
socket ended it.",
+ slog.Any("error", err))
+ c.mtx.Lock()
Review Comment:
this branch keeps the old remembered login: alice signed in, logout dropped,
`LoginUser("bob")` rejected, next dropped request registers alice again. a
delivered logout ends signed out. call `forgetLogin()` here.
##########
foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java:
##########
@@ -457,10 +494,40 @@ private AsyncTcpConnection openConnection(ConnectionInfo
target) {
heartbeatInterval,
maxVsrFrameSize,
this::retryTransientOnLeader,
- routingState::clearAssignments,
+ this::onSessionReset,
this::onConnectionFailure);
}
+ /**
+ * A server-side eviction reached this client. The routing state it cached
+ * belonged to the evicted session, and a stale-client eviction is the
+ * server ending that session authoritatively, like a logout: the
+ * remembered sign-in ends with it, so only credentials configured on the
+ * builder may bring the session back.
+ */
+ private void onSessionReset(int errorCode) {
+ routingState.clearAssignments();
+ if (errorCode != IggyErrorCode.STALE_CLIENT.getCode()) {
+ return;
+ }
+ // Credentials configured on the builder are what every connect of this
+ // client signs in as, so an eviction does not revoke them and the
+ // client recovers on its own. A sign-in a caller ran is different: the
+ // server ended that session deliberately, and reviving it behind the
+ // caller's back is what an explicit logout must not be able to do
+ // either. Dropped in both places, because the connection replays its
+ // own captured login to bring up a replacement channel.
+ if (username.isPresent() && password.isPresent()) {
Review Comment:
with builder credentials A and a hand-run `login(B)`, an eviction revives B
(captured payload) while a connection loss replays A. two different users
depending on the failure kind.
##########
foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java:
##########
@@ -585,42 +654,112 @@ private CompletableFuture<Void> redialAttempt(int
attempt, RetryPolicy policy) {
log.error("Redial gave up after {} attempts, next request will
fail fast", policy.getMaxRetries());
return CompletableFuture.completedFuture(null);
}
- ConnectionInfo target = ReconnectPlan.target(connectionInfo,
seedConnectionInfo, attempt);
- Duration delay = ReconnectPlan.delay(policy, attempt);
+ List<ConnectionInfo> candidates = redialCandidates();
+ // The delay paces rotations, not dials. The first rotation runs at
once
+ // when there is somewhere else to go: pausing before dialing a
survivor
+ // only pushes the failover past the window the caller waits in, and
the
+ // node just lost may be gone for good.
+ Duration delay = attempt == 1 && candidates.size() > 1 ? Duration.ZERO
: ReconnectPlan.delay(policy, attempt);
Executor delayedExecutor =
CompletableFuture.delayedExecutor(delay.toMillis(), TimeUnit.MILLISECONDS);
- return CompletableFuture.supplyAsync(() -> null,
delayedExecutor).thenCompose(ignored -> {
- if (closed) {
- return CompletableFuture.completedFuture(null);
- }
- log.info("Redial attempt {}/{} to {}", attempt,
policy.getMaxRetries(), target.serverAddress());
- return retarget(target)
- .thenCompose(retargeted -> replayLogin())
- .handle((ok, error) -> {
- if (error == null) {
- log.info("Reconnected to {}",
target.serverAddress());
- return
CompletableFuture.<Void>completedFuture(null);
- }
+ return CompletableFuture.supplyAsync(() -> null, delayedExecutor)
+ .thenCompose(ignored -> sweepCandidates(candidates, 0,
attempt, policy));
+ }
+
+ /**
+ * Dials one endpoint of a rotation and, if it does not come up, the next
+ * one. Every endpoint gets its turn inside one attempt, so a full pass
over
+ * the cluster costs one retry rather than one per endpoint: with the
+ * default policy, rotating one endpoint per attempt would first dial a
+ * two-node survivor two delays in.
+ */
+ private CompletableFuture<Void> sweepCandidates(
+ List<ConnectionInfo> candidates, int index, int attempt,
RetryPolicy policy) {
+ if (closed) {
+ return CompletableFuture.completedFuture(null);
+ }
+ if (index >= candidates.size()) {
+ return redialAttempt(attempt + 1, policy);
+ }
+ ConnectionInfo target = candidates.get(index);
+ log.info(
+ "Redial attempt {}/{} to {} ({}/{})",
+ attempt,
+ policy.getMaxRetries(),
+ target.serverAddress(),
+ index + 1,
+ candidates.size());
+ return retarget(target)
+ .handle((retargeted, dialError) -> {
+ if (dialError != null) {
+ log.warn("Redial to {} failed: {}",
target.serverAddress(), dialError.getMessage());
+ return sweepCandidates(candidates, index + 1, attempt,
policy);
+ }
+ return replaySignInOn(target, candidates, index, attempt,
policy);
+ })
+ .thenCompose(Function.identity());
+ }
+
+ /**
+ * Re-establishes the session on an endpoint that just came up.
+ *
+ * A sign-in the server rejected -- a rotated password, an expired token --
+ * ends the redial: the connection is up, no other endpoint would answer
+ * differently, and retrying would tear the working connection down on the
+ * next rotation and leave the client connected but unauthenticated anyway.
+ * The rejected credentials are dropped so nothing replays them.
+ */
+ private CompletableFuture<Void> replaySignInOn(
+ ConnectionInfo target, List<ConnectionInfo> candidates, int index,
int attempt, RetryPolicy policy) {
+ return replayLogin()
+ .handle((ok, loginError) -> {
+ if (loginError == null) {
+ log.info("Reconnected to {}", target.serverAddress());
+ return CompletableFuture.<Void>completedFuture(null);
+ }
+ if (isConnectionLoss(unwrap(loginError))) {
log.warn(
- "Redial attempt {} to {} failed: {}",
- attempt,
+ "The sign-in on {} was lost with the
connection: {}",
Review Comment:
`isConnectionLoss` misses `IggyConnectionException` (channel closed before
the reply), so a survivor that dies mid sign-in counts as a rejected login:
remembered login dropped, client published on a dead node, every later call
fails "not authenticated". reproduced. only a server rejection with a
non-transient code is a rejection; everything else should move on.
##########
foreign/csharp/Iggy_SDK/IggyClient/Implementations/TcpMessageStream.cs:
##########
@@ -1026,7 +1063,17 @@ private async Task TryEstablishConnectionAsync(bool
autoLogin, CancellationToken
// trailing segment of a large request until the previous one
is acked.
socket.NoDelay = true;
- await socket.ConnectAsync(host, port, token);
+ // Neither ConnectAsync nor the TLS handshake has a deadline
of its own, and nothing up the
+ // stack adds one: a node whose syns are dropped - or one that
accepts TCP and then never
+ // answers the ClientHello - would hold the sweep, and the
connection semaphore with it, for
+ // the whole kernel connect timeout while a survivor goes
untried.
+ using var dialCancellation = candidates.Length > 1
+ ? CancellationTokenSource.CreateLinkedTokenSource(token)
+ : null;
+ dialCancellation?.CancelAfter(FailoverDialTimeout);
+ var dialToken = dialCancellation?.Token ?? token;
+
+ await socket.ConnectAsync(host, port, dialToken);
dialed = true;
Review Comment:
set before the handshake, so a handshake the 2s bound cuts short lands in
the fatal filter at 1129: sweep aborted, the caller gets
`OperationCanceledException` with a live token, survivor never dialed. set it
after the handshake or use a separate `established` flag.
##########
foreign/node/src/client/client.connection.ts:
##########
@@ -227,6 +245,38 @@ export class IggyConnection extends EventEmitter {
return connectPromise;
}
+ /**
+ * Waits for one dial, bounded while other endpoints are queued behind it.
+ *
+ * A socket has no connect deadline of its own and there is no 'timeout'
+ * listener on it, so a node whose syns are dropped holds the pass for the
+ * whole OS connect timeout -- and it leads every pass, because the current
+ * endpoint only moves on success. The bound matches the Rust SDK's.
+ */
+ private async _dialWithin(socket: Socket, bounded: boolean): Promise<this> {
Review Comment:
the bound covers the `'connect'` event only. on a tls socket that fires
before `'secureConnect'`, so the handshake is still unbounded - rust, go and c#
bound it too.
##########
foreign/node/src/client/client.connection.ts:
##########
@@ -67,6 +67,16 @@ const getTransport = (config: ClientConfig): Socket => {
}
};
+/** One node of the cluster, as a redial candidate. */
+export type Endpoint = { host: string, port: number };
Review Comment:
one inline `{ host: string, port: number }` left at client.socket.ts:551.
##########
core/sdk/src/leader_aware.rs:
##########
@@ -162,13 +217,41 @@ fn process_cluster_metadata(
/// Check if two addresses refer to the same endpoint
/// Handles various formats like 127.0.0.1:8090 vs localhost:8090
-fn is_same_address(addr1: &str, addr2: &str) -> bool {
+///
+/// A host name and the address it resolves to are one endpoint too: a client
+/// configured as `iggy-server:8090` whose roster advertises `10.0.0.5:8090`
+/// would otherwise dial that node twice per failover sweep, and a single-node
+/// deployment would be treated as a cluster. Resolution is the last resort,
+/// only when the spellings differ and at least one side is not a literal
+/// address, and it is a blocking lookup: this runs on the connect and redirect
+/// paths, which are rare and already wait on the network.
+pub(crate) fn is_same_address(addr1: &str, addr2: &str) -> bool {
match (parse_address(addr1), parse_address(addr2)) {
(Some(sock1), Some(sock2)) => sock1.ip() == sock2.ip() && sock1.port()
== sock2.port(),
- _ => normalize_address(addr1) == normalize_address(addr2),
+ (parsed1, parsed2) => {
+ if normalize_address(addr1) == normalize_address(addr2) {
+ return true;
+ }
+ if parsed1.is_some() && parsed2.is_some() {
Review Comment:
unreachable - the `(Some, Some)` arm above already returns.
##########
foreign/go/client/tcp/tcp_core.go:
##########
@@ -994,6 +1009,77 @@ func (c *IggyTcpClient) Connect(ctx context.Context)
error {
return nil
}
+// awaitReestablish waits out what is left of the reestablishAfter window since
+// the last successful connection, if any.
+func (c *IggyTcpClient) awaitReestablish(connectedAt time.Time) {
+ if connectedAt.IsZero() {
+ return
+ }
+
+ elapsed := time.Since(connectedAt)
+ c.logger.Debug("Elapsed time since last connection",
slog.Duration("elapsed", elapsed))
+ if remaining := c.config.reconnection.reestablishAfter - elapsed;
remaining > 0 {
+ c.logger.Info("Trying to connect to the server",
slog.Duration("remaining", remaining))
+ time.Sleep(remaining)
+ }
+}
+
+// dialCandidate brings one endpoint all the way up, wrapping it in TLS when
+// configured, and records the endpoint that answered: the leader check
+// compares against it and the next reconnect starts from it.
+//
+// bounded caps the whole attempt at failoverDialTimeout, for when other
+// endpoints are queued behind this one. Neither the dial nor the handshake has
+// a deadline of its own, and a node whose syns are dropped -- or one that
+// accepts TCP and then never answers the ClientHello -- would hold the sweep
+// for minutes while a survivor goes untried.
+func (c *IggyTcpClient) dialCandidate(ctx context.Context, address string,
bounded bool) (net.Conn, error) {
+ c.logger.Info("Iggy client is connecting to server...",
slog.String("server_address", address))
+ if bounded {
+ var cancel context.CancelFunc
+ ctx, cancel = context.WithTimeout(ctx, failoverDialTimeout)
+ defer cancel()
+ }
+
+ connection, err := (&net.Dialer{}).DialContext(ctx, "tcp", address)
+ if err != nil {
+ c.logger.Error("Failed to establish TCP connection to the
server", slog.Any("error", err))
+ return nil, ierror.ErrCannotEstablishConnection
+ }
+
+ tc := connection.(*net.TCPConn)
+ if err := tc.SetNoDelay(c.config.noDelay); err != nil {
+ c.logger.Error("Failed to set the nodelay option on the client,
continuing...", slog.Any("error", err))
+ }
+
+ established := connection
+ if c.config.tlsEnabled {
+ tlsConfig, err := c.createTLSConfig()
Review Comment:
`createTLSConfig` takes the sni from `c.currentServerAddress`, which during
a sweep is the endpoint just lost, not the candidate being dialed. with
validation on and no `WithTLSDomain` a failover to a host with another name/ip
fails `x509: certificate is valid for 127.0.0.1, not 127.0.0.2`. pass `address`
in. the new tls tests run with validation off so they don't see it. also: a bad
ca path/cert/domain from here is retried forever under the default `maxRetries`
0; rust aborts on those.
##########
foreign/go/client/tcp/tcp_core.go:
##########
@@ -994,6 +1009,77 @@ func (c *IggyTcpClient) Connect(ctx context.Context)
error {
return nil
}
+// awaitReestablish waits out what is left of the reestablishAfter window since
+// the last successful connection, if any.
+func (c *IggyTcpClient) awaitReestablish(connectedAt time.Time) {
+ if connectedAt.IsZero() {
+ return
+ }
+
+ elapsed := time.Since(connectedAt)
+ c.logger.Debug("Elapsed time since last connection",
slog.Duration("elapsed", elapsed))
+ if remaining := c.config.reconnection.reestablishAfter - elapsed;
remaining > 0 {
+ c.logger.Info("Trying to connect to the server",
slog.Duration("remaining", remaining))
+ time.Sleep(remaining)
Review Comment:
`time.Sleep` ignores `ctx` and `Close`. with the rotation a sweep whose
other candidates are refused reaches this endpoint in ms and then sleeps the
rest of `reestablishAfter` in `Connecting` - 60s under a 5s ctx. select on
`ctx.Done()`.
##########
foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java:
##########
@@ -457,10 +494,40 @@ private AsyncTcpConnection openConnection(ConnectionInfo
target) {
heartbeatInterval,
maxVsrFrameSize,
this::retryTransientOnLeader,
- routingState::clearAssignments,
+ this::onSessionReset,
this::onConnectionFailure);
}
+ /**
+ * A server-side eviction reached this client. The routing state it cached
+ * belonged to the evicted session, and a stale-client eviction is the
+ * server ending that session authoritatively, like a logout: the
+ * remembered sign-in ends with it, so only credentials configured on the
+ * builder may bring the session back.
+ */
+ private void onSessionReset(int errorCode) {
+ routingState.clearAssignments();
+ if (errorCode != IggyErrorCode.STALE_CLIENT.getCode()) {
+ return;
+ }
+ // Credentials configured on the builder are what every connect of this
+ // client signs in as, so an eviction does not revoke them and the
+ // client recovers on its own. A sign-in a caller ran is different: the
+ // server ended that session deliberately, and reviving it behind the
+ // caller's back is what an explicit logout must not be able to do
+ // either. Dropped in both places, because the connection replays its
+ // own captured login to bring up a replacement channel.
+ if (username.isPresent() && password.isPresent()) {
+ return;
+ }
+ log.warn("The server evicted this session as stale; the sign-in it ran
will not be replayed");
+ rememberedLogin = null;
Review Comment:
not pinned: removing this line keeps both failover suites green.
##########
foreign/go/client/tcp/tcp_session_management.go:
##########
@@ -165,14 +187,39 @@ func (c *IggyTcpClient) settleOnLeader(ctx
context.Context, code uint32, body []
// endBoundSession logs out a live session before a re-login, so the server
// drops its client-table entry instead of leaving it to be fenced.
+//
+// The logout runs connect-scoped, and a failure it could recover from is
+// swallowed. Both because this call holds registerMtx: a logout that entered
+// the reconnect path would reconnect, sign in with the remembered credentials,
+// and deadlock on that lock. There is nothing to salvage either way -- a
+// session whose logout cannot be delivered died with its socket, and the
+// server fences what it left behind -- and the sign-in that follows replays
+// through its own reconnect.
func (c *IggyTcpClient) endBoundSession(ctx context.Context) error {
c.mtx.Lock()
bound := c.session.Bound()
c.mtx.Unlock()
if !bound {
return nil
}
- return c.LogoutUser(ctx)
+
+ err := c.LogoutUser(context.WithValue(ctx, connectScoped{}, struct{}{}))
Review Comment:
still deadlocks when the server answers the logout with
`TransientNotAccepted` (it does that when it's not primary): `sendFrame`
follows the redirect with `Connect` (tcp_core.go:646) without
`skipAutoLoginOnce`, so `establishSession` -> `LoginUser` -> `register` waits
on the `registerMtx` this goroutine holds. reproduced with a two-node fake. set
the skip before that `Connect`.
##########
foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java:
##########
@@ -732,7 +874,35 @@ CompletableFuture<Optional<ConnectionInfo>>
findLeaderElsewhere(ConnectionInfo c
if (currentSystemClient == null) {
return CompletableFuture.completedFuture(Optional.empty());
}
- return
LeaderAwareness.findLeaderElsewhere(currentSystemClient::getClusterMetadata,
currentTarget);
+ return
LeaderAwareness.findLeaderElsewhere(currentSystemClient::getClusterMetadata,
currentTarget)
+ .thenApply(lookup -> {
+ // Replaced wholesale rather than merged: the roster is the
+ // cluster's own answer about where its nodes are, so a
node
+ // it dropped stops being dialed. The configured seed is
+ // kept separately and outlives it.
+ if (!lookup.endpoints().isEmpty()) {
Review Comment:
untested: assigning `lookup.endpoints()` unconditionally keeps everything
green.
##########
foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java:
##########
@@ -638,6 +777,9 @@ CompletableFuture<IdentityInfo>
loginOnLeader(Supplier<CompletableFuture<Identit
CompletableFuture<IdentityInfo> callerFuture = new
CompletableFuture<>();
transaction.whenComplete((identity, error) -> {
gate.complete(null);
+ if (error == null) {
+ rememberedLogin = loginAttempt;
Review Comment:
a login still in flight when `close()` nulled this sets it again, so a later
`connect()` + loss replays it.
##########
foreign/csharp/Iggy_SDK/IggyClient/Implementations/TcpMessageStream.cs:
##########
@@ -1095,6 +1146,20 @@ private async Task TryEstablishConnectionAsync(bool
autoLogin, CancellationToken
_logger.LogError(e, "Failed to connect");
+ // Every other endpoint gets its turn before the retry delay:
the node just lost may be gone for
+ // good, and pausing on it helps nothing.
+ if (++candidate < candidates.Length)
+ {
+ _currentAddress = candidates[candidate];
+ continue;
+ }
+
+ candidate = 0;
+ _currentAddress = candidates[0];
+
+ // The sweep is what the reconnection budget applies to, not a
single dial: checked per dial,
Review Comment:
no test for this: moving the checks back above the advance keeps the suite
green.
##########
foreign/node/src/client/client.socket.ts:
##########
@@ -186,21 +203,24 @@ export class CommandResponseStream extends EventEmitter {
if (!this.isAuthenticated && !this.isUnloggedCommand(command))
await this.authenticate(this.options.credentials);
- const response = await new Promise<CommandResponse>(
- (resolve, reject) => {
- const job = {
- command,
- payload,
- handleResponse,
- resolve,
- reject
- };
- if (last)
- this._execQueue.push(job);
- else
- this._execQueue.unshift(job);
- this._processQueue();
- });
+ // The roster read is itself a queued command and the queue is
+ // single-flighted, so the leader re-check cannot happen inside
+ // `_processVsr`. The refusal comes back out here instead, where the
+ // queue is free, and the command is re-issued on the node that now
+ // leads.
+ let response: CommandResponse;
+ for (let move = 0; ; move += 1) {
+ try {
+ response = await this._queueCommand(command, payload, handleResponse,
+ last);
+ break;
+ } catch (error) {
+ if (!(error instanceof LeaderMovedError))
+ throw error;
+ if (move >= MAX_LEADER_REDIRECTS || !await this._followLeaderMove())
+ throw responseError(command, error.refusal.errorCode);
Review Comment:
when the roster still names this node or has no healthy leader the command
now fails with 58 after one 2s window. before it was replayed for 30s, and rust
still keeps replaying on the same connection and re-checking every 2s until the
30s budget. keep that.
##########
foreign/node/src/client/client.socket.ts:
##########
@@ -186,21 +203,24 @@ export class CommandResponseStream extends EventEmitter {
if (!this.isAuthenticated && !this.isUnloggedCommand(command))
await this.authenticate(this.options.credentials);
Review Comment:
after `_followLeaderMove()` the client has no session: `redirect()` emits
`'disconnected'` which resets it, and this step sits above the loop. the
retried replicated command fails `Unauthenticated` client-side, a
non-replicated one goes out with session 0. authenticate again inside the loop
after a move.
##########
foreign/csharp/Iggy_SDK/IggyClient/Implementations/TcpMessageStream.cs:
##########
@@ -1186,7 +1300,10 @@ private async Task<Stream>
CreateSslStreamAndAuthenticate(Socket socket, TlsSett
var stream = new NetworkStream(socket, true);
var sslStream = new SslStream(stream, false,
RemoteCertificateValidationCallback);
- await sslStream.AuthenticateAsClientAsync(tlsSettings.Hostname);
+ // The token carries the dial bound when other endpoints are queued
behind this one: a peer that
Review Comment:
the `SslStream` is never disposed when the handshake fails.
##########
foreign/node/src/client/client.socket.ts:
##########
@@ -216,6 +236,57 @@ export class CommandResponseStream extends EventEmitter {
}
}
+ private _queueCommand(
+ command: number,
+ payload: Buffer,
+ handleResponse: boolean,
+ last: boolean
+ ): Promise<CommandResponse> {
+ return new Promise<CommandResponse>((resolve, reject) => {
+ const job = {
+ command,
+ payload,
+ handleResponse,
+ resolve,
+ reject
+ };
+ if (last)
+ this._execQueue.push(job);
+ else
+ this._execQueue.unshift(job);
+ this._processQueue();
+ });
+ }
+
+ /**
+ * Re-reads the roster and moves to the leader it names.
+ *
+ * @returns Whether the client moved, so the refused request is worth
+ * re-issuing
+ */
+ private async _followLeaderMove(): Promise<boolean> {
Review Comment:
the roster read goes through `sendCommand` and hits the same
`LeaderMovedError` loop with no depth guard. also, concurrent commands during a
move: the first redirect's `'disconnected'` fails the others' roster reads, so
they fail with 58 instead of being re-issued.
##########
foreign/node/src/client/client.connection.ts:
##########
@@ -300,36 +350,60 @@ export class IggyConnection extends EventEmitter {
): Promise<this> {
let lastError = initialError;
let expectedSocket = this.socket;
- let attempt = 0;
+ let firstPass = true;
while (enabled && this.reconnectCount < maxRetries) {
this.connecting = true;
this.reconnectCount += 1;
- await waitForReconnect(interval);
+ const candidates = this._redialCandidates();
+ // The backoff paces retries against a single endpoint. With other
+ // endpoints known there is somewhere else to go, and pausing first only
+ // pushes the failover past the interval a caller is willing to wait;
+ // later passes still back off.
+ if (!firstPass || candidates.length === 1)
+ await waitForReconnect(interval);
+ firstPass = false;
if (this.ending)
throw new Error('connection is closed', { cause: lastError });
// A redirect may replace the socket at any point. Defer to the active
// connection instead of dialing the superseded endpoint.
if (this.connected || this.socket !== expectedSocket)
return this.connect();
- const options = this._reconnectTarget(attempt);
- attempt += 1;
- const socket = this._installSocket(
- getTransport({ ...this.config, options })
- );
- this.socket = socket;
- expectedSocket = socket;
- try {
- await this._waitForConnection(socket);
- if (this.socket !== socket)
+ // Every endpoint gets its turn inside one attempt, so a full pass over
+ // the cluster costs one retry rather than one per endpoint: a pass that
+ // stopped at the first refusal would never reach the survivors of a
+ // client configured for a single retry.
+ for (const options of candidates) {
+ // Re-checked every iteration, not once above the loop: a destroy or a
+ // redirect mid-pass has to stop the pass. Left running, the next
+ // endpoint that answers would leave an open socket nobody closes, a
+ // 'connect' event after the destroy, and the process alive.
+ if (this.ending)
+ throw new Error('connection is closed', { cause: lastError });
+ if (this.socket !== expectedSocket)
return this.connect();
- this.config.options = options;
- return this;
- } catch (error) {
- lastError = error instanceof Error
- ? error
- : new Error(String(error));
- debug('reconnect attempt failed', lastError);
+
+ const socket = this._installSocket(
+ getTransport({ ...this.config, options })
+ );
+ this.socket = socket;
+ expectedSocket = socket;
+ try {
+ await this._dialWithin(socket, candidates.length > 1);
Review Comment:
no test covers the bound - `_dialWithin(socket, false)` keeps the suite
green. same for the first-pass skip at 357.
##########
foreign/node/src/client/client.connection.ts:
##########
@@ -67,6 +67,16 @@ const getTransport = (config: ClientConfig): Socket => {
}
};
+/** One node of the cluster, as a redial candidate. */
+export type Endpoint = { host: string, port: number };
+
+/**
+ * Bound on one dial while other endpoints are queued behind it. A socket has
+ * no connect deadline of its own, so a node whose syns are dropped would hold
+ * the whole pass. Matches the Rust SDK.
+ */
+const FAILOVER_DIAL_TIMEOUT_MS = 2_000;
+
/**
* Default reconnection settings.
Review Comment:
the `DefaultReconnectOption` comment below and `maxRetries` in
`client.type.ts:105` still say attempts; one retry is now a full pass over
every candidate, and the first pass with a roster doesn't wait.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]