From 43891a2f2be5a2868e00ed569be22b372b5c369e Mon Sep 17 00:00:00 2001 From: Mattie Fu Date: Thu, 17 Sep 2026 18:40:57 +0000 Subject: [PATCH 1/2] fix(bigtable): preserve confirmed entries on mid-stream error and detect omitted entries --- .../mutaterows/MutateRowsAttemptCallable.java | 108 ++++++++++++++++- .../MutateRowsRetryingCallable.java | 6 +- .../MutateRowsAttemptCallableTest.java | 110 ++++++++++++++++++ 3 files changed, 222 insertions(+), 2 deletions(-) diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsAttemptCallable.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsAttemptCallable.java index 36c2930bdac6..39b4c6bb6207 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsAttemptCallable.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsAttemptCallable.java @@ -18,12 +18,16 @@ import com.google.api.core.ApiFunction; import com.google.api.core.ApiFuture; import com.google.api.core.ApiFutures; +import com.google.api.core.SettableApiFuture; import com.google.api.gax.grpc.GrpcStatusCode; import com.google.api.gax.retrying.RetryingFuture; import com.google.api.gax.rpc.ApiCallContext; import com.google.api.gax.rpc.ApiException; import com.google.api.gax.rpc.ApiExceptionFactory; +import com.google.api.gax.rpc.ServerStreamingCallable; +import com.google.api.gax.rpc.StateCheckingResponseObserver; import com.google.api.gax.rpc.StatusCode; +import com.google.api.gax.rpc.StreamController; import com.google.api.gax.rpc.UnaryCallable; import com.google.bigtable.v2.MutateRowsRequest; import com.google.bigtable.v2.MutateRowsRequest.Builder; @@ -35,8 +39,10 @@ import com.google.common.base.Preconditions; import com.google.common.collect.ImmutableList; import com.google.common.collect.Lists; +import com.google.common.primitives.Ints; import com.google.common.util.concurrent.MoreExecutors; import com.google.rpc.Code; +import java.util.Collections; import java.util.List; import java.util.Set; import java.util.concurrent.Callable; @@ -212,6 +218,58 @@ public Void call() { return null; } + static class PartialResponseException extends RuntimeException { + private final List partialResponses; + + PartialResponseException(Throwable cause, List partialResponses) { + super(cause); + this.partialResponses = partialResponses; + } + + List getPartialResponses() { + return partialResponses; + } + } + + static UnaryCallable> createStreamingCollector( + final ServerStreamingCallable streamingCallable) { + return new UnaryCallable>() { + @Override + public ApiFuture> futureCall( + MutateRowsRequest request, ApiCallContext context) { + final SettableApiFuture> future = SettableApiFuture.create(); + final List responses = Lists.newArrayList(); + streamingCallable.call( + request, + new StateCheckingResponseObserver() { + @Override + protected void onStartImpl(StreamController controller) {} + + @Override + protected void onResponseImpl(MutateRowsResponse response) { + responses.add(response); + } + + @Override + protected void onErrorImpl(Throwable t) { + if (!responses.isEmpty()) { + future.setException(new PartialResponseException(t, responses)); + } else { + future.setException(t); + } + } + + @Override + protected void onCompleteImpl() { + future.set(responses); + } + }, + context); + return future; + } + }; + } + /** * Handle an RPC level failure by generating a {@link FailedMutation} for each expected entry. The * newly generated {@link FailedMutation}s will be combined with the permanentFailures to give the @@ -219,6 +277,12 @@ public Void call() { * MutateRowsException}. */ private void handleAttemptError(Throwable rpcError) { + List partialResponses = Collections.emptyList(); + if (rpcError instanceof PartialResponseException) { + PartialResponseException partialEx = (PartialResponseException) rpcError; + partialResponses = partialEx.getPartialResponses(); + rpcError = partialEx.getCause(); + } ApiException entryError = createSyntheticErrorForRpcFailure(rpcError); ImmutableList.Builder allFailures = ImmutableList.builder(); MutateRowsRequest lastRequest = currentRequest; @@ -227,8 +291,31 @@ private void handleAttemptError(Throwable rpcError) { Builder builder = lastRequest.toBuilder().clearEntries(); List newOriginalIndexes = Lists.newArrayList(); + boolean[] seenIndices = new boolean[currentRequest.getEntriesCount()]; + + for (MutateRowsResponse response : partialResponses) { + for (Entry entry : response.getEntriesList()) { + seenIndices[Ints.checkedCast(entry.getIndex())] = true; + if (entry.getStatus().getCode() == Code.OK_VALUE) { + continue; + } + int origIndex = getOriginalIndex((int) entry.getIndex()); + FailedMutation failedMutation = + FailedMutation.create(origIndex, createEntryError(entry.getStatus())); + allFailures.add(failedMutation); + if (!failedMutation.getError().isRetryable()) { + permanentFailures.add(failedMutation); + } else { + newOriginalIndexes.add(origIndex); + builder.addEntries(lastRequest.getEntries((int) entry.getIndex())); + } + } + } for (int i = 0; i < currentRequest.getEntriesCount(); i++) { + if (seenIndices[i]) { + continue; + } int origIndex = getOriginalIndex(i); FailedMutation failedMutation = FailedMutation.create(origIndex, entryError); @@ -247,7 +334,7 @@ private void handleAttemptError(Throwable rpcError) { currentRequest = builder.build(); originalIndexes = newOriginalIndexes; - throw new MutateRowsException(rpcError, allFailures.build(), entryError.isRetryable()); + throw new MutateRowsException(rpcError, allFailures.build(), builder.getEntriesCount() > 0); } /** @@ -263,9 +350,11 @@ private void handleAttemptSuccess(List responses) { Builder builder = lastRequest.toBuilder().clearEntries(); List newOriginalIndexes = Lists.newArrayList(); + boolean[] seenIndices = new boolean[currentRequest.getEntriesCount()]; for (MutateRowsResponse response : responses) { for (Entry entry : response.getEntriesList()) { + seenIndices[Ints.checkedCast(entry.getIndex())] = true; if (entry.getStatus().getCode() == Code.OK_VALUE) { continue; } @@ -288,6 +377,23 @@ private void handleAttemptSuccess(List responses) { } } + for (int i = 0; i < seenIndices.length; i++) { + if (seenIndices[i]) { + continue; + } + int origIndex = getOriginalIndex(i); + FailedMutation failedMutation = + FailedMutation.create( + origIndex, + ApiExceptionFactory.createException( + "Missing entry response for entry " + origIndex, + null, + GrpcStatusCode.of(io.grpc.Status.Code.INTERNAL), + false)); + allFailures.add(failedMutation); + permanentFailures.add(failedMutation); + } + currentRequest = builder.build(); originalIndexes = newOriginalIndexes; diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsRetryingCallable.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsRetryingCallable.java index ff0daf78bb2d..24e6ee90d5c8 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsRetryingCallable.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsRetryingCallable.java @@ -60,7 +60,11 @@ public MutateRowsRetryingCallable( public RetryingFuture futureCall(MutateRowsRequest request, ApiCallContext inputContext) { ApiCallContext context = callContextPrototype.nullToSelf(inputContext); MutateRowsAttemptCallable retryCallable = - new MutateRowsAttemptCallable(callable.all(), request, context, retryCodes); + new MutateRowsAttemptCallable( + MutateRowsAttemptCallable.createStreamingCollector(callable), + request, + context, + retryCodes); RetryingFuture retryingFuture = executor.createFuture(retryCallable, context); retryCallable.setExternalFuture(retryingFuture); diff --git a/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsAttemptCallableTest.java b/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsAttemptCallableTest.java index 2df2aaf2b4a9..f1883053c97e 100644 --- a/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsAttemptCallableTest.java +++ b/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsAttemptCallableTest.java @@ -114,6 +114,116 @@ public void testNoRpcTimeout() { assertThat(innerCallable.lastContext.getTimeout()).isNull(); } + @Test + public void partialOmissionMultiEntryTest() throws Exception { + MutateRowsRequest request = + MutateRowsRequest.newBuilder() + .addEntries(Entry.newBuilder().setRowKey(ByteString.copyFromUtf8("0-ok"))) + .addEntries(Entry.newBuilder().setRowKey(ByteString.copyFromUtf8("1-omitted"))) + .addEntries(Entry.newBuilder().setRowKey(ByteString.copyFromUtf8("2-ok"))) + .build(); + innerCallable.response.add( + MutateRowsResponse.newBuilder() + .addEntries( + MutateRowsResponse.Entry.newBuilder().setIndex(0).setStatus(OK_STATUS_PROTO)) + .addEntries( + MutateRowsResponse.Entry.newBuilder().setIndex(2).setStatus(OK_STATUS_PROTO)) + .build()); + + MutateRowsAttemptCallable attemptCallable = + new MutateRowsAttemptCallable(innerCallable, request, callContext, retryCodes); + attemptCallable.setExternalFuture(parentFuture); + attemptCallable.call(); + + Throwable actualError = null; + try { + parentFuture.attemptFuture.get(); + } catch (Throwable t) { + actualError = t.getCause(); + } + + assertThat(actualError).isInstanceOf(MutateRowsException.class); + MutateRowsException mutateRowsException = (MutateRowsException) actualError; + assertThat(mutateRowsException.isRetryable()).isFalse(); + assertThat(mutateRowsException.getFailedMutations()).hasSize(1); + FailedMutation failedMutation = mutateRowsException.getFailedMutations().get(0); + assertThat(failedMutation.getIndex()).isEqualTo(1); + assertThat(failedMutation.getError().getStatusCode().getCode()).isEqualTo(Code.INTERNAL); + assertThat(failedMutation.getError().isRetryable()).isFalse(); + assertThat(failedMutation.getError()) + .hasMessageThat() + .contains("Missing entry response for entry 1"); + } + + @Test + public void midStreamErrorPreservesConfirmedEntriesTest() throws Exception { + MutateRowsRequest request = + MutateRowsRequest.newBuilder() + .addEntries(Entry.newBuilder().setRowKey(ByteString.copyFromUtf8("0-ok"))) + .addEntries(Entry.newBuilder().setRowKey(ByteString.copyFromUtf8("1-unavailable"))) + .build(); + + final UnavailableException rpcError = + new UnavailableException( + "mid-stream error", null, GrpcStatusCode.of(io.grpc.Status.Code.UNAVAILABLE), true); + final List partialResponses = + Lists.newArrayList( + MutateRowsResponse.newBuilder() + .addEntries( + MutateRowsResponse.Entry.newBuilder().setIndex(0).setStatus(OK_STATUS_PROTO)) + .build()); + + final boolean[] firstCall = {true}; + final List capturedRequests = Lists.newArrayList(); + + UnaryCallable> mockStreamingCallable = + new UnaryCallable>() { + @Override + public ApiFuture> futureCall( + MutateRowsRequest req, ApiCallContext context) { + capturedRequests.add(req); + if (firstCall[0]) { + firstCall[0] = false; + return ApiFutures.immediateFailedFuture( + new MutateRowsAttemptCallable.PartialResponseException( + rpcError, partialResponses)); + } + return ApiFutures.immediateFuture( + Lists.newArrayList( + MutateRowsResponse.newBuilder() + .addEntries( + MutateRowsResponse.Entry.newBuilder() + .setIndex(0) + .setStatus(OK_STATUS_PROTO)) + .build())); + } + }; + + MutateRowsAttemptCallable attemptCallable = + new MutateRowsAttemptCallable(mockStreamingCallable, request, callContext, retryCodes); + attemptCallable.setExternalFuture(parentFuture); + + attemptCallable.call(); + Throwable actualError = null; + try { + parentFuture.attemptFuture.get(); + } catch (Throwable t) { + actualError = t.getCause(); + } + assertThat(actualError).isInstanceOf(MutateRowsException.class); + MutateRowsException mutateRowsException = (MutateRowsException) actualError; + assertThat(mutateRowsException.isRetryable()).isTrue(); + assertThat(mutateRowsException.getFailedMutations()).hasSize(1); + assertThat(mutateRowsException.getFailedMutations().get(0).getIndex()).isEqualTo(1); + + attemptCallable.call(); + parentFuture.attemptFuture.get(); + assertThat(capturedRequests).hasSize(2); + assertThat(capturedRequests.get(1).getEntriesCount()).isEqualTo(1); + assertThat(capturedRequests.get(1).getEntries(0).getRowKey().toStringUtf8()) + .isEqualTo("1-unavailable"); + } + @Test public void mixedTest() { // Setup the request & response From 3a15233a075cb152a34bd6d3fc6b4127f17321f5 Mon Sep 17 00:00:00 2001 From: Mattie Fu Date: Fri, 18 Sep 2026 21:04:12 +0000 Subject: [PATCH 2/2] test(bigtable): revert mid-stream response handling and keep omitted entry test --- .../mutaterows/MutateRowsAttemptCallable.java | 108 +----------------- .../MutateRowsRetryingCallable.java | 6 +- .../MutateRowsAttemptCallableTest.java | 70 ------------ 3 files changed, 2 insertions(+), 182 deletions(-) diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsAttemptCallable.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsAttemptCallable.java index 39b4c6bb6207..36c2930bdac6 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsAttemptCallable.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsAttemptCallable.java @@ -18,16 +18,12 @@ import com.google.api.core.ApiFunction; import com.google.api.core.ApiFuture; import com.google.api.core.ApiFutures; -import com.google.api.core.SettableApiFuture; import com.google.api.gax.grpc.GrpcStatusCode; import com.google.api.gax.retrying.RetryingFuture; import com.google.api.gax.rpc.ApiCallContext; import com.google.api.gax.rpc.ApiException; import com.google.api.gax.rpc.ApiExceptionFactory; -import com.google.api.gax.rpc.ServerStreamingCallable; -import com.google.api.gax.rpc.StateCheckingResponseObserver; import com.google.api.gax.rpc.StatusCode; -import com.google.api.gax.rpc.StreamController; import com.google.api.gax.rpc.UnaryCallable; import com.google.bigtable.v2.MutateRowsRequest; import com.google.bigtable.v2.MutateRowsRequest.Builder; @@ -39,10 +35,8 @@ import com.google.common.base.Preconditions; import com.google.common.collect.ImmutableList; import com.google.common.collect.Lists; -import com.google.common.primitives.Ints; import com.google.common.util.concurrent.MoreExecutors; import com.google.rpc.Code; -import java.util.Collections; import java.util.List; import java.util.Set; import java.util.concurrent.Callable; @@ -218,58 +212,6 @@ public Void call() { return null; } - static class PartialResponseException extends RuntimeException { - private final List partialResponses; - - PartialResponseException(Throwable cause, List partialResponses) { - super(cause); - this.partialResponses = partialResponses; - } - - List getPartialResponses() { - return partialResponses; - } - } - - static UnaryCallable> createStreamingCollector( - final ServerStreamingCallable streamingCallable) { - return new UnaryCallable>() { - @Override - public ApiFuture> futureCall( - MutateRowsRequest request, ApiCallContext context) { - final SettableApiFuture> future = SettableApiFuture.create(); - final List responses = Lists.newArrayList(); - streamingCallable.call( - request, - new StateCheckingResponseObserver() { - @Override - protected void onStartImpl(StreamController controller) {} - - @Override - protected void onResponseImpl(MutateRowsResponse response) { - responses.add(response); - } - - @Override - protected void onErrorImpl(Throwable t) { - if (!responses.isEmpty()) { - future.setException(new PartialResponseException(t, responses)); - } else { - future.setException(t); - } - } - - @Override - protected void onCompleteImpl() { - future.set(responses); - } - }, - context); - return future; - } - }; - } - /** * Handle an RPC level failure by generating a {@link FailedMutation} for each expected entry. The * newly generated {@link FailedMutation}s will be combined with the permanentFailures to give the @@ -277,12 +219,6 @@ protected void onCompleteImpl() { * MutateRowsException}. */ private void handleAttemptError(Throwable rpcError) { - List partialResponses = Collections.emptyList(); - if (rpcError instanceof PartialResponseException) { - PartialResponseException partialEx = (PartialResponseException) rpcError; - partialResponses = partialEx.getPartialResponses(); - rpcError = partialEx.getCause(); - } ApiException entryError = createSyntheticErrorForRpcFailure(rpcError); ImmutableList.Builder allFailures = ImmutableList.builder(); MutateRowsRequest lastRequest = currentRequest; @@ -291,31 +227,8 @@ private void handleAttemptError(Throwable rpcError) { Builder builder = lastRequest.toBuilder().clearEntries(); List newOriginalIndexes = Lists.newArrayList(); - boolean[] seenIndices = new boolean[currentRequest.getEntriesCount()]; - - for (MutateRowsResponse response : partialResponses) { - for (Entry entry : response.getEntriesList()) { - seenIndices[Ints.checkedCast(entry.getIndex())] = true; - if (entry.getStatus().getCode() == Code.OK_VALUE) { - continue; - } - int origIndex = getOriginalIndex((int) entry.getIndex()); - FailedMutation failedMutation = - FailedMutation.create(origIndex, createEntryError(entry.getStatus())); - allFailures.add(failedMutation); - if (!failedMutation.getError().isRetryable()) { - permanentFailures.add(failedMutation); - } else { - newOriginalIndexes.add(origIndex); - builder.addEntries(lastRequest.getEntries((int) entry.getIndex())); - } - } - } for (int i = 0; i < currentRequest.getEntriesCount(); i++) { - if (seenIndices[i]) { - continue; - } int origIndex = getOriginalIndex(i); FailedMutation failedMutation = FailedMutation.create(origIndex, entryError); @@ -334,7 +247,7 @@ private void handleAttemptError(Throwable rpcError) { currentRequest = builder.build(); originalIndexes = newOriginalIndexes; - throw new MutateRowsException(rpcError, allFailures.build(), builder.getEntriesCount() > 0); + throw new MutateRowsException(rpcError, allFailures.build(), entryError.isRetryable()); } /** @@ -350,11 +263,9 @@ private void handleAttemptSuccess(List responses) { Builder builder = lastRequest.toBuilder().clearEntries(); List newOriginalIndexes = Lists.newArrayList(); - boolean[] seenIndices = new boolean[currentRequest.getEntriesCount()]; for (MutateRowsResponse response : responses) { for (Entry entry : response.getEntriesList()) { - seenIndices[Ints.checkedCast(entry.getIndex())] = true; if (entry.getStatus().getCode() == Code.OK_VALUE) { continue; } @@ -377,23 +288,6 @@ private void handleAttemptSuccess(List responses) { } } - for (int i = 0; i < seenIndices.length; i++) { - if (seenIndices[i]) { - continue; - } - int origIndex = getOriginalIndex(i); - FailedMutation failedMutation = - FailedMutation.create( - origIndex, - ApiExceptionFactory.createException( - "Missing entry response for entry " + origIndex, - null, - GrpcStatusCode.of(io.grpc.Status.Code.INTERNAL), - false)); - allFailures.add(failedMutation); - permanentFailures.add(failedMutation); - } - currentRequest = builder.build(); originalIndexes = newOriginalIndexes; diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsRetryingCallable.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsRetryingCallable.java index 24e6ee90d5c8..ff0daf78bb2d 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsRetryingCallable.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsRetryingCallable.java @@ -60,11 +60,7 @@ public MutateRowsRetryingCallable( public RetryingFuture futureCall(MutateRowsRequest request, ApiCallContext inputContext) { ApiCallContext context = callContextPrototype.nullToSelf(inputContext); MutateRowsAttemptCallable retryCallable = - new MutateRowsAttemptCallable( - MutateRowsAttemptCallable.createStreamingCollector(callable), - request, - context, - retryCodes); + new MutateRowsAttemptCallable(callable.all(), request, context, retryCodes); RetryingFuture retryingFuture = executor.createFuture(retryCallable, context); retryCallable.setExternalFuture(retryingFuture); diff --git a/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsAttemptCallableTest.java b/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsAttemptCallableTest.java index f1883053c97e..b15be131edb0 100644 --- a/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsAttemptCallableTest.java +++ b/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/mutaterows/MutateRowsAttemptCallableTest.java @@ -141,7 +141,6 @@ public void partialOmissionMultiEntryTest() throws Exception { } catch (Throwable t) { actualError = t.getCause(); } - assertThat(actualError).isInstanceOf(MutateRowsException.class); MutateRowsException mutateRowsException = (MutateRowsException) actualError; assertThat(mutateRowsException.isRetryable()).isFalse(); @@ -155,75 +154,6 @@ public void partialOmissionMultiEntryTest() throws Exception { .contains("Missing entry response for entry 1"); } - @Test - public void midStreamErrorPreservesConfirmedEntriesTest() throws Exception { - MutateRowsRequest request = - MutateRowsRequest.newBuilder() - .addEntries(Entry.newBuilder().setRowKey(ByteString.copyFromUtf8("0-ok"))) - .addEntries(Entry.newBuilder().setRowKey(ByteString.copyFromUtf8("1-unavailable"))) - .build(); - - final UnavailableException rpcError = - new UnavailableException( - "mid-stream error", null, GrpcStatusCode.of(io.grpc.Status.Code.UNAVAILABLE), true); - final List partialResponses = - Lists.newArrayList( - MutateRowsResponse.newBuilder() - .addEntries( - MutateRowsResponse.Entry.newBuilder().setIndex(0).setStatus(OK_STATUS_PROTO)) - .build()); - - final boolean[] firstCall = {true}; - final List capturedRequests = Lists.newArrayList(); - - UnaryCallable> mockStreamingCallable = - new UnaryCallable>() { - @Override - public ApiFuture> futureCall( - MutateRowsRequest req, ApiCallContext context) { - capturedRequests.add(req); - if (firstCall[0]) { - firstCall[0] = false; - return ApiFutures.immediateFailedFuture( - new MutateRowsAttemptCallable.PartialResponseException( - rpcError, partialResponses)); - } - return ApiFutures.immediateFuture( - Lists.newArrayList( - MutateRowsResponse.newBuilder() - .addEntries( - MutateRowsResponse.Entry.newBuilder() - .setIndex(0) - .setStatus(OK_STATUS_PROTO)) - .build())); - } - }; - - MutateRowsAttemptCallable attemptCallable = - new MutateRowsAttemptCallable(mockStreamingCallable, request, callContext, retryCodes); - attemptCallable.setExternalFuture(parentFuture); - - attemptCallable.call(); - Throwable actualError = null; - try { - parentFuture.attemptFuture.get(); - } catch (Throwable t) { - actualError = t.getCause(); - } - assertThat(actualError).isInstanceOf(MutateRowsException.class); - MutateRowsException mutateRowsException = (MutateRowsException) actualError; - assertThat(mutateRowsException.isRetryable()).isTrue(); - assertThat(mutateRowsException.getFailedMutations()).hasSize(1); - assertThat(mutateRowsException.getFailedMutations().get(0).getIndex()).isEqualTo(1); - - attemptCallable.call(); - parentFuture.attemptFuture.get(); - assertThat(capturedRequests).hasSize(2); - assertThat(capturedRequests.get(1).getEntriesCount()).isEqualTo(1); - assertThat(capturedRequests.get(1).getEntries(0).getRowKey().toStringUtf8()) - .isEqualTo("1-unavailable"); - } - @Test public void mixedTest() { // Setup the request & response