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


##########
foreign/go/client/tcp/tcp_failover_test.go:
##########
@@ -0,0 +1,190 @@
+// 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"
+       "log/slog"
+       "sync/atomic"
+       "testing"
+       "time"
+
+       "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 the server ending the session authoritatively,
+// like a logout: the remembered sign-in must not resurrect it, so the evicted
+// request surfaces the loss instead of reconnecting into a fresh session.
+func TestFailover_ServerEvictionForgetsTheRememberedSignIn(t *testing.T) {
+       var server *testListener
+       var evict atomic.Bool
+       server = listenVSR(t, nil, func(_, _ int, read request) []byte {
+               if read.operation() == vsr.OperationRegister {
+                       return registerReplyFrame(7, 128)
+               }
+               if evict.Load() {
+                       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)
+       connectionsBefore := server.connections()
+
+       evict.Store(true)
+       require.Error(t, client.Ping(ctx), "the evicted request surfaces the 
loss")
+
+       _, remembered := client.signInCredentials()
+       assert.False(t, remembered, "the eviction forgot the remembered 
sign-in")
+       assert.Equal(t, connectionsBefore, server.connections(),
+               "no reconnect dial resurrected the evicted session")
+}
+
+// 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
+       server = listenVSR(t, nil, singleNodeHandler(t, func() string { return 
server.address() }))
+
+       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()
+       assert.False(t, ok, "the sign-out forgot them")
+       assert.Empty(t, credentials.username)
+}
+
+func TestFailover_LeavesTheReestablishPauseToSingleEndpointClients(t 
*testing.T) {

Review Comment:
   Replaced with two tests that drive `Connect` and measure the elapsed time:
   
   - `TestFailover_DoesNotSpendTheLostEndpointsPauseOnAnotherEndpoint`: dead 
current endpoint, live roster endpoint, a one-minute `reestablishAfter` — 
connects in under 2s under a 5s context.
   - `TestFailover_KeepsTheReestablishPauseForTheEndpointThatWasLost`: live 
current endpoint, dead roster endpoint, a 500ms window — takes at least 350ms.
   
   Go also picked up the Rust semantics from the `reestablish_after` thread: 
instead of skipping the pause whenever there is more than one candidate, the 
paced endpoint rotates to the end of the list and waits out what is left of its 
window there. Deleting that rotation turns the second test red.



##########
foreign/go/client/tcp/tcp_failover_test.go:
##########
@@ -0,0 +1,190 @@
+// 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"
+       "log/slog"
+       "sync/atomic"
+       "testing"
+       "time"
+
+       "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 the server ending the session authoritatively,
+// like a logout: the remembered sign-in must not resurrect it, so the evicted
+// request surfaces the loss instead of reconnecting into a fresh session.
+func TestFailover_ServerEvictionForgetsTheRememberedSignIn(t *testing.T) {
+       var server *testListener
+       var evict atomic.Bool
+       server = listenVSR(t, nil, func(_, _ int, read request) []byte {
+               if read.operation() == vsr.OperationRegister {
+                       return registerReplyFrame(7, 128)
+               }
+               if evict.Load() {
+                       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)
+       connectionsBefore := server.connections()
+
+       evict.Store(true)
+       require.Error(t, client.Ping(ctx), "the evicted request surfaces the 
loss")
+
+       _, remembered := client.signInCredentials()
+       assert.False(t, remembered, "the eviction forgot the remembered 
sign-in")
+       assert.Equal(t, connectionsBefore, server.connections(),
+               "no reconnect dial resurrected the evicted session")
+}
+
+// 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) {

Review Comment:
   Rewritten. The handler drops the socket after the logout, a `Ping` fails 
(pinning that a signed-out client does not replay through a sign-in), then 
`Connect` + `Ping` recover the transport, and the test asserts exactly one 
`OperationRegister` reached the server for the whole run.



##########
foreign/go/client/tcp/tcp_session_management.go:
##########
@@ -33,15 +33,25 @@ func (c *IggyTcpClient) LoginUser(ctx context.Context, 
username string, password
        if err != nil {
                return nil, err
        }
-       return c.register(ctx, uint32(command.LoginRegisterCode), body)
+       identity, err := c.register(ctx, uint32(command.LoginRegisterCode), 
body)
+       if err != nil {
+               return nil, err
+       }
+       c.rememberLogin(NewUsernamePasswordCredentials(username, password))

Review Comment:
   Fixed. `LoginUser` and `LoginWithPersonalAccessToken` now pass the 
credentials into `register`, which remembers them after `settleOnLeader` while 
it still holds `registerMtx`. `rememberLogin`'s doc says so, so the next caller 
does not move it back out.



##########
foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java:
##########
@@ -585,7 +616,7 @@ 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);
+        ConnectionInfo target = ReconnectPlan.target(redialCandidates(), 
attempt);

Review Comment:
   Fixed. `redialAttempt` now plans a rotation and `sweepCandidates` walks 
every candidate inside it; the policy delay is spent once per rotation, and the 
first rotation runs immediately when more than one endpoint is known (the node 
just lost may be gone for good, and the other SDKs do the same). 
`ReconnectPlan.target` is gone, along with its tests, since a rotation no 
longer picks a single endpoint.
   
   `shouldResumeOnASurvivorAfterTheSignedInNodeDies` now uses a 5s fixed delay 
with a 4s resume budget, so it can only pass if the survivor is reached inside 
the first rotation. Reverting the sweep makes it fail.



##########
foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java:
##########
@@ -585,7 +616,7 @@ 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);
+        ConnectionInfo target = ReconnectPlan.target(redialCandidates(), 
attempt);
         Duration delay = ReconnectPlan.delay(policy, attempt);

Review Comment:
   Fixed. `redialCandidates` dedups with a new 
`LeaderAwareness.isSameSpelling`, which compares canonicalized host and port 
and never resolves. `isSameAddress` keeps the resolving comparison for the 
leader check, which is one comparison off the dial path, and its javadoc now 
says why the cheap half exists. Covered by `LeaderAwarenessTest.SameAddress`.



##########
foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java:
##########
@@ -585,7 +616,7 @@ 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);
+        ConnectionInfo target = ReconnectPlan.target(redialCandidates(), 
attempt);
         Duration delay = ReconnectPlan.delay(policy, attempt);
         Executor delayedExecutor = 
CompletableFuture.delayedExecutor(delay.toMillis(), TimeUnit.MILLISECONDS);

Review Comment:
   Split into the three outcomes. A dial failure moves to the next candidate. A 
login failure that is a connection loss also moves on (the endpoint died 
mid-sign-in). Anything else — a rejected password, an expired token — ends the 
redial with the connection left standing, logs at error, and drops 
`rememberedLogin` so nothing replays the rejected credentials. No more tearing 
down a freshly published good connection twelve times to arrive 
connected-but-unauthenticated.



-- 
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