From 918b2afb886da8ad234e99b0c31e7f393f999e98 Mon Sep 17 00:00:00 2001 From: Rahul Yadav Date: Wed, 17 Jun 2026 14:04:51 +0530 Subject: [PATCH 1/2] test(spanner): deflake location-aware retry harness test --- ...nAwareSharedBackendReplicaHarnessTest.java | 183 +++++++++++++----- 1 file changed, 138 insertions(+), 45 deletions(-) diff --git a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/LocationAwareSharedBackendReplicaHarnessTest.java b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/LocationAwareSharedBackendReplicaHarnessTest.java index 6ea248288ac1..4c6e600d372f 100644 --- a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/LocationAwareSharedBackendReplicaHarnessTest.java +++ b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/LocationAwareSharedBackendReplicaHarnessTest.java @@ -57,13 +57,11 @@ import org.junit.AfterClass; import org.junit.Before; import org.junit.BeforeClass; -import org.junit.Ignore; import org.junit.Test; import org.junit.runner.RunWith; import org.junit.runners.JUnit4; @RunWith(JUnit4.class) -@Ignore("Flaky test, tracked in b/509979376") public class LocationAwareSharedBackendReplicaHarnessTest { private static final String PROJECT = "fake-project"; @@ -437,8 +435,7 @@ public void singleUseReadMidStreamRecvFailureWithoutRetryInfoRetriesForBypassTra DatabaseClient client = spanner.getDatabaseClient(DatabaseId.of(PROJECT, INSTANCE, DATABASE)); seedLocationMetadata(client); - int firstReplicaIndex = waitForReplicaRoutedRead(client, harness); - int secondReplicaIndex = 1 - firstReplicaIndex; + waitForReplicaRoutedRead(client, harness); harness.clearRequests(); harness.backend.setStreamingReadExecutionTime( @@ -459,54 +456,30 @@ public void singleUseReadMidStreamRecvFailureWithoutRetryInfoRetriesForBypassTra } assertEquals(Arrays.asList("b", "c", "d"), rows); + String diagnostics = routingDiagnostics(harness); + List attempts = streamingReadAttempts(harness); assertEquals( - 1, - harness - .replicas - .get(firstReplicaIndex) - .getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .size()); - assertEquals( - 1, - harness - .replicas - .get(secondReplicaIndex) - .getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .size()); + "Expected original request and retry request on replica endpoints.\n" + diagnostics, + 2, + attempts.size()); assertEquals( + "Default replica should not receive bypass traffic.\n" + diagnostics, 0, harness .defaultReplica .getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ) .size()); - ReadRequest replicaARequest = - (ReadRequest) - harness - .replicas - .get(firstReplicaIndex) - .getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .get(0); - ReadRequest replicaBRequest = - (ReadRequest) - harness - .replicas - .get(secondReplicaIndex) - .getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .get(0); - assertTrue(replicaARequest.getResumeToken().isEmpty()); - assertEquals(RESUME_TOKEN_AFTER_FIRST_ROW, replicaBRequest.getResumeToken()); + ReplicaRequestAttempt originalRequest = + findSingleAttemptWithResumeToken(harness, attempts, ByteString.EMPTY); + ReplicaRequestAttempt retryRequest = + findSingleAttemptWithResumeToken(harness, attempts, RESUME_TOKEN_AFTER_FIRST_ROW); + assertNotEquals( + "Retry should reroute to a different replica.\n" + diagnostics, + originalRequest.replicaIndex, + retryRequest.replicaIndex); assertRetriedOnSameLogicalRequest( - harness - .replicas - .get(firstReplicaIndex) - .getRequestIds(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .get(0), - harness - .replicas - .get(secondReplicaIndex) - .getRequestIds(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .get(0)); + originalRequest.requestId, retryRequest.requestId, diagnostics); } } @@ -784,6 +757,106 @@ private static RoutingHint exactReadRoutingHint(RecipeList recipes) { return request.getRoutingHint(); } + private static final class ReplicaRequestAttempt { + private final int replicaIndex; + private final ReadRequest request; + private final String requestId; + + private ReplicaRequestAttempt(int replicaIndex, ReadRequest request, String requestId) { + this.replicaIndex = replicaIndex; + this.request = request; + this.requestId = requestId; + } + } + + private static List streamingReadAttempts( + SharedBackendReplicaHarness harness) { + List attempts = new ArrayList<>(); + for (int replicaIndex = 0; replicaIndex < harness.replicas.size(); replicaIndex++) { + SharedBackendReplicaHarness.HookedReplicaSpannerService replica = + harness.replicas.get(replicaIndex); + List requests = + replica.getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ); + List requestIds = + replica.getRequestIds(SharedBackendReplicaHarness.METHOD_STREAMING_READ); + for (int requestIndex = 0; requestIndex < requests.size(); requestIndex++) { + attempts.add( + new ReplicaRequestAttempt( + replicaIndex, + (ReadRequest) requests.get(requestIndex), + requestIndex < requestIds.size() ? requestIds.get(requestIndex) : null)); + } + } + return attempts; + } + + private static ReplicaRequestAttempt findSingleAttemptWithResumeToken( + SharedBackendReplicaHarness harness, + List attempts, + ByteString resumeToken) { + ReplicaRequestAttempt match = null; + for (ReplicaRequestAttempt attempt : attempts) { + if (attempt.request.getResumeToken().equals(resumeToken)) { + if (match != null) { + fail( + "Expected exactly one StreamingRead request with resume token " + + resumeToken.toStringUtf8() + + ", but found multiple.\n" + + routingDiagnostics(harness)); + } + match = attempt; + } + } + if (match == null) { + fail( + "Expected exactly one StreamingRead request with resume token " + + resumeToken.toStringUtf8() + + ", but found none.\n" + + routingDiagnostics(harness)); + } + return match; + } + + private static String routingDiagnostics(SharedBackendReplicaHarness harness) { + StringBuilder builder = new StringBuilder(); + appendStreamingReadDiagnostics(builder, "default", -1, harness.defaultReplica); + for (int replicaIndex = 0; replicaIndex < harness.replicas.size(); replicaIndex++) { + appendStreamingReadDiagnostics( + builder, "replica", replicaIndex, harness.replicas.get(replicaIndex)); + } + return builder.toString(); + } + + private static void appendStreamingReadDiagnostics( + StringBuilder builder, + String label, + int replicaIndex, + SharedBackendReplicaHarness.HookedReplicaSpannerService replica) { + List requests = + replica.getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ); + List requestIds = + replica.getRequestIds(SharedBackendReplicaHarness.METHOD_STREAMING_READ); + builder + .append(label) + .append(replicaIndex >= 0 ? "[" + replicaIndex + "]" : "") + .append(" requestCount=") + .append(requests.size()) + .append('\n'); + for (int requestIndex = 0; requestIndex < requests.size(); requestIndex++) { + ReadRequest request = (ReadRequest) requests.get(requestIndex); + builder + .append(" #") + .append(requestIndex) + .append(" requestId=") + .append(requestIndex < requestIds.size() ? requestIds.get(requestIndex) : "") + .append(" resumeToken=") + .append(request.getResumeToken().toStringUtf8()) + .append(" routingHint=") + .append(request.getRoutingHint()) + .append('\n'); + } + } + private static io.grpc.StatusRuntimeException resourceExhaustedWithRetryInfo(String description) { Metadata trailers = new Metadata(); trailers.put( @@ -807,10 +880,30 @@ private static StatusRuntimeException unavailable(String description) { private static void assertRetriedOnSameLogicalRequest( String firstRequestId, String secondRequestId) { + assertRetriedOnSameLogicalRequest(firstRequestId, secondRequestId, ""); + } + + private static void assertRetriedOnSameLogicalRequest( + String firstRequestId, String secondRequestId, String diagnostics) { + if (firstRequestId == null || secondRequestId == null) { + fail( + "Missing x-goog-spanner-request-id header. firstRequestId=" + + firstRequestId + + ", secondRequestId=" + + secondRequestId + + "\n" + + diagnostics); + } XGoogSpannerRequestId first = XGoogSpannerRequestId.of(firstRequestId); XGoogSpannerRequestId second = XGoogSpannerRequestId.of(secondRequestId); - assertEquals(first.getLogicalRequestKey(), second.getLogicalRequestKey()); - assertEquals(first.getAttempt() + 1, second.getAttempt()); + assertEquals( + "Retry should use same logical request.\n" + diagnostics, + first.getLogicalRequestKey(), + second.getLogicalRequestKey()); + assertEquals( + "Retry should increment request attempt.\n" + diagnostics, + first.getAttempt() + 1, + second.getAttempt()); } private static com.google.spanner.v1.ResultSet singleRowReadResultSet(String value) { From d1b08629d9f73f55fc5fce950fafb920f3305578 Mon Sep 17 00:00:00 2001 From: Rahul Yadav Date: Wed, 17 Jun 2026 15:04:56 +0530 Subject: [PATCH 2/2] format asserts to handle out of order request on servers --- ...nAwareSharedBackendReplicaHarnessTest.java | 368 +++++++++--------- 1 file changed, 179 insertions(+), 189 deletions(-) diff --git a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/LocationAwareSharedBackendReplicaHarnessTest.java b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/LocationAwareSharedBackendReplicaHarnessTest.java index 4c6e600d372f..84643ad21207 100644 --- a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/LocationAwareSharedBackendReplicaHarnessTest.java +++ b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/LocationAwareSharedBackendReplicaHarnessTest.java @@ -121,16 +121,11 @@ public void singleUseReadReroutesOnResourceExhaustedForBypassTraffic() throws Ex DatabaseClient client = spanner.getDatabaseClient(DatabaseId.of(PROJECT, INSTANCE, DATABASE)); seedLocationMetadata(client); - int firstReplicaIndex = waitForReplicaRoutedRead(client, harness); - int secondReplicaIndex = 1 - firstReplicaIndex; + waitForReplicaRoutedRead(client, harness); harness.clearRequests(); - harness - .replicas - .get(firstReplicaIndex) - .putMethodErrors( - SharedBackendReplicaHarness.METHOD_STREAMING_READ, - resourceExhausted("busy-routed-replica")); + harness.backend.setStreamingReadExecutionTime( + SimulatedExecutionTime.ofStreamException(resourceExhausted("busy-routed-replica"), 0L)); try (ResultSet resultSet = client @@ -143,45 +138,15 @@ public void singleUseReadReroutesOnResourceExhaustedForBypassTraffic() throws Ex assertTrue(resultSet.next()); } + String diagnostics = routingDiagnostics(harness); + List attempts = streamingReadAttempts(harness); assertEquals( - 1, - harness - .replicas - .get(firstReplicaIndex) - .getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .size()); - assertEquals( - 1, - harness - .replicas - .get(secondReplicaIndex) - .getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .size()); - assertEquals( - 0, - harness - .defaultReplica - .getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .size()); - ReadRequest replicaARequest = - (ReadRequest) - harness - .replicas - .get(firstReplicaIndex) - .getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .get(0); - assertTrue(replicaARequest.getResumeToken().isEmpty()); - assertRetriedOnSameLogicalRequest( - harness - .replicas - .get(firstReplicaIndex) - .getRequestIds(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .get(0), - harness - .replicas - .get(secondReplicaIndex) - .getRequestIds(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .get(0)); + "Expected original request and retry request on replica endpoints.\n" + diagnostics, + 2, + attempts.size()); + assertDefaultReplicaUnused(harness, diagnostics); + assertAllResumeTokensEmpty(attempts, diagnostics); + assertRetriedOnDifferentReplicasInEitherOrder(attempts.get(0), attempts.get(1), diagnostics); } } @@ -192,16 +157,12 @@ public void singleUseReadCooldownSkipsReplicaOnNextRequestForBypassTraffic() thr DatabaseClient client = spanner.getDatabaseClient(DatabaseId.of(PROJECT, INSTANCE, DATABASE)); seedLocationMetadata(client); - int firstReplicaIndex = waitForReplicaRoutedRead(client, harness); - int secondReplicaIndex = 1 - firstReplicaIndex; + waitForReplicaRoutedRead(client, harness); harness.clearRequests(); - harness - .replicas - .get(firstReplicaIndex) - .putMethodErrors( - SharedBackendReplicaHarness.METHOD_STREAMING_READ, - resourceExhaustedWithRetryInfo("busy-routed-replica")); + harness.backend.setStreamingReadExecutionTime( + SimulatedExecutionTime.ofStreamException( + resourceExhaustedWithRetryInfo("busy-routed-replica"), 0L)); try (ResultSet firstRead = client @@ -225,49 +186,23 @@ public void singleUseReadCooldownSkipsReplicaOnNextRequestForBypassTraffic() thr assertTrue(secondRead.next()); } + String diagnostics = routingDiagnostics(harness); + List attempts = streamingReadAttempts(harness); assertEquals( - 1, - harness - .replicas - .get(firstReplicaIndex) - .getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .size()); - assertEquals( - 2, - harness - .replicas - .get(secondReplicaIndex) - .getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .size()); - assertEquals( - 0, - harness - .defaultReplica - .getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .size()); - List replicaBRequests = - harness - .replicas - .get(secondReplicaIndex) - .getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ); - for (AbstractMessage request : replicaBRequests) { - assertTrue(((ReadRequest) request).getResumeToken().isEmpty()); - } - List replicaBRequestIds = - harness - .replicas - .get(secondReplicaIndex) - .getRequestIds(SharedBackendReplicaHarness.METHOD_STREAMING_READ); + "Expected failed original request, retry request, and next request.\n" + diagnostics, + 3, + attempts.size()); + assertDefaultReplicaUnused(harness, diagnostics); + assertAllResumeTokensEmpty(attempts, diagnostics); + + int originalReplicaIndex = findSingleReplicaWithRequestCount(harness, 1, diagnostics); + int routedReplicaIndex = findSingleReplicaWithRequestCount(harness, 2, diagnostics); + ReplicaRequestAttempt originalAttempt = + attemptsForReplica(attempts, originalReplicaIndex).get(0); + List routedAttempts = attemptsForReplica(attempts, routedReplicaIndex); assertRetriedOnSameLogicalRequest( - harness - .replicas - .get(firstReplicaIndex) - .getRequestIds(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .get(0), - replicaBRequestIds.get(0)); - assertNotEquals( - XGoogSpannerRequestId.of(replicaBRequestIds.get(0)).getLogicalRequestKey(), - XGoogSpannerRequestId.of(replicaBRequestIds.get(1)).getLogicalRequestKey()); + originalAttempt.requestId, routedAttempts.get(0).requestId, diagnostics); + assertDifferentLogicalRequests(routedAttempts.get(0), routedAttempts.get(1), diagnostics); } } @@ -278,15 +213,11 @@ public void singleUseReadReroutesOnUnavailableForBypassTraffic() throws Exceptio DatabaseClient client = spanner.getDatabaseClient(DatabaseId.of(PROJECT, INSTANCE, DATABASE)); seedLocationMetadata(client); - int firstReplicaIndex = waitForReplicaRoutedRead(client, harness); - int secondReplicaIndex = 1 - firstReplicaIndex; + waitForReplicaRoutedRead(client, harness); harness.clearRequests(); - harness - .replicas - .get(firstReplicaIndex) - .putMethodErrors( - SharedBackendReplicaHarness.METHOD_STREAMING_READ, unavailable("isolated-replica")); + harness.backend.setStreamingReadExecutionTime( + SimulatedExecutionTime.ofStreamException(unavailable("isolated-replica"), 0L)); try (ResultSet resultSet = client @@ -299,45 +230,15 @@ public void singleUseReadReroutesOnUnavailableForBypassTraffic() throws Exceptio assertTrue(resultSet.next()); } + String diagnostics = routingDiagnostics(harness); + List attempts = streamingReadAttempts(harness); assertEquals( - 1, - harness - .replicas - .get(firstReplicaIndex) - .getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .size()); - assertEquals( - 1, - harness - .replicas - .get(secondReplicaIndex) - .getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .size()); - assertEquals( - 0, - harness - .defaultReplica - .getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .size()); - ReadRequest replicaARequest = - (ReadRequest) - harness - .replicas - .get(firstReplicaIndex) - .getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .get(0); - assertTrue(replicaARequest.getResumeToken().isEmpty()); - assertRetriedOnSameLogicalRequest( - harness - .replicas - .get(firstReplicaIndex) - .getRequestIds(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .get(0), - harness - .replicas - .get(secondReplicaIndex) - .getRequestIds(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .get(0)); + "Expected original request and retry request on replica endpoints.\n" + diagnostics, + 2, + attempts.size()); + assertDefaultReplicaUnused(harness, diagnostics); + assertAllResumeTokensEmpty(attempts, diagnostics); + assertRetriedOnDifferentReplicasInEitherOrder(attempts.get(0), attempts.get(1), diagnostics); } } @@ -349,15 +250,11 @@ public void singleUseReadCooldownSkipsUnavailableReplicaOnNextRequestForBypassTr DatabaseClient client = spanner.getDatabaseClient(DatabaseId.of(PROJECT, INSTANCE, DATABASE)); seedLocationMetadata(client); - int firstReplicaIndex = waitForReplicaRoutedRead(client, harness); - int secondReplicaIndex = 1 - firstReplicaIndex; + waitForReplicaRoutedRead(client, harness); harness.clearRequests(); - harness - .replicas - .get(firstReplicaIndex) - .putMethodErrors( - SharedBackendReplicaHarness.METHOD_STREAMING_READ, unavailable("isolated-replica")); + harness.backend.setStreamingReadExecutionTime( + SimulatedExecutionTime.ofStreamException(unavailable("isolated-replica"), 0L)); try (ResultSet firstRead = client @@ -381,49 +278,23 @@ public void singleUseReadCooldownSkipsUnavailableReplicaOnNextRequestForBypassTr assertTrue(secondRead.next()); } + String diagnostics = routingDiagnostics(harness); + List attempts = streamingReadAttempts(harness); assertEquals( - 1, - harness - .replicas - .get(firstReplicaIndex) - .getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .size()); - assertEquals( - 2, - harness - .replicas - .get(secondReplicaIndex) - .getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .size()); - assertEquals( - 0, - harness - .defaultReplica - .getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .size()); - List replicaBRequests = - harness - .replicas - .get(secondReplicaIndex) - .getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ); - for (AbstractMessage request : replicaBRequests) { - assertTrue(((ReadRequest) request).getResumeToken().isEmpty()); - } - List replicaBRequestIds = - harness - .replicas - .get(secondReplicaIndex) - .getRequestIds(SharedBackendReplicaHarness.METHOD_STREAMING_READ); + "Expected failed original request, retry request, and next request.\n" + diagnostics, + 3, + attempts.size()); + assertDefaultReplicaUnused(harness, diagnostics); + assertAllResumeTokensEmpty(attempts, diagnostics); + + int originalReplicaIndex = findSingleReplicaWithRequestCount(harness, 1, diagnostics); + int routedReplicaIndex = findSingleReplicaWithRequestCount(harness, 2, diagnostics); + ReplicaRequestAttempt originalAttempt = + attemptsForReplica(attempts, originalReplicaIndex).get(0); + List routedAttempts = attemptsForReplica(attempts, routedReplicaIndex); assertRetriedOnSameLogicalRequest( - harness - .replicas - .get(firstReplicaIndex) - .getRequestIds(SharedBackendReplicaHarness.METHOD_STREAMING_READ) - .get(0), - replicaBRequestIds.get(0)); - assertNotEquals( - XGoogSpannerRequestId.of(replicaBRequestIds.get(0)).getLogicalRequestKey(), - XGoogSpannerRequestId.of(replicaBRequestIds.get(1)).getLogicalRequestKey()); + originalAttempt.requestId, routedAttempts.get(0).requestId, diagnostics); + assertDifferentLogicalRequests(routedAttempts.get(0), routedAttempts.get(1), diagnostics); } } @@ -779,17 +650,111 @@ private static List streamingReadAttempts( replica.getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ); List requestIds = replica.getRequestIds(SharedBackendReplicaHarness.METHOD_STREAMING_READ); + assertEquals( + "StreamingRead requests and request ids should be recorded in pairs.\n" + + routingDiagnostics(harness), + requests.size(), + requestIds.size()); for (int requestIndex = 0; requestIndex < requests.size(); requestIndex++) { attempts.add( new ReplicaRequestAttempt( replicaIndex, (ReadRequest) requests.get(requestIndex), - requestIndex < requestIds.size() ? requestIds.get(requestIndex) : null)); + requestIds.get(requestIndex))); } } return attempts; } + private static void assertDefaultReplicaUnused( + SharedBackendReplicaHarness harness, String diagnostics) { + assertEquals( + "Expected default endpoint to receive no StreamingRead requests.\n" + diagnostics, + 0, + harness + .defaultReplica + .getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ) + .size()); + } + + private static void assertAllResumeTokensEmpty( + List attempts, String diagnostics) { + for (ReplicaRequestAttempt attempt : attempts) { + assertTrue( + "Expected all retry/retry-info requests to restart without resume token.\n" + diagnostics, + attempt.request.getResumeToken().isEmpty()); + } + } + + private static int findSingleReplicaWithRequestCount( + SharedBackendReplicaHarness harness, int expectedCount, String diagnostics) { + int match = -1; + for (int replicaIndex = 0; replicaIndex < harness.replicas.size(); replicaIndex++) { + int requestCount = + harness + .replicas + .get(replicaIndex) + .getRequests(SharedBackendReplicaHarness.METHOD_STREAMING_READ) + .size(); + if (requestCount == expectedCount) { + if (match >= 0) { + fail( + "Expected exactly one replica with " + + expectedCount + + " StreamingRead requests, but found multiple.\n" + + diagnostics); + } + match = replicaIndex; + } + } + if (match < 0) { + fail( + "Expected exactly one replica with " + + expectedCount + + " StreamingRead requests, but found none.\n" + + diagnostics); + } + return match; + } + + private static List attemptsForReplica( + List attempts, int replicaIndex) { + List result = new ArrayList<>(); + for (ReplicaRequestAttempt attempt : attempts) { + if (attempt.replicaIndex == replicaIndex) { + result.add(attempt); + } + } + return result; + } + + private static void assertRetriedOnDifferentReplicasInEitherOrder( + ReplicaRequestAttempt firstAttempt, ReplicaRequestAttempt secondAttempt, String diagnostics) { + assertNotEquals( + "Retry should reroute to different replica.\n" + diagnostics, + firstAttempt.replicaIndex, + secondAttempt.replicaIndex); + assertRetriedOnSameLogicalRequestInEitherOrder( + firstAttempt.requestId, secondAttempt.requestId, diagnostics); + } + + private static void assertDifferentLogicalRequests( + ReplicaRequestAttempt firstAttempt, ReplicaRequestAttempt secondAttempt, String diagnostics) { + if (firstAttempt.requestId == null || secondAttempt.requestId == null) { + fail( + "Missing x-goog-spanner-request-id header. firstRequestId=" + + firstAttempt.requestId + + ", secondRequestId=" + + secondAttempt.requestId + + "\n" + + diagnostics); + } + assertNotEquals( + "Expected separate reads to use distinct logical request ids.\n" + diagnostics, + XGoogSpannerRequestId.of(firstAttempt.requestId).getLogicalRequestKey(), + XGoogSpannerRequestId.of(secondAttempt.requestId).getLogicalRequestKey()); + } + private static ReplicaRequestAttempt findSingleAttemptWithResumeToken( SharedBackendReplicaHarness harness, List attempts, @@ -841,6 +806,8 @@ private static void appendStreamingReadDiagnostics( .append(replicaIndex >= 0 ? "[" + replicaIndex + "]" : "") .append(" requestCount=") .append(requests.size()) + .append(" requestIdCount=") + .append(requestIds.size()) .append('\n'); for (int requestIndex = 0; requestIndex < requests.size(); requestIndex++) { ReadRequest request = (ReadRequest) requests.get(requestIndex); @@ -878,6 +845,29 @@ private static StatusRuntimeException unavailable(String description) { return Status.UNAVAILABLE.withDescription(description).asRuntimeException(); } + private static void assertRetriedOnSameLogicalRequestInEitherOrder( + String firstRequestId, String secondRequestId, String diagnostics) { + if (firstRequestId == null || secondRequestId == null) { + fail( + "Missing x-goog-spanner-request-id header. firstRequestId=" + + firstRequestId + + ", secondRequestId=" + + secondRequestId + + "\n" + + diagnostics); + } + XGoogSpannerRequestId first = XGoogSpannerRequestId.of(firstRequestId); + XGoogSpannerRequestId second = XGoogSpannerRequestId.of(secondRequestId); + assertEquals( + "Retry should use same logical request.\n" + diagnostics, + first.getLogicalRequestKey(), + second.getLogicalRequestKey()); + assertEquals( + "Retry attempts should be consecutive.\n" + diagnostics, + 1, + Math.abs(first.getAttempt() - second.getAttempt())); + } + private static void assertRetriedOnSameLogicalRequest( String firstRequestId, String secondRequestId) { assertRetriedOnSameLogicalRequest(firstRequestId, secondRequestId, "");