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


##########
foreign/go/client/tcp/tcp_core.go:
##########
@@ -881,38 +931,73 @@ func (c *IggyTcpClient) Connect(ctx context.Context) 
error {
                c.logger.Debug("Client is already connected.", 
slog.String("client_address", clientAddress))
                return nil
        case iggcon.TransportStateConnecting:
+               inFlight := c.connectInFlight
                c.mtx.Unlock()
-               c.logger.Debug("Client is already connecting.")
-               return nil
+               c.logger.Debug("Client is already connecting; waiting for that 
attempt.")
+               if inFlight == nil {
+                       return nil
+               }
+               select {

Review Comment:
   this wait can deadlock. `register` holds `registerMtx` and on the redirect 
and logout paths calls `disconnect()` then `Connect(suppressAutoLogin(ctx))`. a 
plain request that slips into `Connect` between those two becomes the attempt 
owner with auto-login on (exchange only suppresses when `login`), runs 
`establishSession` -> `LoginUser` -> `register` and blocks on `registerMtx`, 
while the sign-in goroutine sits here waiting on that attempt. only 
`ctx.Done()` or `Close()` break the cycle and callers pass `Background`. before 
this change the `Connecting` case returned nil, so there was no cycle. a caller 
whose ctx carries `skipAutoLogin` must not wait on a foreign attempt, or the 
owner's auto-login has to run outside `registerMtx`.



##########
foreign/go/client/tcp/tcp_core.go:
##########
@@ -881,38 +931,73 @@ func (c *IggyTcpClient) Connect(ctx context.Context) 
error {
                c.logger.Debug("Client is already connected.", 
slog.String("client_address", clientAddress))
                return nil
        case iggcon.TransportStateConnecting:
+               inFlight := c.connectInFlight
                c.mtx.Unlock()
-               c.logger.Debug("Client is already connecting.")
-               return nil
+               c.logger.Debug("Client is already connecting; waiting for that 
attempt.")
+               if inFlight == nil {
+                       return nil
+               }
+               select {
+               case <-inFlight:
+               case <-ctx.Done():
+                       return ctx.Err()
+               case <-c.closed:
+                       return ierror.ErrClientShutdown
+               }
+               c.mtx.Lock()
+               attemptErr := c.connectErr

Review Comment:
   read the error from the attempt you waited on, not the shared field. if a 
new attempt started after the channel closed, `connectErr` is already nil and 
this returns nil while the client is `Connecting`, so the caller's next request 
fails `ErrNotConnected`. keep `{done, err}` per attempt and read the snapshot.



##########
foreign/go/client/tcp/tcp_core.go:
##########
@@ -975,25 +1056,147 @@ func (c *IggyTcpClient) Connect(ctx context.Context) 
error {
        }
 
        c.mtx.Lock()
+       if c.transportState != iggcon.TransportStateConnecting {
+               // Disconnected or shut down while this attempt was dialing: the
+               // connection it just made is not wanted, and installing it 
would
+               // resurrect a client somebody asked to stop.
+               state := c.transportState
+               c.mtx.Unlock()
+               _ = conn.Close()
+               c.logger.Debug("The connect was superseded while dialing; 
dropping the connection.")
+               if state == iggcon.TransportStateShutdown {
+                       return ierror.ErrClientShutdown
+               }
+               return ierror.ErrNotConnected

Review Comment:
   state can be `Connected` here because another attempt won, and this hands 
`ErrNotConnected` to a caller whose client is connected. share the success 
instead of treating every non-`Connecting` state as 'asked to stop'. the 
`previous.Close()` below is unreachable - every path into `Connecting` starts 
from `Disconnected`, which already nils `conn`.



##########
foreign/go/client/tcp/tcp_core.go:
##########
@@ -975,25 +1056,147 @@ func (c *IggyTcpClient) Connect(ctx context.Context) 
error {
        }
 
        c.mtx.Lock()
+       if c.transportState != iggcon.TransportStateConnecting {
+               // Disconnected or shut down while this attempt was dialing: the
+               // connection it just made is not wanted, and installing it 
would
+               // resurrect a client somebody asked to stop.
+               state := c.transportState
+               c.mtx.Unlock()
+               _ = conn.Close()
+               c.logger.Debug("The connect was superseded while dialing; 
dropping the connection.")
+               if state == iggcon.TransportStateShutdown {
+                       return ierror.ErrClientShutdown
+               }
+               return ierror.ErrNotConnected
+       }
+       // A connection installed over another one leaks its socket: two 
Connects
+       // can race here, and the loser's conn would otherwise stay open with
+       // nothing left holding it.
+       if previous := c.conn; previous != nil {
+               _ = previous.Close()
+       }
        c.conn = conn
        c.reader = bufio.NewReaderSize(conn, 64*1024)
        c.transportState = iggcon.TransportStateConnected
        c.connectedAt = time.Now()
        // The server fence does not survive the old socket, so the new 
connection
        // starts from a fresh client identity.
        c.session.Reset()
-       skipAutoLogin := c.skipAutoLoginOnce
-       c.skipAutoLoginOnce = false
-       c.logger.Info("Iggy client has connected to the Iggy server", 
slog.String("client_address", c.clientAddress), slog.String("server_address", 
c.currentServerAddress))
+       clientAddress := c.clientAddress
+       serverAddress := c.currentServerAddress
        c.mtx.Unlock()
+       c.logger.Info("Iggy client has connected to the Iggy server",
+               slog.String("client_address", clientAddress),
+               slog.String("server_address", serverAddress))
 
-       if err := c.establishSession(ctx, skipAutoLogin); err != nil {
+       if err := c.establishSession(ctx, ctx.Value(skipAutoLogin{}) != nil); 
err != nil {
                _ = c.disconnect()
                return err
        }
        return nil
 }
 
+// isTLSConfigFault reports whether a dial failed for a reason that says the
+// client's own TLS configuration is wrong -- an unreadable or unparsable CA
+// file, a domain that yields no server name, or a certificate this client will
+// never accept. None of those change on a retry.
+func isTLSConfigFault(err error) bool {
+       if errors.Is(err, ierror.ErrInvalidTlsCertificatePath) ||
+               errors.Is(err, ierror.ErrInvalidTlsCertificate) ||
+               errors.Is(err, ierror.ErrInvalidTlsDomain) {
+               return true
+       }
+
+       var certificateError *tls.CertificateVerificationError
+       var recordError tls.RecordHeaderError
+       return errors.As(err, &certificateError) || errors.As(err, &recordError)

Review Comment:
   `tls.RecordHeaderError` makes the whole connect unrecoverable, so one 
plaintext endpoint in a mixed roster ends the sweep even when the others are 
only down for a moment. a bad cert is a config fault, a plaintext peer is one 
bad endpoint - skip it and keep retrying the rest.



##########
foreign/go/client/tcp/tcp_core.go:
##########
@@ -505,22 +539,20 @@ func (c *IggyTcpClient) exchange(ctx context.Context, 
code uint32, frame []byte)
        if disconnectErr := c.disconnect(); disconnectErr != nil {

Review Comment:
   two requests failing at once still start two attempts: `disconnect()` does 
not early-out on `Connecting`, so the second one stomps the first's state, 
takes the `default` branch in `Connect` and overwrites `connectInFlight` - the 
first attempt's defer then closes the second's channel with the wrong error. 
skip the disconnect (or make it a no-op) while another attempt is connecting.



##########
foreign/go/client/tcp/tcp_core.go:
##########
@@ -616,7 +648,20 @@ func (c *IggyTcpClient) sendFrame(ctx context.Context, 
code uint32, frame []byte
                                return nil, redirectErr
                        }
                        if redirect {
-                               if connectErr := c.Connect(ctx); connectErr != 
nil {
+                               // A connect-scoped request is issued from 
inside the sign-in
+                               // transaction, which holds registerMtx: the 
automatic sign-in
+                               // on the reconnect path would wait on that 
lock forever. The
+                               // transaction signs in itself on the node it 
lands on, so the
+                               // reconnect must not.
+                               redirectCtx := ctx
+                               if ctx.Value(connectScoped{}) != nil {
+                                       // Issued from inside the sign-in 
transaction, which holds

Review Comment:
   same paragraph as the one right above it.



##########
foreign/node/src/client/client.socket.ts:
##########
@@ -153,12 +186,40 @@ export class CommandResponseStream extends EventEmitter {
     });
     this.connection.on('disconnected', () => {
       this._resetSession();
+      if (this.connection.redirecting) {
+        // The client is moving to the leader, which is its own doing: a queued
+        // command has not been written, so it belongs on the node being moved
+        // to rather than in an error.
+        this._reissueQueue();
+        return;
+      }
       this._failQueue(
         new Error('connection closed before queued commands were sent')
       );
     });
   }
 
+  /**
+   * Re-submits queued commands through the full send path, so each one
+   * reconnects, re-authenticates and re-checks the leader as if it had just
+   * been called.
+   *
+   * Only for a drop the client caused. Nothing here was written, so there is 
no
+   * outcome in doubt: a command still in the queue when the socket is replaced
+   * would otherwise fail with a lost-connection error the caller can do 
nothing
+   * about.
+   */
+  private _reissueQueue(): void {
+    const queued = this._execQueue;
+    this._execQueue = [];
+    for (const job of queued) {
+      debug('re-issuing a queued command after a leader move', job.command);
+      this.sendCommand(job.command, job.payload, {

Review Comment:
   this drops `job.deadline` and `sendCommand` mints a fresh 30s, so a command 
caught in a move can take 60s - the comment at line 263 promises the opposite. 
it also loses `last` and `followsLeaderMoves`: a roster read queued with 
`followsLeaderMoves: false` comes back with it on, which is the recursion the 
check above guards against. carry the whole job through.



##########
foreign/node/src/client/client.socket.ts:
##########
@@ -228,15 +411,22 @@ export class CommandResponseStream extends EventEmitter {
     while (this._execQueue.length > 0 && this.connection.socket.writable) {

Review Comment:
   the realistic case is not covered: when the roster read resolves this loop 
keeps going synchronously and writes the next queued job to the old socket 
before `_followLeaderMove` reaches `redirect()`. that job is then in flight, 
dies with the plain 'connection closed while waiting for response' error and 
resets the session. the `redirecting` check at 425 is also unpinned - swapping 
it back to `_failQueue` keeps the suite green.



##########
foreign/node/src/client/client.connection.test.ts:
##########
@@ -393,6 +393,309 @@ describe('IggyConnection', () => {
     }
   );
 
+  it('rotates a redial through the roster it learned while connected',
+    async () => {
+      const seed = await startServer();
+      const seedPort = (seed.address() as AddressInfo).port;
+      const connection = new IggyConnection(connectionConfig(seed));
+      connection.on('error', () => undefined);
+      try {
+        connection.rememberRoster([
+          { host: '127.0.0.1', port: seedPort },
+          { host: '127.0.0.1', port: seedPort + 1 },
+          { host: '127.0.0.1', port: seedPort + 2 }
+        ]);
+        // The endpoint the client is on leads, the roster follows, and the
+        // roster's copy of that endpoint does not earn a second attempt.
+        assert.deepEqual(
+          connection._redialCandidates().map((options) => options.port),
+          [seedPort, seedPort + 1, seedPort + 2]
+        );
+      } finally {
+        connection._destroy();
+        await new Promise<void>((resolve) => seed.close(() => resolve()));
+      }
+    }
+  );
+
+  it('dials the endpoint it is on, then the seed, then the roster',
+    async () => {
+      const seed = await startServer();
+      const seedPort = (seed.address() as AddressInfo).port;
+      const connection = new IggyConnection(connectionConfig(seed));
+      connection.on('error', () => undefined);
+      try {
+        // A redirect moves the client off its seed; the seed is still the one
+        // endpoint the caller vouched for, so it comes before a roster the
+        // cluster may have reshaped since.
+        connection.config.options = {
+          ...connection.config.options,
+          port: seedPort + 9
+        };
+        connection.rememberRoster([{ host: '127.0.0.1', port: seedPort + 5 }]);
+
+        assert.deepEqual(
+          connection._redialCandidates().map((options) => options.port),
+          [seedPort + 9, seedPort, seedPort + 5]
+        );
+      } finally {
+        connection._destroy();
+        await new Promise<void>((resolve) => seed.close(() => resolve()));
+      }
+    }
+  );
+
+  it('counts endpoints that only differ in spelling once',
+    async () => {
+      const seed = await startServer();
+      const seedPort = (seed.address() as AddressInfo).port;
+      const connection = new IggyConnection(connectionConfig(seed));
+      connection.on('error', () => undefined);
+      try {
+        // The loopback aliases and an IPv4-mapped address all name the 
endpoint
+        // the client is already on, so none of them earns a dial of its own.
+        connection.rememberRoster([
+          { host: 'localhost', port: seedPort },
+          { host: '::1', port: seedPort },
+          { host: '::ffff:127.0.0.1', port: seedPort },
+          { host: '127.0.0.1', port: seedPort + 1 }
+        ]);
+
+        assert.deepEqual(
+          connection._redialCandidates().map((options) => options.port),
+          [seedPort, seedPort + 1]
+        );
+      } finally {
+        connection._destroy();
+        await new Promise<void>((resolve) => seed.close(() => resolve()));
+      }
+    }
+  );
+
+  it('does not redial at all when reconnection is disabled',
+    async () => {
+      // `enabled: false` is what a caller says to opt out. The retry budget is
+      // whatever the defaults hold, so a loop that reads it without checking
+      // this flag would run every one of those passes - and with the backoff
+      // gated on the same flag, back to back.
+      //
+      // The endpoint accepts and hangs up, so the drop that would start a
+      // redial happens and every dial of it is counted.
+      const hangup = await startServer();
+      const hangupPort = (hangup.address() as AddressInfo).port;
+      let accepted = 0;
+      hangup.on('connection', (socket) => {
+        accepted += 1;
+        socket.destroy();
+      });
+
+      const connection = new IggyConnection({
+        transport: 'TCP',
+        options: { host: '127.0.0.1', port: hangupPort },
+        credentials: { username: 'iggy', password: 'iggy' },
+        reconnect: { enabled: false, interval: 10, maxRetries: 12 },
+        maxResponseFrameSize: FRAME_LIMIT
+      });
+      connection.on('error', () => undefined);
+      try {
+        await connection.connect().catch(() => undefined);
+        await new Promise<void>((resolve) => setTimeout(resolve, 300));
+
+        assert.equal(accepted, 1,
+          'a client that turned reconnection off redialed anyway'
+        );
+        assert.equal(connection.connected, false);
+      } finally {
+        connection._destroy();
+        await new Promise<void>((resolve) => hangup.close(() => resolve()));
+      }
+    }
+  );
+
+  it('sweeps the endpoints it knows once when reconnection is disabled',
+    async () => {
+      // Opting out of retries is not opting out of the endpoints: with more
+      // than one known, they get exactly one pass and no backoff, as in the
+      // other SDKs.
+      const dead = await startServer();
+      const deadPort = (dead.address() as AddressInfo).port;
+      await new Promise<void>((resolve) => dead.close(() => resolve()));
+      const live = await startServer();
+      const livePort = (live.address() as AddressInfo).port;
+      let accepted = 0;
+      live.on('connection', () => { accepted += 1; });
+
+      const connection = new IggyConnection({
+        transport: 'TCP',
+        options: { host: '127.0.0.1', port: deadPort },
+        credentials: { username: 'iggy', password: 'iggy' },
+        reconnect: { enabled: false, interval: 10, maxRetries: 12 },
+        maxResponseFrameSize: FRAME_LIMIT
+      });
+      connection.on('error', () => undefined);
+      try {
+        connection.rememberRoster([{ host: '127.0.0.1', port: livePort }]);
+        await connection.connect().catch(() => undefined);
+        await new Promise<void>((resolve) => setTimeout(resolve, 200));
+
+        assert.equal(accepted, 1,

Review Comment:
   green without the `enabled &&` guard too - pass 1 lands on the live node 
either way. the case that regressed was every endpoint down with reconnection 
off (12 back-to-back passes) and that one still has no test.



##########
foreign/node/src/client/client.socket.test.ts:
##########
@@ -864,6 +992,224 @@ describe('VSR client socket', () => {
     }
   });
 
+  it('re-issues every request refused by a demoted node, not just the first',
+    async () => {
+      // One demotion, several refused requests: each of them re-checking on
+      // its own would move the client once per request, and the first
+      // redirect's drop would fail the others' roster reads - reporting a
+      // refusal they never had to.
+      const leader = await startVsrServer((frame, socket) => {
+        singleNodeHandler(leader.port)(frame, socket);
+      });
+      // Leader at login, so the settlement leaves the client here, then
+      // demoted: the refusals below are what tells the client to look again.
+      let demotedYet = false;
+      const demoted = await startVsrServer((frame, socket) => {
+        const code = frame.readUInt32LE(REQUEST_OFFSET.reserved);
+        if (code === COMMAND_CODE.GetClusterMetadata) {
+          socket.write(replyFrame(
+            Operation.NonReplicated,
+            demotedYet
+              ? twoNodeMetadataBody(demoted.port, leader.port)
+              : twoNodeMetadataBody(leader.port, demoted.port)
+          ));
+          return;
+        }
+        if (code === 60_032) {
+          socket.write(replyFrame(Operation.NonReplicated, Buffer.alloc(0), 
58));
+          return;
+        }
+        singleNodeHandler(demoted.port)(frame, socket);
+      });
+
+      const client = new CommandResponseStream(vsrConfig(demoted.port));
+      try {
+        await client.authenticate(vsrConfig(demoted.port).credentials);
+        demotedYet = true;
+        const rosterReadsBefore = demoted.frames.filter(
+          (frame) => frame.readUInt32LE(REQUEST_OFFSET.reserved) ===
+            COMMAND_CODE.GetClusterMetadata
+        ).length;
+
+        // Two refusals, one re-check: the second caller shares the move the
+        // first started rather than starting its own or being told 58.
+        const moves = await Promise.all([

Review Comment:
   this calls `_followLeaderMove` directly and never sends a command, so the 
`60_032 -> 58` branch is dead and nothing is re-issued despite the title. send 
two real commands that both get 58 and assert one roster read and two re-issues.



##########
foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClientTransientFailoverTest.java:
##########
@@ -61,6 +64,7 @@ class AsyncIggyTcpClientTransientFailoverTest {
     private static final int OPERATION_CREATE_STREAM = 128;
     private static final int GET_CLUSTER_METADATA_CODE = 12;
     private static final int CREATE_STREAM_CODE = 202;
+    private static final int PING_CODE = 1;

Review Comment:
   unused.



##########
foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java:
##########
@@ -732,7 +919,62 @@ 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 -> {
+                    rememberRoster(lookup);
+                    return lookup.redirect();
+                });
+    }
+
+    /**
+     * Keeps what a leader check learned about where the cluster's nodes are.
+     *
+     * Replaced wholesale rather than merged: the roster is the cluster's own
+     * answer, so a node it dropped stops being dialed. The configured seed is
+     * kept separately and outlives it.
+     *
+     * An inconclusive check -- an unreadable roster, a metadata read that
+     * failed -- names no endpoint, and that must leave the last roster
+     * standing: assigning it anyway would empty the redial candidates exactly
+     * when the cluster is unreachable, which is when they are needed.
+     */
+    void rememberRoster(LeaderAwareness.LeaderLookup lookup) {
+        if (!lookup.endpoints().isEmpty()) {
+            rosterTargets = lookup.endpoints();
+        }
+    }
+
+    /** The roster this client would redial, for tests in this package. */
+    List<ConnectionInfo> rosterTargets() {
+        return rosterTargets;
+    }
+
+    /** Whether a sign-in is remembered for replay, for tests in this package. 
*/
+    boolean hasRememberedLogin() {
+        return rememberedLogin != null;
+    }
+
+    /** The live connection, for tests in this package. */
+    AsyncTcpConnection currentConnection() {

Review Comment:
   no callers left since the test that used it went away.



##########
core/common/src/types/configuration/tcp_config/tcp_client_config_builder.rs:
##########
@@ -44,6 +45,13 @@ impl TcpClientConfigBuilder {
         self
     }
 
+    /// Sets the addresses of other nodes of the same cluster, dialed in order
+    /// when `server_address` cannot be reached.
+    pub fn with_failover_addresses(mut self, failover_addresses: Vec<String>) 
-> Self {

Review Comment:
   still the only way to set seeds: `IggyClientBuilder::with_tcp()` has no 
`with_failover_addresses`, the connection string has no option, 
`client_provider.rs:140` hardcodes an empty vec so the cli cannot, and python 
has nothing. in scope or a follow-up? say which in the PR description.



##########
foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncTcpConnection.java:
##########
@@ -827,10 +831,14 @@ private void handlePostResponse(Channel channel, int 
commandCode, boolean isLogi
      * channel. Bumping the generation makes the replacement channel re-run
      * login and Register. The fresh session invalidates cached routing state
      * such as consumer-group assignments.
+     *
+     * The reason travels to the listener, which owns the question of whether
+     * the session may be re-established at all: only it knows whether the

Review Comment:
   the listener no longer decides anything - the session is kept whichever way 
the sign-in was made. same for the comment in `VsrResponseHandler` at line 166.



##########
foreign/csharp/Iggy_SDK_Tests/VsrTests/EndpointFailoverTests.cs:
##########
@@ -0,0 +1,570 @@
+// 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.
+
+using System.Buffers.Binary;
+using System.Net;
+using System.Net.Sockets;
+using System.Text;
+using Apache.Iggy.Configuration;
+using Apache.Iggy.Contracts.Tcp;
+using Apache.Iggy.Enums;
+using Apache.Iggy.Exceptions;
+using Apache.Iggy.IggyClient;
+using Apache.Iggy.IggyClient.Implementations;
+using Apache.Iggy.Vsr;
+using Microsoft.Extensions.Logging.Abstractions;
+
+namespace Apache.Iggy.Tests.VsrTests;
+
+/// <summary>
+///     The node a client signed in on dies; its next request has to complete 
on a survivor the roster named,
+///     under a session established there. Mirrors
+///     <c>core/integration/tests/cluster/failover_client_continuity.rs</c>.
+/// </summary>
+public sealed class EndpointFailoverTests
+{
+    private const int HeaderSize = 256;
+    private const int SizeOffset = 48;
+    private const int CommandOffset = 60;
+    private const int RequestIdOffset = 168;
+    private const int RequestOperationOffset = 176;
+    private const int RequestReservedOffset = 196;
+    private const int ReplyRequestIdOffset = 200;
+    private const int ReplyOperationOffset = 208;
+    private const int ReplyStatusOffset = 216;
+
+    private const byte CommandReply = 8;
+    private const byte CommandEviction = 13;
+    private const int EvictionReasonOffset = 255;
+    private const byte EvictionStaleClient = 13;
+    private const byte OperationRegister = 1;
+    private const byte OperationNonReplicated = 2;
+    private const int GetClusterMetadataCode = 12;
+    private const int PingCode = 1;
+
+    [Fact]
+    public async Task ResumesOnASurvivorAfterTheSignedInNodeDies()
+    {
+        using var primary = new MockNode();
+        using var survivor = new MockNode();
+
+        // The primary leads, so the sign-in settles there and the roster is 
only remembered - not acted on -
+        // until the node dies.
+        primary.Serve(request => request.Code == GetClusterMetadataCode
+            ? Reply(OperationNonReplicated, ClusterMetadata(primary.Port, 
survivor.Port, primary.Port))
+            : Answer(request));
+        survivor.Serve(request => request.Code == GetClusterMetadataCode
+            ? Reply(OperationNonReplicated, ClusterMetadata(primary.Port, 
survivor.Port, survivor.Port))
+            : Answer(request));
+
+        var configuration = new IggyClientConfigurator
+        {
+            BaseAddress = $"127.0.0.1:{primary.Port}",
+            Protocol = Protocol.Tcp,
+            ReconnectionSettings = new ReconnectionSettings
+            {
+                Enabled = true,
+                MaxRetries = 4,
+                InitialDelay = TimeSpan.FromMilliseconds(20)
+            }
+        };
+        using var client = new TcpMessageStream(configuration, 
NullLoggerFactory.Instance);
+
+        await client.ConnectAsync(TestContext.Current.CancellationToken);
+        // No auto login: the credentials come from the caller's own sign-in, 
which is the shape that could not
+        // reconnect at all before.
+        await client.LoginUserAsync("iggy", "iggy", 
TestContext.Current.CancellationToken);
+        await client.PingAsync(TestContext.Current.CancellationToken);
+        Assert.Equal(1, primary.Pings);
+
+        primary.Kill();
+
+        // The request in flight when the node died is allowed to fail; what 
is not allowed is never completing
+        // one, which is what a client that only knows the dead endpoint does.
+        var (resumed, lastError) = await ResumedWithin(client, 
TimeSpan.FromSeconds(10));
+        Assert.True(resumed,
+            $"the client has to resume on the survivor the roster named 
({lastError}, survivor saw " +
+            $"{survivor.Registrations} registrations and {survivor.Pings} 
pings)");
+        Assert.True(survivor.Registrations >= 1, "the remembered credentials 
signed in again on the survivor");
+        Assert.True(survivor.Pings >= 1, "the request landed on the survivor");
+    }
+
+    /// <summary>
+    ///     Mirrors the integration contract (HeartbeatTests
+    ///     EvictedClient_WithoutAutoLogin_Should_ReestablishItsSession): an 
eviction comes off the server's
+    ///     heartbeat timer rather than from the caller, so the sign-in this 
client remembered survives it and
+    ///     the reconnect re-establishes the session. Same rule in every SDK.
+    /// </summary>
+    [Fact]
+    public async Task ServerEvictionReplaysTheRememberedSignIn()
+    {
+        using var node = new MockNode();
+        var evict = false;
+        node.Serve(request =>
+        {
+            if (request.Operation == OperationRegister)
+            {
+                return Reply(OperationRegister, RegisterBody(session: 128));
+            }
+
+            if (evict)
+            {
+                evict = false;
+                return EvictionFrame(EvictionStaleClient);
+            }
+
+            return Reply(OperationNonReplicated, request.Code == 
GetClusterMetadataCode
+                ? ClusterMetadata(node.Port, node.Port, node.Port)
+                : []);
+        });
+
+        var configuration = new IggyClientConfigurator
+        {
+            BaseAddress = $"127.0.0.1:{node.Port}",
+            Protocol = Protocol.Tcp,
+            ReconnectionSettings = new ReconnectionSettings
+            {
+                Enabled = true,
+                MaxRetries = 2,
+                InitialDelay = TimeSpan.FromMilliseconds(20)
+            }
+        };
+        using var client = new TcpMessageStream(configuration, 
NullLoggerFactory.Instance);
+
+        await client.ConnectAsync(TestContext.Current.CancellationToken);
+        await client.LoginUserAsync("iggy", "iggy", 
TestContext.Current.CancellationToken);
+        await client.PingAsync(TestContext.Current.CancellationToken);
+        var registrationsBeforeEviction = node.Registrations;
+
+        // A ping is replay-safe, so the eviction is absorbed: the reconnect 
signs in again with the
+        // remembered credentials and the request completes over the session 
it re-established.
+        evict = true;
+        await client.PingAsync(TestContext.Current.CancellationToken);
+
+        Assert.True(node.Registrations > registrationsBeforeEviction,
+            "the reconnect signed in again with the remembered credentials");
+        await client.PingAsync(TestContext.Current.CancellationToken);
+    }
+
+    /// <summary>
+    ///     The same rule when the eviction lands on a replicated write: that 
request is reported as
+    ///     outcome-unknown, because its own outcome is unknown, but the 
session behind it is still
+    ///     re-established for the requests that follow.
+    /// </summary>
+    [Fact]
+    public async Task 
ServerEvictionDuringAReplicatedWriteReplaysTheRememberedSignIn()
+    {
+        using var node = new MockNode();
+        var evict = false;
+        node.Serve(request =>
+        {
+            if (request.Operation == OperationRegister)
+            {
+                return Reply(OperationRegister, RegisterBody(session: 128));
+            }
+
+            if (evict)
+            {
+                evict = false;
+                return EvictionFrame(EvictionStaleClient);
+            }
+
+            return Reply(OperationNonReplicated, request.Code == 
GetClusterMetadataCode
+                ? ClusterMetadata(node.Port, node.Port, node.Port)
+                : []);
+        });
+
+        var configuration = new IggyClientConfigurator
+        {
+            BaseAddress = $"127.0.0.1:{node.Port}",
+            Protocol = Protocol.Tcp,
+            ReconnectionSettings = new ReconnectionSettings
+            {
+                Enabled = true,
+                MaxRetries = 2,
+                InitialDelay = TimeSpan.FromMilliseconds(20)
+            }
+        };
+        using var client = new TcpMessageStream(configuration, 
NullLoggerFactory.Instance);
+
+        await client.ConnectAsync(TestContext.Current.CancellationToken);
+        await client.LoginUserAsync("iggy", "iggy", 
TestContext.Current.CancellationToken);
+        await client.PingAsync(TestContext.Current.CancellationToken);
+        var registrationsBeforeEviction = node.Registrations;
+
+        evict = true;
+        await Assert.ThrowsAsync<VsrRequestOutcomeUnknownException>(() =>
+            client.CreateStreamAsync("evicted-mid-write", token: 
TestContext.Current.CancellationToken));
+
+        await client.PingAsync(TestContext.Current.CancellationToken);
+        Assert.True(node.Registrations > registrationsBeforeEviction,
+            "the reconnect signed in again with the remembered credentials");
+    }
+
+    /// <summary>
+    ///     A survivor that is not listening yet when its node dies still has 
to be found: the client keeps
+    ///     rotating over every endpoint it knows, so one that comes up while 
it is retrying is dialed on a
+    ///     later pass rather than only on the first.
+    /// </summary>
+    [Fact]
+    public async Task ResumesOnASurvivorThatComesUpWhileTheClientIsRetrying()
+    {
+        // A port nothing listens on: the survivor binds it only after the 
client has already failed on both
+        // endpoints, so the first pass cannot be the one that finds it.
+        var probe = new TcpListener(IPAddress.Loopback, 0);
+        probe.Start();
+        var survivorPort = (ushort)((IPEndPoint)probe.LocalEndpoint).Port;
+        probe.Stop();
+
+        using var primary = new MockNode();
+        primary.Serve(request => request.Code == GetClusterMetadataCode
+            ? Reply(OperationNonReplicated, ClusterMetadata(primary.Port, 
survivorPort, primary.Port))
+            : Answer(request));
+
+        var configuration = new IggyClientConfigurator
+        {
+            BaseAddress = $"127.0.0.1:{primary.Port}",
+            Protocol = Protocol.Tcp,
+            AutoLoginSettings = new AutoLoginSettings { Enabled = true, 
Username = "iggy", Password = "iggy" },
+            ReconnectionSettings = new ReconnectionSettings
+            {
+                Enabled = true,
+                MaxRetries = 4,

Review Comment:
   with 4 the second sweep starts after the 400ms backoff, when the survivor 
has been up since 300ms, so moving the budget check back above the advance 
stays green. `MaxRetries = 1` is the value that pins it. the docstring also 
lost the 'counts rotations, not dials' claim, which was the point of the test.



##########
core/common/src/types/configuration/tcp_config/tcp_client_reconnection_config.rs:
##########
@@ -20,10 +20,23 @@ 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, or `None` for unlimited.
+    ///
+    /// Passes, not dials: one pass tries the endpoint the client is on, the
+    /// addresses it was configured with, and every node the roster named, so a
+    /// survivor is reached inside the first pass rather than one delay per
+    /// endpoint.
     pub max_retries: Option<u32>,
-    /// Delay between connection attempts.
+    /// Delay between passes. The first pass runs at once when the client knows
+    /// more than one endpoint.
     pub interval: NonZeroIggyDuration,
+    /// Cooldown before redialing the endpoint that was just lost. It is owed 
to

Review Comment:
   measured from `connected_at`, not from the loss (`tcp_client.rs:969-981`): a 
session older than the interval redials that endpoint with no wait. the old 
builder wording 'after a previously successful connection' was closer.



##########
foreign/node/src/client/client.connection.ts:
##########
@@ -108,6 +124,13 @@ export class IggyConnection extends EventEmitter {
   public connecting: boolean;
   /** Whether the connection is being intentionally closed */
   public ending: boolean;
+  /**
+   * Whether the socket is being replaced by a deliberate leader redirect
+   * rather than lost. The drop looks the same from the outside, but nothing a
+   * caller submitted is in doubt: work waiting to be sent belongs on the node
+   * the client moves to, not in an error.
+   */
+  public redirecting: boolean;
   /** Reconnection configuration */
   private reconnectOption: ReconnectOption;
   /** Number of reconnection attempts made */

Review Comment:
   counts passes now, not attempts.



##########
foreign/go/client/tcp/tcp_failover_test.go:
##########
@@ -0,0 +1,614 @@
+// 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.
+
+package tcp
+
+import (
+       "context"
+       "crypto/tls"
+       "log/slog"
+       "net"
+       "sync/atomic"
+       "testing"
+       "time"
+
+       ierror "github.com/apache/iggy/foreign/go/errors"
+       "github.com/apache/iggy/foreign/go/internal/command"
+       "github.com/apache/iggy/foreign/go/internal/vsr"
+       "github.com/stretchr/testify/assert"
+       "github.com/stretchr/testify/require"
+)
+
+// The node a client signed in on dies; its next request has to complete on a
+// survivor the roster named, under the identity a fresh sign-in binds there.
+// Mirrors `core/integration/tests/cluster/failover_client_continuity.rs`.
+func TestFailover_ResumesOnASurvivorAfterTheSignedInNodeDies(t *testing.T) {
+       var survivor *testListener
+       var primary *testListener
+       var primaryDead atomic.Bool
+
+       survivor = listenVSR(t, nil, func(_, _ int, read request) []byte {
+               switch {
+               case read.code() == uint32(command.GetClusterMetadataCode):
+                       return clusterMetadataFrame(t, 1, primary.address(), 
survivor.address())
+               case read.operation() == vsr.OperationRegister:
+                       return registerReplyFrame(7, 512)
+               default:
+                       return replyFrame(vsr.OperationNonReplicated, nil)
+               }
+       })
+
+       primary = listenVSR(t, nil, func(_, _ int, read request) []byte {
+               // A dead node answers nothing; returning nil drops the 
connection the
+               // way a killed process does.
+               if primaryDead.Load() {
+                       return nil
+               }
+               switch {
+               case read.code() == uint32(command.GetClusterMetadataCode):
+                       // The primary leads, so the sign-in settles here and 
the roster is
+                       // only remembered -- not acted on -- until the node 
dies.
+                       return clusterMetadataFrame(t, 0, primary.address(), 
survivor.address())
+               case read.operation() == vsr.OperationRegister:
+                       return registerReplyFrame(7, 128)
+               default:
+                       return replyFrame(vsr.OperationNonReplicated, nil)
+               }
+       })
+
+       // No auto-login: the credentials come from the caller's own sign-in, 
which
+       // is the shape that could not reconnect at all before.
+       client := newDialingClient(t, primary.address())
+       ctx := context.Background()
+       require.NoError(t, client.Connect(ctx))
+       _, err := client.LoginUser(ctx, "iggy", "iggy")
+       require.NoError(t, err)
+       require.NoError(t, client.Ping(ctx), "the live primary answers")
+       require.Equal(t, primary.address(), client.currentServerAddress)
+
+       primaryDead.Store(true)
+       require.NoError(t, primary.listener.Close(), "stop accepting, so a 
redial is refused")
+
+       require.NoError(t, client.Ping(ctx),
+               "the client has to resume on the survivor the roster named")
+
+       assert.Equal(t, survivor.address(), client.currentServerAddress,
+               "the client moved off the dead endpoint")
+       assert.True(t, client.session.Bound(), "the session was re-established")
+
+       var registers int
+       for _, read := range survivor.recorded() {
+               if read.operation() == vsr.OperationRegister {
+                       registers++
+               }
+       }
+       assert.Equal(t, 1, registers,
+               "the remembered credentials signed in again on the survivor")
+}
+
+// Without any credentials there is nothing to sign in with, so a request on a
+// dead node fails instead of reconnecting into an unauthenticated session.
+func TestFailover_FailsFastWhenNothingEverSignedIn(t *testing.T) {
+       var server *testListener
+       var dead atomic.Bool
+       server = listenVSR(t, nil, func(_, _ int, read request) []byte {
+               if dead.Load() {
+                       return nil
+               }
+               return singleNodeHandler(t, func() string { return 
server.address() })(0, 0, read)
+       })
+
+       client := newDialingClient(t, server.address())
+       ctx := context.Background()
+       require.NoError(t, client.Connect(ctx))
+       require.NoError(t, client.Ping(ctx))
+
+       dead.Store(true)
+       require.NoError(t, server.listener.Close())
+
+       assert.Error(t, client.Ping(ctx),
+               "a client that never signed in cannot restore a session by 
reconnecting")
+}
+
+// A stale-client eviction is not caller intent: the heartbeat verifier sends 
it
+// after a gc pause or a laptop sleep, and a client that signed in by hand has
+// to recover from it exactly like one with a configured auto-login. Same rule
+// in every SDK.
+func TestFailover_ServerEvictionReplaysTheRememberedSignIn(t *testing.T) {
+       var server *testListener
+       var evict atomic.Bool
+       var registers atomic.Int32
+       server = listenVSR(t, nil, func(_, _ int, read request) []byte {
+               if read.operation() == vsr.OperationRegister {
+                       registers.Add(1)
+                       return registerReplyFrame(7, 128)
+               }
+               if evict.CompareAndSwap(true, false) {
+                       return evictionFrame(vsr.EvictionStaleClient, 0, 0)
+               }
+               if read.code() == uint32(command.GetClusterMetadataCode) {
+                       return clusterMetadataFrame(t, 0, server.address())
+               }
+               return replyFrame(vsr.OperationNonReplicated, nil)
+       })
+
+       client := newDialingClient(t, server.address())
+       ctx := context.Background()
+       require.NoError(t, client.Connect(ctx))
+       _, err := client.LoginUser(ctx, "iggy", "iggy")
+       require.NoError(t, err)
+       registersBefore := registers.Load()
+
+       evict.Store(true)
+       // A ping is non-replicated, so the eviction is absorbed: the reconnect 
it
+       // triggers signs in again with the credentials the sign-in remembered 
and
+       // the request completes over the session it re-established.
+       require.NoError(t, client.Ping(ctx), "the evicted request was not 
recovered")
+
+       _, remembered := client.signInCredentials()
+       assert.True(t, remembered, "an eviction is not a sign-out; the 
credentials stay")
+       require.NoError(t, client.Ping(ctx), "the session came back on its own")
+       assert.Greater(t, registers.Load(), registersBefore,
+               "the reconnect re-established the session")
+       assert.True(t, client.session.Bound())
+}
+
+// An explicit sign-out is caller intent: the reconnect must not sign back in
+// with the credentials the earlier sign-in used.
+func TestFailover_DoesNotResurrectASignedOutSession(t *testing.T) {
+       var server *testListener
+       var dropSocket atomic.Bool
+       server = listenVSR(t, nil, func(_, _ int, read request) []byte {
+               // A dropped connection is what makes the client reconnect at 
all; nil
+               // ends it the way a killed process does.
+               if dropSocket.Load() {
+                       return nil
+               }
+               return singleNodeHandler(t, func() string { return 
server.address() })(0, 0, read)
+       })
+
+       client := newDialingClient(t, server.address())
+       ctx := context.Background()
+       require.NoError(t, client.Connect(ctx))
+       _, err := client.LoginUser(ctx, "iggy", "iggy")
+       require.NoError(t, err)
+       require.NoError(t, client.LogoutUser(ctx))
+
+       credentials, ok := client.signInCredentials()
+       require.False(t, ok, "the sign-out forgot them")
+       require.Empty(t, credentials.username)
+
+       // The socket dies under a signed-out client: the reconnect has nothing 
to
+       // restore and must not invent a session.
+       dropSocket.Store(true)
+       assert.Error(t, client.Ping(ctx), "a signed-out client cannot replay 
through a sign-in")
+       dropSocket.Store(false)
+
+       require.NoError(t, client.Connect(ctx))
+       require.NoError(t, client.Ping(ctx), "the transport recovers on its 
own")
+
+       var registers int
+       for _, read := range server.recorded() {
+               if read.operation() == vsr.OperationRegister {
+                       registers++
+               }
+       }
+       assert.Equal(t, 1, registers,
+               "only the caller's own sign-in registered; the reconnect added 
none")
+}
+
+// A re-login over a dropped transport has to complete. The logout that ends
+// the old session runs while the sign-in lock is held, so a logout that enters
+// the reconnect path would reconnect, sign in with the remembered credentials,
+// and deadlock on that same lock.
+func TestFailover_ReLoginSurvivesALogoutTheTransportSwallowed(t *testing.T) {
+       var server *testListener
+       var dropLogout atomic.Bool
+       server = listenVSR(t, nil, func(_, _ int, read request) []byte {
+               if dropLogout.Load() && read.operation() == vsr.OperationLogout 
{
+                       // The frame is swallowed and the connection ends, 
exactly as a
+                       // node that dies mid-logout leaves it.
+                       return nil
+               }
+               return singleNodeHandler(t, func() string { return 
server.address() })(0, 0, read)
+       })
+
+       client := newDialingClient(t, server.address())
+       ctx := context.Background()
+       require.NoError(t, client.Connect(ctx))
+       _, err := client.LoginUser(ctx, "iggy", "iggy")
+       require.NoError(t, err)
+
+       dropLogout.Store(true)
+       relogin := make(chan error, 1)
+       go func() {
+               _, err := client.LoginUser(ctx, "iggy", "iggy")
+               relogin <- err
+       }()
+       select {
+       case err := <-relogin:
+               require.NoError(t, err, "the sign-in has to replay on the new 
connection")
+       case <-time.After(15 * time.Second):
+               t.Fatal("the re-login deadlocked on the sign-in lock")
+       }
+
+       assert.True(t, client.session.Bound(), "the replayed sign-in bound a 
session")
+       require.NoError(t, client.Ping(ctx))
+}
+
+// The other way a logout fails to land: the node answers it as not-admitted,
+// which is what a node that stopped being primary does. The redirect that
+// follows must not sign in on its own -- this goroutine holds the sign-in 
lock,
+// and the reconnect's automatic sign-in would wait on it forever.
+func TestFailover_ReLoginSurvivesALogoutTheOldPrimaryRefused(t *testing.T) {
+       var leader *testListener
+       var follower *testListener
+       var demoted atomic.Bool
+
+       // The node the client is on: leader until the logout, then a follower 
that
+       // refuses it as not-admitted and points at the survivor.
+       follower = listenVSR(t, nil, func(_, _ int, read request) []byte {
+               switch {
+               case read.code() == uint32(command.GetClusterMetadataCode):
+                       if demoted.Load() {
+                               return clusterMetadataFrame(t, 1, 
follower.address(), leader.address())
+                       }
+                       return clusterMetadataFrame(t, 0, follower.address(), 
leader.address())
+               case read.operation() == vsr.OperationRegister:
+                       return registerReplyFrame(7, 128)
+               case read.operation() == vsr.OperationLogout:
+                       demoted.Store(true)
+                       return statusReplyFrame(vsr.OperationLogout,
+                               uint32(ierror.TransientNotAcceptedCode), nil)
+               default:
+                       return replyFrame(vsr.OperationNonReplicated, nil)
+               }
+       })
+       leader = listenVSR(t, nil, func(_, _ int, read request) []byte {
+               switch {
+               case read.code() == uint32(command.GetClusterMetadataCode):
+                       return clusterMetadataFrame(t, 1, follower.address(), 
leader.address())
+               case read.operation() == vsr.OperationRegister:
+                       return registerReplyFrame(7, 256)
+               default:
+                       return replyFrame(vsr.OperationNonReplicated, nil)
+               }
+       })
+
+       client := newDialingClient(t, follower.address())
+       ctx := context.Background()
+       require.NoError(t, client.Connect(ctx))
+       _, err := client.LoginUser(ctx, "iggy", "iggy")
+       require.NoError(t, err)
+
+       relogin := make(chan error, 1)
+       go func() {
+               _, err := client.LoginUser(ctx, "iggy", "iggy")
+               relogin <- err
+       }()
+       select {
+       case err := <-relogin:
+               require.NoError(t, err, "the sign-in has to settle on the node 
that leads")
+       case <-time.After(15 * time.Second):
+               t.Fatal("the re-login deadlocked on the sign-in lock")
+       }
+
+       assert.True(t, client.session.Bound(), "the replayed sign-in bound a 
session")
+       assert.Equal(t, leader.address(), client.currentServerAddress)
+}
+
+// A logout that never landed still ended the session it belonged to, so the
+// credentials that established it must not outlive it: a sign-in that then
+// fails would otherwise leave them for the next dropped request to replay,
+// signing the old user back in after the caller asked for another one.
+func TestFailover_ARejectedReLoginDoesNotResurrectThePreviousUser(t 
*testing.T) {
+       var server *testListener
+       var dropLogout atomic.Bool
+       var dropSocket atomic.Bool
+       var rejectLogin atomic.Bool
+       var registeredUsers atomic.Int32
+       server = listenVSR(t, nil, func(_, _ int, read request) []byte {
+               if dropSocket.Load() {
+                       return nil
+               }
+               if dropLogout.Load() && read.operation() == vsr.OperationLogout 
{
+                       return nil
+               }
+               if read.operation() == vsr.OperationRegister {
+                       registeredUsers.Add(1)
+                       if rejectLogin.Load() {
+                               return statusReplyFrame(vsr.OperationRegister,
+                                       uint32(ierror.InvalidCredentialsCode), 
nil)
+                       }
+                       return registerReplyFrame(7, 128)
+               }
+               return singleNodeHandler(t, func() string { return 
server.address() })(0, 0, read)
+       })
+
+       client := newDialingClient(t, server.address())
+       ctx := context.Background()
+       require.NoError(t, client.Connect(ctx))
+       _, err := client.LoginUser(ctx, "alice", "alice")
+       require.NoError(t, err)
+
+       // The logout is swallowed and the sign-in that follows is rejected, so 
the
+       // client ends up with no session and no credentials it may use.
+       dropLogout.Store(true)
+       rejectLogin.Store(true)
+       _, err = client.LoginUser(ctx, "bob", "bob")
+       require.Error(t, err)
+
+       _, remembered := client.signInCredentials()
+       assert.False(t, remembered, "the ended session's credentials must not 
survive it")
+
+       // The socket dies with nothing remembered: the reconnect has no 
session to
+       // restore, and must not invent one out of the user who was signed in
+       // before.
+       dropLogout.Store(false)
+       dropSocket.Store(true)
+       registersBefore := registeredUsers.Load()
+       assert.Error(t, client.Ping(ctx), "there is no session left to restore")
+       assert.Equal(t, registersBefore, registeredUsers.Load(),
+               "the reconnect signed the previous user back in")
+}
+
+// reestablishAfter is a cooldown on redialing the endpoint that was lost. It
+// is owed to that endpoint alone, so a failover to another one must not sit
+// through it.
+func TestFailover_DoesNotSpendTheLostEndpointsPauseOnAnotherEndpoint(t 
*testing.T) {
+       var survivor *testListener
+       survivor = listenVSR(t, nil, singleNodeHandler(t, func() string { 
return survivor.address() }))
+
+       client := newDialingClient(t, deadAddress(t))
+       client.config.reconnection.reestablishAfter = time.Minute
+       client.knownServerAddresses = []string{survivor.address()}
+       client.connectedAt = time.Now()
+
+       ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
+       defer cancel()
+       started := time.Now()
+       require.NoError(t, client.Connect(ctx))
+
+       assert.Equal(t, survivor.address(), client.currentServerAddress)
+       assert.Less(t, time.Since(started), 2*time.Second,
+               "the failover waited out a pause it owed only the lost 
endpoint")
+}
+
+// The other half of the same promise: WithReestablishAfter is a cooldown on
+// the endpoint that was lost, and a known roster does not cancel it.
+func TestFailover_KeepsTheReestablishPauseForTheEndpointThatWasLost(t 
*testing.T) {
+       var current *testListener
+       current = listenVSR(t, nil, singleNodeHandler(t, func() string { return 
current.address() }))
+
+       client := newDialingClient(t, current.address())
+       client.config.reconnection.reestablishAfter = 500 * time.Millisecond
+       client.knownServerAddresses = []string{deadAddress(t)}
+       client.connectedAt = time.Now()
+
+       started := time.Now()
+       require.NoError(t, client.Connect(context.Background()))
+
+       assert.Equal(t, current.address(), client.currentServerAddress)
+       assert.GreaterOrEqual(t, time.Since(started), 350*time.Millisecond,
+               "the cooldown on the endpoint that was lost was skipped")
+}
+
+// The cooldown is a pace limit, not a commitment: a caller that gave the
+// connect a deadline has to get an answer inside it, and Close has to end the
+// wait too.
+func TestFailover_TheReestablishPauseHonoursTheCallersDeadline(t *testing.T) {
+       current := listenVSR(t, nil, func(_, _ int, read request) []byte {
+               return singleNodeHandler(t, func() string { return 
"127.0.0.1:8090" })(0, 0, read)
+       })
+
+       client := newDialingClient(t, current.address())
+       client.config.reconnection.reestablishAfter = time.Minute
+       client.connectedAt = time.Now()
+
+       ctx, cancel := context.WithTimeout(context.Background(), 
300*time.Millisecond)
+       defer cancel()
+       started := time.Now()
+       _ = client.Connect(ctx)
+
+       assert.Less(t, time.Since(started), 5*time.Second,
+               "the cooldown outlived the deadline the caller gave the 
connect")
+}
+
+// A node whose syns are dropped must not hold the sweep: without a bound on
+// the dial the survivors behind it are never reached. A black-holed address
+// cannot be arranged portably, so this pins the bound itself.
+func TestFailover_BoundsTheDialWhenOtherEndpointsAreQueuedBehindIt(t 
*testing.T) {
+       assert.Equal(t, 2*time.Second, failoverDialTimeout,
+               "the dial bound has to match the other SDKs")
+
+       var survivor *testListener
+       survivor = listenVSR(t, nil, singleNodeHandler(t, func() string { 
return survivor.address() }))
+
+       // A listener that accepts and never answers: the dial completes out of 
the
+       // backlog, so only the bound ends the attempt.
+       silent, err := net.Listen("tcp", "127.0.0.1:0")
+       require.NoError(t, err)
+       t.Cleanup(func() { _ = silent.Close() })
+
+       client := newDialingClient(t, silent.Addr().String(),
+               WithTLS(WithTLSValidateCertificate(false)))
+       client.knownServerAddresses = []string{survivor.address()}
+
+       done := make(chan error, 1)
+       go func() { done <- client.Connect(context.Background()) }()
+       select {
+       case <-done:
+       case <-time.After(3 * failoverDialTimeout):
+               t.Fatal("the sweep never got past an endpoint that answers 
nothing")
+       }
+}
+
+// An endpoint that accepts TCP but fails the handshake is not where this
+// client lives: recording it would make the next pass lead with it and shadow
+// every endpoint behind it.
+func TestFailover_DoesNotSettleOnAnEndpointThatFailedTheHandshake(t 
*testing.T) {
+       hangup, err := net.Listen("tcp", "127.0.0.1:0")
+       require.NoError(t, err)
+       t.Cleanup(func() { _ = hangup.Close() })
+       go func() {
+               for {
+                       connection, err := hangup.Accept()
+                       if err != nil {
+                               return
+                       }
+                       // Plain TCP behind a TLS client: the dial succeeds, 
the handshake
+                       // cannot.
+                       _ = connection.Close()
+               }
+       }()
+
+       configured := deadAddress(t)
+       client := newDialingClient(t, configured, 
WithTLS(WithTLSValidateCertificate(false)))
+       client.config.reconnection.enabled = false
+       client.knownServerAddresses = []string{hangup.Addr().String()}
+
+       require.Error(t, client.Connect(context.Background()))
+       assert.Equal(t, configured, client.currentServerAddress,
+               "the endpoint that failed the handshake became the current one")
+}
+
+// The SNI of a dial belongs to the endpoint being dialed. Taken from the
+// endpoint the client just lost, a failover to a node the certificate does not
+// cover fails the handshake -- which is every failover, once the addresses
+// differ.
+func TestFailover_UsesTheDialedEndpointAsTheServerName(t *testing.T) {
+       certificate, caPath := selfSignedCert(t)
+       survivor := listenVSR(t,
+               func(conn net.Conn) net.Conn {
+                       return tls.Server(conn, &tls.Config{Certificates: 
[]tls.Certificate{certificate}})
+               },
+               singleNodeHandler(t, func() string { return "127.0.0.1:8090" }))
+
+       // The certificate covers 127.0.0.1, and the endpoint the client starts 
on
+       // is 127.0.0.2, where nothing listens: with the server name taken from
+       // that endpoint, the handshake on the survivor is checked against the
+       // address that died.
+       _, port, err := net.SplitHostPort(survivor.address())
+       require.NoError(t, err)
+       client := newDialingClient(t, "127.0.0.2:"+port,
+               WithTLS(WithTLSCAFile(caPath), 
WithTLSValidateCertificate(true)))
+       client.knownServerAddresses = []string{"127.0.0.1:" + port}
+
+       require.NoError(t, client.Connect(context.Background()))
+       assert.Equal(t, "127.0.0.1:"+port, client.currentServerAddress)
+}
+
+// A TLS configuration the client itself cannot satisfy says the same thing on
+// every attempt, so it has to reach the caller instead of being redialed every
+// interval forever -- which is what the default unlimited retries did with it.
+func TestFailover_AConfigFaultEndsTheConnectInsteadOfRetryingForever(t 
*testing.T) {
+       certificate, _ := selfSignedCert(t)
+       server := listenVSR(t,
+               func(conn net.Conn) net.Conn {
+                       return tls.Server(conn, &tls.Config{Certificates: 
[]tls.Certificate{certificate}})
+               },
+               singleNodeHandler(t, func() string { return "127.0.0.1:8090" }))
+
+       // A CA the server's certificate was not signed by: no retry makes that
+       // certificate acceptable.
+       _, unrelatedCA := selfSignedCert(t)
+       client := NewIggyTcpClient(slog.New(slog.DiscardHandler),
+               WithServerAddress(server.address()),
+               WithTLS(WithTLSCAFile(unrelatedCA), 
WithTLSValidateCertificate(true)))
+       t.Cleanup(func() { _ = client.Close() })
+       client.config.reconnection.maxRetries = 0 // unlimited
+       client.config.reconnection.interval = 10 * time.Millisecond
+
+       done := make(chan error, 1)
+       go func() { done <- client.Connect(context.Background()) }()
+       select {
+       case err := <-done:
+               require.Error(t, err)
+       case <-time.After(5 * time.Second):
+               t.Fatal("a connect that can never succeed has to end instead of 
retrying forever")
+       }
+}
+
+// A client with nothing to dial must say so: reporting success would leave
+// every request answering ErrNotConnected while Connect keeps claiming a
+// connection.
+func TestFailover_RejectsAConnectWithNoEndpointToDial(t *testing.T) {
+       client := NewIggyTcpClient(slog.New(slog.DiscardHandler), 
WithServerAddress(""))
+       t.Cleanup(func() { _ = client.Close() })
+
+       require.ErrorIs(t, client.Connect(context.Background()), 
ierror.ErrCannotEstablishConnection)
+       assert.Error(t, client.Ping(context.Background()))
+}
+
+// Concurrent Connects are one attempt, and a caller that did not run it still
+// gets a client it can use the moment its Connect returns. Told "connected"
+// while the attempt is still signing in, its next request fails
+// ErrNotConnected for no reason of its own.
+func TestConnect_ConcurrentCallersShareOneAttempt(t *testing.T) {

Review Comment:
   covers direct `Connect()` racing only. the case that hurts is two requests 
reconnecting through `exchange` at once, each calling `disconnect()` first - 
that path still dials twice, add it here.



##########
foreign/node/src/client/client.socket.test.ts:
##########
@@ -864,6 +992,224 @@ describe('VSR client socket', () => {
     }
   });
 
+  it('re-issues every request refused by a demoted node, not just the first',
+    async () => {
+      // One demotion, several refused requests: each of them re-checking on
+      // its own would move the client once per request, and the first
+      // redirect's drop would fail the others' roster reads - reporting a
+      // refusal they never had to.
+      const leader = await startVsrServer((frame, socket) => {
+        singleNodeHandler(leader.port)(frame, socket);
+      });
+      // Leader at login, so the settlement leaves the client here, then
+      // demoted: the refusals below are what tells the client to look again.
+      let demotedYet = false;
+      const demoted = await startVsrServer((frame, socket) => {
+        const code = frame.readUInt32LE(REQUEST_OFFSET.reserved);
+        if (code === COMMAND_CODE.GetClusterMetadata) {
+          socket.write(replyFrame(
+            Operation.NonReplicated,
+            demotedYet
+              ? twoNodeMetadataBody(demoted.port, leader.port)
+              : twoNodeMetadataBody(leader.port, demoted.port)
+          ));
+          return;
+        }
+        if (code === 60_032) {
+          socket.write(replyFrame(Operation.NonReplicated, Buffer.alloc(0), 
58));
+          return;
+        }
+        singleNodeHandler(demoted.port)(frame, socket);
+      });
+
+      const client = new CommandResponseStream(vsrConfig(demoted.port));
+      try {
+        await client.authenticate(vsrConfig(demoted.port).credentials);
+        demotedYet = true;
+        const rosterReadsBefore = demoted.frames.filter(
+          (frame) => frame.readUInt32LE(REQUEST_OFFSET.reserved) ===
+            COMMAND_CODE.GetClusterMetadata
+        ).length;
+
+        // Two refusals, one re-check: the second caller shares the move the
+        // first started rather than starting its own or being told 58.
+        const moves = await Promise.all([
+          followLeaderMove(client),
+          followLeaderMove(client)
+        ]);
+
+        assert.deepEqual(moves, [true, true]);
+        const rosterReads = demoted.frames.filter(
+          (frame) => frame.readUInt32LE(REQUEST_OFFSET.reserved) ===
+            COMMAND_CODE.GetClusterMetadata
+        ).length - rosterReadsBefore;
+        assert.equal(rosterReads, 1,
+          'each refusal re-read the roster on its own'
+        );
+        const connection = (client as unknown as {
+          connection: { isConnectedTo: (host: string, port: number) => boolean 
}
+        }).connection;
+        assert.equal(connection.isConnectedTo('127.0.0.1', leader.port), true);
+      } finally {
+        client.destroy();
+        await leader.close();
+        await demoted.close();
+      }
+    }
+  );
+
+  it('re-issues a command queued behind a leader move instead of failing it',
+    async () => {
+      // A move replaces the socket, which looks like a drop to everything
+      // waiting in the queue. Nothing queued was written, though, so it 
belongs
+      // on the node the client moves to rather than in a lost-connection error
+      // the caller can do nothing about.
+      const leader = await startVsrServer((frame, socket) => {
+        singleNodeHandler(leader.port)(frame, socket);
+      });
+      const demoted = await startVsrServer((frame, socket) => {
+        singleNodeHandler(demoted.port)(frame, socket);
+      });
+
+      const client = new CommandResponseStream(vsrConfig(demoted.port));
+      try {
+        await client.authenticate(vsrConfig(demoted.port).credentials);
+
+        // Parked the way a command is while something else holds the queue.
+        const queued = new Promise<CommandResponse>((resolve, reject) => {
+          execQueue(client).push({
+            command: 60_034,
+            payload: Buffer.alloc(0),
+            handleResponse: true,
+            deadline: Date.now() + 30_000,
+            resolve,
+            reject
+          });
+        });
+
+        await connectionOf(client).redirect('127.0.0.1', leader.port);
+        await queued;
+
+        const landedOnLeader = leader.frames.some(
+          (frame) => frame.readUInt32LE(REQUEST_OFFSET.reserved) === 60_034
+        );
+        assert.ok(landedOnLeader,
+          'the queued command never reached the node the client moved to'
+        );
+      } finally {
+        client.destroy();
+        await leader.close();
+        await demoted.close();
+      }
+    }
+  );
+
+  it('surfaces the refusal rather than a timeout when the budget runs out',
+    async () => {
+      const server = await startVsrServer((frame, socket) => {
+        if (frame.readUInt32LE(REQUEST_OFFSET.reserved) === 60_036) {
+          socket.write(replyFrame(Operation.NonReplicated, Buffer.alloc(0), 
58));
+          return;
+        }
+        singleNodeHandler(server.port)(frame, socket);
+      });
+      const client = new CommandResponseStream(vsrConfig(server.port));
+      const realNow = Date.now;
+      try {
+        await client.authenticate(vsrConfig(server.port).credentials);
+        // The request's own budget, then the clock jumped past it: what is 
left
+        // cannot carry another attempt, so the caller has to see the answer 
the
+        // server gave and not the timeout a doomed re-issue would produce.
+        // The request's budget, the first exchange, the window that hands the
+        // refusal out, and then a clock 10ms short of the deadline: too little
+        // to carry another attempt, so the caller has to see the answer the
+        // server gave rather than the timeout a doomed re-issue would produce.
+        const times = [0, 1, 2_001, 29_990];

Review Comment:
   passes on the old code too: at 29_990 the old check let the roster read 
through and then hit the `remaining <= 0` path, also a typed 58. what changed 
is a positive remainder under 50ms after the roster read - clock it there. the 
two comment paragraphs above also say the same thing.



##########
foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java:
##########
@@ -452,15 +495,50 @@ private AsyncTcpConnection openConnection(ConnectionInfo 
target) {
                 enableTls,
                 tlsCertificate,
                 poolConfig,
-                connectionTimeout,
+                dialTimeout(),
                 requestTimeout,
                 heartbeatInterval,
                 maxVsrFrameSize,
                 this::retryTransientOnLeader,
-                routingState::clearAssignments,
+                this::onSessionReset,
                 this::onConnectionFailure);
     }
 
+    /**
+     * How long one dial may take.
+     *
+     * With other endpoints queued behind this one, a node whose syns are
+     * dropped must not hold the rotation, so the wait is capped at
+     * {@link #FAILOVER_DIAL_TIMEOUT} - the bound the other SDKs use. A caller
+     * who configured a connection timeout gets exactly that, and a client that
+     * knows one endpoint keeps the ordinary default.
+     */
+    private Optional<Duration> dialTimeout() {
+        if (connectionTimeout.isPresent() || redialCandidates().size() < 2) {
+            return connectionTimeout;

Review Comment:
   a configured `connectionTimeout` skips the cap, while rust, go, c# and node 
cap every dial once more than one endpoint is known, whatever the configured 
timeout. it is also tcp connect only - the tls handler has no handshake 
timeout, so netty's 10s applies. and nothing pins it: the killed `ServerSocket` 
refuses instantly, so passing `connectionTimeout` here again stays green.



##########
foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java:
##########
@@ -581,46 +661,147 @@ private CompletableFuture<Void> redialAttempt(int 
attempt, RetryPolicy policy) {
         if (closed) {
             return CompletableFuture.completedFuture(null);
         }
-        if (attempt > policy.getMaxRetries()) {
+        List<ConnectionInfo> candidates = redialCandidates();
+        // The retry budget bounds the rotations, not the endpoints. A policy 
of
+        // zero retries with several endpoints known still gets one rotation:
+        // those endpoints - the address the client was configured with, the
+        // nodes the roster named - were made known in order to be tried, and
+        // the other SDKs sweep them once too. With one endpoint known, zero
+        // retries redials nothing, which is what it asked for.
+        boolean sweepOnce = attempt == 1 && candidates.size() > 1;
+        if (attempt > policy.getMaxRetries() && !sweepOnce) {

Review Comment:
   no test drives this with `noRetry()` and two candidates. the log lines for 
that pass also read 'attempt 1/0' and 'gave up after 0 attempts'.



##########
foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java:
##########
@@ -581,46 +661,147 @@ private CompletableFuture<Void> redialAttempt(int 
attempt, RetryPolicy policy) {
         if (closed) {
             return CompletableFuture.completedFuture(null);
         }
-        if (attempt > policy.getMaxRetries()) {
+        List<ConnectionInfo> candidates = redialCandidates();
+        // The retry budget bounds the rotations, not the endpoints. A policy 
of
+        // zero retries with several endpoints known still gets one rotation:
+        // those endpoints - the address the client was configured with, the
+        // nodes the roster named - were made known in order to be tried, and
+        // the other SDKs sweep them once too. With one endpoint known, zero
+        // retries redials nothing, which is what it asked for.
+        boolean sweepOnce = attempt == 1 && candidates.size() > 1;
+        if (attempt > policy.getMaxRetries() && !sweepOnce) {
             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);
+        // 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 (!isSignInRejection(unwrap(loginError))) {
                         log.warn(
-                                "Redial attempt {} to {} failed: {}",
-                                attempt,
+                                "The sign-in on {} did not complete: {}",
                                 target.serverAddress(),
-                                error.getMessage());
-                        return redialAttempt(attempt + 1, policy);
-                    })
-                    .thenCompose(Function.identity());
-        });
+                                loginError.getMessage());
+                        return sweepCandidates(candidates, index + 1, attempt, 
policy);
+                    }
+                    log.error(
+                            "Reconnected to {} but the sign-in was rejected: 
{}. The connection stands"
+                                    + " unauthenticated until the caller signs 
in again.",
+                            target.serverAddress(),
+                            loginError.getMessage());
+                    rememberedLogin = null;
+                    return CompletableFuture.<Void>completedFuture(null);
+                })
+                .thenCompose(Function.identity());
+    }
+
+    /**
+     * Whether the server answered the sign-in with a verdict no other endpoint
+     * would change: a rotated password, an expired token.
+     *
+     * Only that ends a redial. Everything else -- the channel closing before
+     * the reply, a timeout, a transient refusal from a node that is not the
+     * primary -- says nothing about the credentials, and treating it as a
+     * rejection would drop them and leave the client published on a node that
+     * is already gone, with every later call failing "not authenticated".
+     */
+    static boolean isSignInRejection(Throwable error) {
+        if (!(error instanceof IggyServerException serverError)) {
+            return false;
+        }
+        int code = serverError.getRawErrorCode();
+        return code != AsyncTcpConnection.TRANSIENT_NOT_ACCEPTED && code != 
AsyncTcpConnection.TRANSIENT_NOT_COMMITTED;
+    }
+
+    private static Throwable unwrap(Throwable error) {
+        return error instanceof CompletionException && error.getCause() != 
null ? error.getCause() : error;
     }
 
     /**
-     * Replays the builder credentials on the freshly published connection.
+     * Re-establishes the session on the freshly published connection with the
+     * sign-in that last succeeded, falling back to the credentials configured
+     * on the builder when no login has run yet.
+     *
+     * The last sign-in rather than the configured one, which is where this
+     * differs from the Rust and Go SDKs: the connection re-authenticates a
+     * replacement channel from the login payload it captured, which is that
+     * same last sign-in. Replaying a different user here would make the same
+     * eviction land on a different session depending on whether the channel
+     * or the redial got there first. A client that only ever used its
+     * configured credentials remembers exactly those, so nothing changes for
+     * it.
+     *
      * The login runs through the users client, so leader discovery retargets
      * again before Register when the redialed node is not the leader.
      */
     private CompletableFuture<Void> replayLogin() {

Review Comment:
   this is now the opposite of rust and go, and the new test pins it. pick one 
rule for every sdk: either configured credentials win everywhere (then the 
channel re-auth here has to use them too) or the last sign-in wins everywhere. 
as it stands the same failover replays different users depending on the sdk. 
the eviction case with both a configured and a hand-run login has no test at 
all.



##########
core/common/src/types/configuration/tcp_config/tcp_client_reconnection_config.rs:
##########
@@ -20,10 +20,23 @@ 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, or `None` for unlimited.

Review Comment:
   `Some(0)` still runs one full pass, so this is the number of passes after 
the first.



##########
core/sdk/src/tcp/tcp_client.rs:
##########
@@ -231,6 +357,60 @@ impl iggy_common::VsrSessionControl for TcpClient {
         Ok(())
     }
 
+    async fn remember_session_credentials(&self, credentials: Credentials, 
user_id: u32) {
+        self.session_credentials
+            .lock()
+            .await
+            .replace(RememberedSignIn {
+                credentials,
+                user_id,
+            });
+    }
+
+    async fn forget_session_credentials(&self) {
+        self.session_credentials.lock().await.take();
+    }
+
+    async fn refresh_session_password(&self, user: &Identifier, new_password: 
&str) {
+        let mut remembered = self.session_credentials.lock().await;
+        let Some(sign_in) = remembered.as_mut() else {
+            return;
+        };
+        // A personal access token is not derived from the password.
+        let Credentials::UsernamePassword(username, password) = &mut 
sign_in.credentials else {
+            return;
+        };
+
+        let targets_session_user = match user.kind {
+            IdKind::Numeric => user.get_u32_value().is_ok_and(|id| id == 
sign_in.user_id),
+            IdKind::String => user
+                .get_cow_str_value()
+                .is_ok_and(|name| name.as_ref() == username),
+        };
+        if !targets_session_user {
+            return;
+        }
+
+        *password = SecretString::from(new_password.to_owned());
+        // The configured credentials cannot be rewritten -- the config is
+        // shared and immutable -- and the password they carry will never work
+        // again, so the new one is kept here instead. Kept outside the
+        // remembered sign-in on purpose: that record is replaced wholesale by
+        // every later login, so a marker on it would survive exactly one
+        // reconnect and the one after that would replay the dead password.
+        let configured_user_changed = matches!(
+            &self.config.auto_login,
+            AutoLogin::Enabled(Credentials::UsernamePassword(configured, _))

Review Comment:
   the override only lands when the remembered sign-in is the configured user. 
sign in by hand as another user (or with a pat), then change the configured 
user's password, and every later reconnect still auto-logins with the old one.



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

To unsubscribe, e-mail: [email protected]

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

Reply via email to