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]
