numinnex commented on code in PR #3944:
URL: https://github.com/apache/iggy/pull/3944#discussion_r3852355035
##########
foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java:
##########
@@ -129,10 +132,38 @@ public class AsyncIggyTcpClient {
private final Optional<File> tlsCertificate;
private final TcpConnectionPoolConfig poolConfig;
private final ClientRoutingState routingState = new ClientRoutingState();
+ private final LoginRoutingHook loginRoutingHook = new LoginRoutingHook() {
+
+ @Override
+ public CompletableFuture<IdentityInfo>
loginOnLeader(Supplier<CompletableFuture<IdentityInfo>> loginAttempt) {
+ return AsyncIggyTcpClient.this.loginOnLeader(loginAttempt);
+ }
+
+ @Override
+ public void forgetLogin() {
+ rememberedLogin = null;
Review Comment:
Fixed: `close()` nulls it, with a comment naming the `login(A) -> close() ->
connect()` path you described.
##########
foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java:
##########
@@ -638,6 +675,9 @@ CompletableFuture<IdentityInfo>
loginOnLeader(Supplier<CompletableFuture<Identit
CompletableFuture<IdentityInfo> callerFuture = new
CompletableFuture<>();
transaction.whenComplete((identity, error) -> {
gate.complete(null);
+ if (error == null) {
+ rememberedLogin = loginAttempt;
Review Comment:
Plumbed. `VsrResponseHandler` takes an `IntConsumer` and reports the
eviction's error code, `AsyncTcpConnection.onSessionEvicted(int)` forwards it,
and the client's `onSessionReset(int)` drops both the remembered sign-in and
the connection's captured payload through a new
`AsyncTcpConnection.forgetCapturedLogin()`.
One scoping decision worth your eyes: the drop only applies to clients with
no configured credentials. `shouldReplayTransientImplicitLoginAfterEviction`
asserts that a client built with credentials re-authenticates after a
stale-client eviction, which is the auto-login recovery you describe in the
Rust thread, and dropping the payload unconditionally broke it. So configured
credentials still recover; a sign-in the caller ran by hand is not revived.
New test `shouldNotReviveAHandRunSignInAfterAStaleClientEviction`, red
without the fix.
##########
foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java:
##########
@@ -617,6 +648,12 @@ private CompletableFuture<Void> redialAttempt(int attempt,
RetryPolicy policy) {
* again before Register when the redialed node is not the leader.
*/
private CompletableFuture<Void> replayLogin() {
+ Supplier<CompletableFuture<IdentityInfo>> replay = rememberedLogin;
Review Comment:
Flipped: configured credentials first, the remembered sign-in second, with
the reasoning in the javadoc — the remembered sign-in exists to make a hand-run
login as reconnectable as a configured one, not to override it.
##########
foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClient.java:
##########
@@ -617,6 +648,12 @@ private CompletableFuture<Void> redialAttempt(int attempt,
RetryPolicy policy) {
* again before Register when the redialed node is not the leader.
Review Comment:
Rewritten to match the code: a rotation sweeps every endpoint the client
knows, only a rotation that reaches none of them waits out the policy delay,
and the sign-in is replayed whether it was configured on the builder or run by
the caller, personal access tokens included.
##########
foreign/java/java-sdk/src/main/java/org/apache/iggy/client/async/tcp/LeaderAwareness.java:
##########
@@ -248,6 +265,24 @@ private static boolean
reachesOnlyLocalMachine(InetAddress[] addresses) {
/**
Review Comment:
Fixed: the javadoc moved back onto `LeaderCheck`, and `LeaderLookup` keeps
its own.
##########
foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/LeaderAwarenessTest.java:
##########
@@ -276,6 +281,22 @@ void shouldGiveUpWhenMetadataFetchThrowsSynchronously() {
assertThat(leader).isEmpty();
}
+ @Test
+ void shouldRememberEveryNodeTheRosterNamesEvenWhileLeaderless() {
Review Comment:
Added a `NodeTargets` group: the port-0 skip, unhealthy nodes kept (a node
that is down is still where the cluster says it lives), and
`LeaderLookup.inconclusive()` naming no endpoint.
No accessor was needed for the second half: the client keeps its last roster
precisely because an inconclusive lookup returns no endpoints and
`findLeaderElsewhere` only assigns `rosterTargets` when the list is non-empty,
so the invariant is testable on the record itself.
##########
foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClientEndpointFailoverTest.java:
##########
@@ -0,0 +1,318 @@
+/*
+ * 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 org.apache.iggy.client.async.tcp;
+
+import io.netty.buffer.ByteBuf;
+import io.netty.buffer.Unpooled;
+import org.apache.iggy.config.RetryPolicy;
+import org.junit.jupiter.api.Test;
+
+import java.io.EOFException;
+import java.io.IOException;
+import java.io.InputStream;
+import java.io.OutputStream;
+import java.net.InetAddress;
+import java.net.ServerSocket;
+import java.net.Socket;
+import java.nio.ByteBuffer;
+import java.nio.ByteOrder;
+import java.nio.charset.StandardCharsets;
+import java.time.Duration;
+import java.util.List;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/**
+ * 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
+ * {@code core/integration/tests/cluster/failover_client_continuity.rs}. The
+ * mock VSR framing matches {@link AsyncIggyTcpClientTransientFailoverTest},
+ * kept separate so a death mid-connection cannot disturb that suite's server.
+ */
+class AsyncIggyTcpClientEndpointFailoverTest {
+ private static final int HEADER_SIZE = 256;
+ private static final int SIZE_OFFSET = 48;
+ private static final int COMMAND_OFFSET = 60;
+ private static final int REQUEST_ID_OFFSET = 168;
+ private static final int REQUEST_OPERATION_OFFSET = 176;
+ private static final int REQUEST_CODE_OFFSET = 196;
+ private static final int REPLY_REQUEST_ID_OFFSET = 200;
+ private static final int REPLY_OPERATION_OFFSET = 208;
+ private static final int REPLY_STATUS_OFFSET = 216;
+
+ private static final int COMMAND_REPLY = 8;
+ private static final int OPERATION_REGISTER = 1;
+ private static final int OPERATION_NON_REPLICATED = 2;
+ private static final int GET_CLUSTER_METADATA_CODE = 12;
+ private static final int PING_CODE = 1;
+
+ @Test
+ void shouldResumeOnASurvivorAfterTheSignedInNodeDies() throws Exception {
+ InetAddress loopback = InetAddress.getLoopbackAddress();
+ try (ServerSocket primarySocket = new ServerSocket(0, 4, loopback);
+ ServerSocket survivorSocket = new ServerSocket(0, 4,
loopback)) {
+ int primaryPort = primarySocket.getLocalPort();
+ int survivorPort = survivorSocket.getLocalPort();
+ AtomicInteger survivorRegistrations = new AtomicInteger();
+ AtomicInteger survivorPings = new AtomicInteger();
+
+ // The primary leads, so the sign-in settles there and the roster
is
+ // only remembered -- not acted on -- until the node dies.
+ MockNode primary = MockNode.serve(primarySocket, request -> {
+ if (request.is(GET_CLUSTER_METADATA_CODE,
OPERATION_NON_REPLICATED)) {
+ return Response.success(
+ OPERATION_NON_REPLICATED,
clusterMetadata(primaryPort, survivorPort, primaryPort));
+ }
+ if (request.operation() == OPERATION_REGISTER) {
+ return Response.success(OPERATION_REGISTER,
registerBody(1));
+ }
+ return Response.success(OPERATION_NON_REPLICATED,
Unpooled.EMPTY_BUFFER);
+ });
+ MockNode survivor = MockNode.serve(survivorSocket, request -> {
+ if (request.is(GET_CLUSTER_METADATA_CODE,
OPERATION_NON_REPLICATED)) {
+ return Response.success(
+ OPERATION_NON_REPLICATED,
clusterMetadata(primaryPort, survivorPort, survivorPort));
+ }
+ if (request.operation() == OPERATION_REGISTER) {
+ survivorRegistrations.incrementAndGet();
+ return Response.success(OPERATION_REGISTER,
registerBody(2));
+ }
+ if (request.is(PING_CODE, OPERATION_NON_REPLICATED)) {
+ survivorPings.incrementAndGet();
+ }
+ return Response.success(OPERATION_NON_REPLICATED,
Unpooled.EMPTY_BUFFER);
+ });
+
+ AsyncIggyTcpClient client = AsyncIggyTcpClient.builder()
+ .host(loopback.getHostAddress())
+ .port(primaryPort)
+ .credentials("iggy", "iggy")
Review Comment:
Both done. The client in `shouldResumeOnASurvivorAfterTheSignedInNodeDies`
is built without credentials and signs in through `users().login("iggy",
"iggy")`, so only the remembered sign-in can restore the session. And
`shouldNotResurrectASignedOutSessionOnASurvivor` signs in, logs out, kills the
node, lets every attempt fail and asserts zero registrations on the survivor.
##########
foreign/csharp/Iggy_SDK/IggyClient/Implementations/TcpMessageStream.cs:
##########
@@ -999,15 +1027,18 @@ private async Task TryEstablishConnectionAsync(bool
autoLogin, CancellationToken
var retryCount = 0;
var redirects = 0;
var delay = _configuration.ReconnectionSettings.InitialDelay;
+
+ if (string.IsNullOrEmpty(_currentAddress))
+ {
+ _currentAddress = _configuration.BaseAddress;
+ }
+
+ var candidates = DialCandidates();
Review Comment:
Fixed. When more than one candidate is queued the dial runs under a linked
CTS with `CancelAfter(FailoverDialTimeout)` (2s, like Rust), and the same token
is passed to `CreateSslStreamAndAuthenticate` — via the
`SslClientAuthenticationOptions` overload of `AuthenticateAsClientAsync` —
because the handshake had no deadline either. The fatal-exception filter now
reads `e is OperationCanceledException && token.IsCancellationRequested`, so a
bound that expired advances the candidate instead of throwing out of the sweep.
No functional C# test: a black-holed address is not portable, and a C#
client only ever gets a second candidate from a learned roster, so the setup
would be several mocks deep. The same behaviour is pinned in the Rust
(`an_endpoint_that_never_answers_the_handshake_does_not_hold_up_the_sweep`) and
Go (`TestFailover_BoundsTheDialWhenOtherEndpointsAreQueuedBehindIt`) suites.
Say the word if you want it here as well and I will build the roster fixture.
--
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]