Skip to content

Commit 151e5c0

Browse files
test(spanner): unflake LocationAwareSharedBackendReplicaHarnessTest (#14396)
Wait for all backend replicas to establish active transport connections before issuing the warmup query in waitForReplicaRoutedRead. Previously, waitForReplicaRoutedRead exited as soon as the first replica handled the initial warmup read. In fast test executions, secondary replicas could still be completing their background Netty channel handshakes initiated by EndpointLifecycleManager probing. When the primary replica subsequently failed, secondary replicas were evaluated as unhealthy (channel state not yet READY) and skipped, causing traffic to unexpectedly fall back to defaultReplica. Tracking active transport connections via ServerTransportFilter in SharedBackendReplicaHarness ensures all replica channels are connected and ready before test assertions begin.
1 parent 305f47d commit 151e5c0

2 files changed

Lines changed: 46 additions & 0 deletions

File tree

java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/LocationAwareSharedBackendReplicaHarnessTest.java

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616

1717
package com.google.cloud.spanner;
1818

19+
import static org.awaitility.Awaitility.await;
1920
import static org.junit.Assert.assertEquals;
2021
import static org.junit.Assert.assertNotEquals;
2122
import static org.junit.Assert.assertTrue;
@@ -49,6 +50,7 @@
4950
import io.grpc.Status;
5051
import io.grpc.StatusRuntimeException;
5152
import io.grpc.protobuf.ProtoUtils;
53+
import java.time.Duration;
5254
import java.util.ArrayList;
5355
import java.util.Arrays;
5456
import java.util.List;
@@ -573,8 +575,13 @@ private static boolean drain(ResultSet resultSet) {
573575
return sawRow;
574576
}
575577

578+
private static void waitForAllReplicasConnected(SharedBackendReplicaHarness harness) {
579+
await().atMost(Duration.ofSeconds(10)).until(harness::allReplicasConnected);
580+
}
581+
576582
private static int waitForReplicaRoutedRead(
577583
DatabaseClient client, SharedBackendReplicaHarness harness) throws InterruptedException {
584+
waitForAllReplicasConnected(harness);
578585
long deadlineNanos = System.nanoTime() + TimeUnit.SECONDS.toNanos(10);
579586
while (System.nanoTime() < deadlineNanos) {
580587
try (ResultSet resultSet =

java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SharedBackendReplicaHarness.java

Lines changed: 39 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -34,12 +34,14 @@
3434
import com.google.spanner.v1.Session;
3535
import com.google.spanner.v1.SpannerGrpc;
3636
import com.google.spanner.v1.Transaction;
37+
import io.grpc.Attributes;
3738
import io.grpc.Metadata;
3839
import io.grpc.Server;
3940
import io.grpc.ServerCall;
4041
import io.grpc.ServerCallHandler;
4142
import io.grpc.ServerInterceptor;
4243
import io.grpc.ServerInterceptors;
44+
import io.grpc.ServerTransportFilter;
4345
import io.grpc.netty.shaded.io.grpc.netty.NettyServerBuilder;
4446
import io.grpc.stub.StreamObserver;
4547
import java.io.Closeable;
@@ -50,6 +52,7 @@
5052
import java.util.HashMap;
5153
import java.util.List;
5254
import java.util.Map;
55+
import java.util.concurrent.atomic.AtomicInteger;
5356

5457
/** Shared-backend replica harness for end-to-end location-aware routing tests. */
5558
final class SharedBackendReplicaHarness implements Closeable {
@@ -71,11 +74,24 @@ static final class HookedReplicaSpannerService extends SpannerGrpc.SpannerImplBa
7174
private final Map<String, ArrayDeque<Throwable>> methodErrors = new HashMap<>();
7275
private final Map<String, List<AbstractMessage>> requests = new HashMap<>();
7376
private final Map<String, List<String>> requestIds = new HashMap<>();
77+
private final AtomicInteger activeConnections = new AtomicInteger();
7478

7579
private HookedReplicaSpannerService(MockSpannerServiceImpl backend) {
7680
this.backend = backend;
7781
}
7882

83+
void recordConnectionReady() {
84+
activeConnections.incrementAndGet();
85+
}
86+
87+
void recordConnectionTerminated() {
88+
activeConnections.decrementAndGet();
89+
}
90+
91+
boolean hasConnected() {
92+
return activeConnections.get() > 0;
93+
}
94+
7995
synchronized void putMethodErrors(String method, Throwable... errors) {
8096
ArrayDeque<Throwable> queue = new ArrayDeque<>();
8197
for (Throwable error : errors) {
@@ -278,12 +294,35 @@ public <ReqT, RespT> ServerCall.Listener<ReqT> interceptCall(
278294
Server server =
279295
NettyServerBuilder.forAddress(address)
280296
.addService(ServerInterceptors.intercept(service, interceptor))
297+
.addTransportFilter(
298+
new ServerTransportFilter() {
299+
@Override
300+
public Attributes transportReady(Attributes transportAttrs) {
301+
service.recordConnectionReady();
302+
return super.transportReady(transportAttrs);
303+
}
304+
305+
@Override
306+
public void transportTerminated(Attributes transportAttrs) {
307+
service.recordConnectionTerminated();
308+
super.transportTerminated(transportAttrs);
309+
}
310+
})
281311
.build()
282312
.start();
283313
servers.add(server);
284314
return "localhost:" + server.getPort();
285315
}
286316

317+
boolean allReplicasConnected() {
318+
for (HookedReplicaSpannerService replica : replicas) {
319+
if (!replica.hasConnected()) {
320+
return false;
321+
}
322+
}
323+
return true;
324+
}
325+
287326
void clearRequests() {
288327
defaultReplica.clearRequests();
289328
for (HookedReplicaSpannerService replica : replicas) {

0 commit comments

Comments
 (0)