diff --git a/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ChunkAttemptCallable.java b/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ChunkAttemptCallable.java index d63d9cc34969..ae3dffc2182c 100644 --- a/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ChunkAttemptCallable.java +++ b/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ChunkAttemptCallable.java @@ -72,6 +72,7 @@ class ChunkAttemptCallable implements Callable implements Callable implements Callable implements Callable implements Callable> currentRetryingFuture) { this.currentCommand = ResumableUploadCommand.QUERY; + progressTracker.onRecovering(lastFailure); + // Per GAX-R7: query uses sensible unary defaults trimmed to the remaining global deadline. long remainingNanos = deadlineNanos == Long.MAX_VALUE @@ -259,6 +266,7 @@ private void handleQuerySuccess( // Normal path: realign buffer to committedOffset, compact and top up. try { buffer.realignTo(committedOffset); + progressTracker.onOffsetReceived(committedOffset); } catch (Throwable e) { failAttempt(attemptFuture, e); return; diff --git a/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ResumableUploadChunkCoordinator.java b/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ResumableUploadChunkCoordinator.java index 09423381b97c..9df50ce4869a 100644 --- a/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ResumableUploadChunkCoordinator.java +++ b/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ResumableUploadChunkCoordinator.java @@ -87,6 +87,7 @@ final class ResumableUploadChunkCoordinator { private final long deadlineNanos; private final ApiCallContext callContext; private final ClientContext clientContext; + private final UploadProgressTracker progressTracker; private final SettableApiFuture result = SettableApiFuture.create(); private volatile @Nullable ApiFuture currentChunkFuture; @@ -98,7 +99,8 @@ final class ResumableUploadChunkCoordinator { ResumableUploadCallSettings settings, long deadlineNanos, ApiCallContext callContext, - ClientContext clientContext) { + ClientContext clientContext, + UploadProgressTracker progressTracker) { this.uploadChunkCallable = checkNotNull(uploadChunkCallable, "uploadChunkCallable must not be null"); this.queryStatusCallable = @@ -109,6 +111,7 @@ final class ResumableUploadChunkCoordinator { this.deadlineNanos = deadlineNanos; this.callContext = checkNotNull(callContext, "callContext must not be null"); this.clientContext = checkNotNull(clientContext, "clientContext must not be null"); + this.progressTracker = checkNotNull(progressTracker, "progressTracker must not be null"); this.buffer = new RewindableStreamBuffer(payload, settings.getChunkSize(), uploadUrl); } @@ -200,6 +203,7 @@ private void transmitSingleChunk(long currentOffset) { chunkRequest, chunkCallContext, command, + progressTracker, deadlineNanos, clientContext.getClock()); @@ -224,6 +228,7 @@ public void onSuccess(ChunkUploadResponse response) { return; } long nextOffset = buffer.getBufferBaseOffset() + buffer.getPayloadLength(); + progressTracker.onChunkUploaded(nextOffset); if (response.getUploadStatus() == ResumableUploadStatus.FINAL) { result.set(response.getResponse()); } else if (buffer.isFinal()) { diff --git a/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ResumableUploadFuture.java b/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ResumableUploadFuture.java index 4282ef89fa33..1868a69d2998 100644 --- a/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ResumableUploadFuture.java +++ b/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ResumableUploadFuture.java @@ -31,6 +31,7 @@ import com.google.api.core.ApiFuture; import com.google.api.core.BetaApi; +import java.util.concurrent.Executor; import org.jspecify.annotations.NullMarked; import org.jspecify.annotations.Nullable; @@ -48,4 +49,18 @@ public interface ResumableUploadFuture extends ApiFuture { /** Returns the upload session URL, or {@code null} if session initiation is in progress. */ @Nullable String getUploadSessionUrl(); + + /** + * Registers a listener to receive progress and state transition notifications for this upload. + * + *

A snapshot of the current upload status is dispatched to the listener immediately upon + * subscription on the provided executor. Subsequent status updates are delivered in order. + * + * @param listener callback listener to receive progress notifications + * @param executor executor on which the listener callbacks are dispatched + */ + void addProgressListener(ResumableUploadProgressListener listener, Executor executor); + + /** Returns the current progress snapshot of the upload session. */ + ResumableUploadProgress getStatus(); } diff --git a/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ResumableUploadFutureImpl.java b/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ResumableUploadFutureImpl.java index ff44d984081d..01fc25185f7e 100644 --- a/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ResumableUploadFutureImpl.java +++ b/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ResumableUploadFutureImpl.java @@ -77,6 +77,7 @@ final class ResumableUploadFutureImpl implements ResumableUploadFutur private final ApiCallContext callContext; private final ClientContext clientContext; private final ScheduledExecutorService executor; + private final UploadProgressTracker progressTracker = new UploadProgressTracker(); private final SettableApiFuture resultFuture = SettableApiFuture.create(); private volatile @Nullable String uploadSessionUrl; @@ -151,6 +152,7 @@ public void onSuccess(ResumableUploadSession session) { return; } uploadSessionUrl = session.getUploadUrl(); + progressTracker.onStarted(uploadSessionUrl); ResumableUploadChunkCoordinator coordinator = new ResumableUploadChunkCoordinator<>( uploadChunkCallable, @@ -160,7 +162,8 @@ public void onSuccess(ResumableUploadSession session) { settings, deadlineNanos, callContext, - clientContext); + clientContext, + progressTracker); ApiFuture uploadFuture; try { uploadFuture = coordinator.start(); @@ -224,6 +227,7 @@ private void succeed(@Nullable ResponseT result) { if (timeout != null) { timeout.cancel(false); } + progressTracker.onFinalized(progressTracker.getStatus().getBytesUploaded()); closePayload(); resultFuture.set(result); } @@ -243,6 +247,7 @@ private void fail(Throwable t) { if (inFlight != null) { inFlight.cancel(true); } + progressTracker.onFailed(t, uploadSessionUrl); closePayload(); resultFuture.setException(t); } @@ -260,6 +265,16 @@ private void closePayload() { return uploadSessionUrl; } + @Override + public void addProgressListener(ResumableUploadProgressListener listener, Executor executor) { + progressTracker.addListener(listener, executor); + } + + @Override + public ResumableUploadProgress getStatus() { + return progressTracker.getStatus(); + } + @Override public void addListener(Runnable listener, Executor executor) { resultFuture.addListener(listener, executor); @@ -283,6 +298,7 @@ public boolean cancel(boolean mayInterruptIfRunning) { if (inFlight != null) { inFlight.cancel(mayInterruptIfRunning); } + progressTracker.onFailed(new CancellationException("Upload was cancelled"), uploadSessionUrl); closePayload(); return cancelled; } diff --git a/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ChunkAttemptCallableTest.java b/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ChunkAttemptCallableTest.java index 365ae6fb02a7..cd6ad47a6575 100644 --- a/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ChunkAttemptCallableTest.java +++ b/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ChunkAttemptCallableTest.java @@ -131,7 +131,8 @@ void call_successfulChunk_setsAttemptFuture() throws Exception { "https://upload.url/test", request, callContext, - ResumableUploadCommand.UPLOAD); + ResumableUploadCommand.UPLOAD, + new UploadProgressTracker()); callable.setRetryingFuture(mockExternalFuture); ChunkUploadResponse callResult = callable.call(); @@ -179,7 +180,8 @@ void call_returnsWithoutBlocking_andPropagatesCancellation() throws Exception { "https://upload.url/test", request, callContext, - ResumableUploadCommand.UPLOAD); + ResumableUploadCommand.UPLOAD, + new UploadProgressTracker()); List listeners = new ArrayList<>(); doAnswer( @@ -253,7 +255,8 @@ void call_perAttemptDeadline_appliesRpcTimeoutToCallContext() throws Exception { "https://upload.url/test", request, callContext, - ResumableUploadCommand.UPLOAD); + ResumableUploadCommand.UPLOAD, + new UploadProgressTracker()); callable.setRetryingFuture(mockExternalFuture); callable.call(); @@ -294,7 +297,8 @@ void call_nonBlockingExecution_callingThreadMakesImmediateProgress() throws Exce "https://upload.url/test", request, callContext, - ResumableUploadCommand.UPLOAD); + ResumableUploadCommand.UPLOAD, + new UploadProgressTracker()); callable.setRetryingFuture(mockExternalFuture); diff --git a/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ResumableUploadCallableImplTest.java b/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ResumableUploadCallableImplTest.java index b035017060b5..86f5512fa088 100644 --- a/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ResumableUploadCallableImplTest.java +++ b/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ResumableUploadCallableImplTest.java @@ -58,15 +58,20 @@ import java.io.InputStream; import java.nio.charset.StandardCharsets; import java.time.Duration; +import java.util.ArrayList; import java.util.Arrays; import java.util.List; import java.util.concurrent.CancellationException; +import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; import org.jspecify.annotations.Nullable; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; @@ -1210,6 +1215,295 @@ void testGlobalTimeout_coversStartSessionTimeout() throws Exception { assertThat(hungStartFuture.isCancelled()).isTrue(); } + @Test + void testProgressListener_snapshotOnSubscribe_postsExactlyOneImmediateUpdate() throws Exception { + SettableApiFuture hungStartFuture = SettableApiFuture.create(); + when(mockStartCallable.futureCall(any(), any())).thenReturn(hungStartFuture); + + ResumableUploadFuture future = + callable.futureCall("resource-path", streamOf("hello"), null); + + List statuses = new CopyOnWriteArrayList<>(); + CountDownLatch latch = new CountDownLatch(1); + future.addProgressListener( + status -> { + statuses.add(status); + latch.countDown(); + }, + executor); + + assertThat(latch.await(5, TimeUnit.SECONDS)).isTrue(); + assertThat(statuses).hasSize(1); + ResumableUploadProgress snapshot = statuses.get(0); + assertThat(snapshot.getState()).isEqualTo(ResumableUploadProgress.State.STARTING); + assertThat(snapshot.getBytesUploaded()).isEqualTo(0L); + assertThat(snapshot.getUploadUrl()).isNull(); + } + + @Test + void testProgressListener_prescribedStateTransitions() throws Exception { + SettableApiFuture startFuture = SettableApiFuture.create(); + when(mockStartCallable.futureCall(any(), any())).thenReturn(startFuture); + + // 20 bytes with chunkSize = 8 -> 3 chunks: [0..8), [8..16), [16..20) + // Chunk 1 succeeds -> [0..8) + // Chunk 2 fails with Cat-2 400 + // Query succeeds -> committed offset = 8 + // Chunk 2 resend succeeds -> [8..16) + // Chunk 3 succeeds and finalizes -> [16..20) + when(mockChunkCallable.futureCall(any(ChunkUploadRequest.class), any())) + .thenReturn( + ApiFutures.immediateFuture( + ChunkUploadResponse.create(ResumableUploadStatus.ACTIVE, null))) + .thenReturn( + ApiFutures.immediateFailedFuture( + createApiException(400, StatusCode.Code.INVALID_ARGUMENT))) + .thenReturn( + ApiFutures.immediateFuture( + ChunkUploadResponse.create(ResumableUploadStatus.ACTIVE, null))) + .thenReturn( + ApiFutures.immediateFuture( + ChunkUploadResponse.create(ResumableUploadStatus.FINAL, "done"))); + + when(mockQueryCallable.futureCall(any(QueryStatusRequest.class), any())) + .thenReturn( + ApiFutures.immediateFuture( + QueryStatusResponse.newBuilder() + .setCommittedOffset(8L) + .setUploadStatus(ResumableUploadStatus.ACTIVE) + .build())); + + List receivedStatuses = new CopyOnWriteArrayList<>(); + CountDownLatch finalizedLatch = new CountDownLatch(1); + + ResumableUploadFuture future = + callable.futureCall("resource-path", streamOf("01234567890123456789"), null); + future.addProgressListener( + status -> { + receivedStatuses.add(status); + if (status.getState() == ResumableUploadProgress.State.FINALIZED) { + finalizedLatch.countDown(); + } + }, + executor); + + startFuture.set( + ResumableUploadSession.newBuilder() + .setUploadUrl("https://upload.url/progress-transitions") + .build()); + + assertThat(future.get()).isEqualTo("done"); + assertThat(finalizedLatch.await(5, TimeUnit.SECONDS)).isTrue(); + + List states = new ArrayList<>(); + for (ResumableUploadProgress s : receivedStatuses) { + states.add(s.getState()); + } + + assertThat(states) + .containsAtLeast( + ResumableUploadProgress.State.STARTED, + ResumableUploadProgress.State.UPLOADING, + ResumableUploadProgress.State.RECOVERING, + ResumableUploadProgress.State.OFFSET_RECEIVED, + ResumableUploadProgress.State.FINALIZED) + .inOrder(); + + long lastBytes = 0; + for (ResumableUploadProgress s : receivedStatuses) { + assertThat(s.getBytesUploaded()).isAtLeast(lastBytes); + lastBytes = s.getBytesUploaded(); + if (s.getState() != ResumableUploadProgress.State.STARTING) { + assertThat(s.getUploadUrl()).isEqualTo("https://upload.url/progress-transitions"); + } + } + assertThat(lastBytes).isEqualTo(20L); + } + + @Test + void testProgressListener_throwingListener_doesNotBreakUpload() throws Exception { + stubStartSession("https://upload.url/throwing-listener"); + when(mockChunkCallable.futureCall(any(ChunkUploadRequest.class), any())) + .thenReturn( + ApiFutures.immediateFuture( + ChunkUploadResponse.create(ResumableUploadStatus.FINAL, "ok"))); + + ResumableUploadFuture future = + callable.futureCall("resource-path", streamOf("hello"), null); + + future.addProgressListener( + status -> { + throw new RuntimeException("boom from listener"); + }, + executor); + + assertThat(future.get()).isEqualTo("ok"); + assertThat(future.isDone()).isTrue(); + } + + @Test + void testProgressListener_subscribingAfterCompletion_yieldsOneTerminalSnapshot() + throws Exception { + stubStartSession("https://upload.url/post-completion"); + when(mockChunkCallable.futureCall(any(ChunkUploadRequest.class), any())) + .thenReturn( + ApiFutures.immediateFuture( + ChunkUploadResponse.create(ResumableUploadStatus.FINAL, "completed-ok"))); + + ResumableUploadFuture future = + callable.futureCall("resource-path", streamOf("hello"), null); + assertThat(future.get()).isEqualTo("completed-ok"); + + List postStatuses = new CopyOnWriteArrayList<>(); + CountDownLatch latch = new CountDownLatch(1); + future.addProgressListener( + status -> { + postStatuses.add(status); + latch.countDown(); + }, + executor); + + assertThat(latch.await(5, TimeUnit.SECONDS)).isTrue(); + executor.submit(() -> {}).get(5, TimeUnit.SECONDS); + assertThat(postStatuses).hasSize(1); + ResumableUploadProgress snapshot = postStatuses.get(0); + assertThat(snapshot.getState()).isEqualTo(ResumableUploadProgress.State.FINALIZED); + assertThat(snapshot.getUploadUrl()).isEqualTo("https://upload.url/post-completion"); + assertThat(snapshot.getBytesUploaded()).isEqualTo(5L); + } + + @Test + void testProgressListener_subscribingAfterFailure_yieldsOneTerminalFailedSnapshot() + throws Exception { + when(mockStartCallable.futureCall(any(), any())) + .thenReturn( + ApiFutures.immediateFailedFuture( + createApiException(401, StatusCode.Code.UNAUTHENTICATED))); + + ResumableUploadFuture future = + callable.futureCall("resource-path", streamOf("hello"), null); + assertThrows(ExecutionException.class, future::get); + + List postStatuses = new CopyOnWriteArrayList<>(); + CountDownLatch latch = new CountDownLatch(1); + future.addProgressListener( + status -> { + postStatuses.add(status); + latch.countDown(); + }, + executor); + + assertThat(latch.await(5, TimeUnit.SECONDS)).isTrue(); + executor.submit(() -> {}).get(5, TimeUnit.SECONDS); + assertThat(postStatuses).hasSize(1); + ResumableUploadProgress snapshot = postStatuses.get(0); + assertThat(snapshot.getState()).isEqualTo(ResumableUploadProgress.State.FAILED); + assertThat(snapshot.getException()).isInstanceOf(ApiException.class); + } + + @Test + void testProgressListener_futureCancelFromInsideListenerBody_worksWithoutDeadlock() + throws Exception { + stubStartSession("https://upload.url/cancel-inside-listener"); + SettableApiFuture> hungChunk = SettableApiFuture.create(); + when(mockChunkCallable.futureCall(any(ChunkUploadRequest.class), any())).thenReturn(hungChunk); + + ResumableUploadFuture future = + callable.futureCall("resource-path", streamOf("hello"), null); + + CountDownLatch cancelAttemptedLatch = new CountDownLatch(1); + future.addProgressListener( + status -> { + if (status.getState() == ResumableUploadProgress.State.STARTED) { + future.cancel(true); + cancelAttemptedLatch.countDown(); + } + }, + executor); + + assertThat(cancelAttemptedLatch.await(5, TimeUnit.SECONDS)).isTrue(); + assertThat(future.isCancelled()).isTrue(); + assertThat(hungChunk.isCancelled()).isTrue(); + } + + @Test + void testProgressListener_getStatus_reflectsCurrentState() throws Exception { + stubStartSession("https://upload.url/get-status"); + when(mockChunkCallable.futureCall(any(ChunkUploadRequest.class), any())) + .thenReturn( + ApiFutures.immediateFuture( + ChunkUploadResponse.create(ResumableUploadStatus.FINAL, "ok"))); + + ResumableUploadFuture future = + callable.futureCall("resource-path", streamOf("hello"), null); + + assertThat(future.get()).isEqualTo("ok"); + ResumableUploadProgress status = future.getStatus(); + assertThat(status.getState()).isEqualTo(ResumableUploadProgress.State.FINALIZED); + assertThat(status.getUploadUrl()).isEqualTo("https://upload.url/get-status"); + assertThat(status.getBytesUploaded()).isEqualTo(5L); + } + + @Test + void testProgressListener_orderingUnderConcurrency_pinsSequentialExecutor() throws Exception { + ExecutorService multiThreadedExecutor = Executors.newFixedThreadPool(8); + try { + stubStartSession("https://upload.url/concurrency-order"); + ChunkUploadResponse activeResponse = + ChunkUploadResponse.create(ResumableUploadStatus.ACTIVE, null); + when(mockChunkCallable.futureCall(any(ChunkUploadRequest.class), any())) + .thenReturn(ApiFutures.immediateFuture(activeResponse)) + .thenReturn(ApiFutures.immediateFuture(activeResponse)) + .thenReturn(ApiFutures.immediateFuture(activeResponse)) + .thenReturn(ApiFutures.immediateFuture(activeResponse)) + .thenReturn( + ApiFutures.immediateFuture( + ChunkUploadResponse.create(ResumableUploadStatus.FINAL, "finished"))); + + List events = new CopyOnWriteArrayList<>(); + AtomicInteger concurrentExecutions = new AtomicInteger(0); + AtomicBoolean concurrencyDetected = new AtomicBoolean(false); + CountDownLatch finalizedLatch = new CountDownLatch(1); + + ResumableUploadFuture future = + callable.futureCall( + "resource-path", streamOf("0123456789012345678901234567890123456789"), null); + + future.addProgressListener( + status -> { + int inProgress = concurrentExecutions.incrementAndGet(); + if (inProgress > 1) { + concurrencyDetected.set(true); + } + try { + Thread.sleep(10); + events.add(status); + if (status.getState() == ResumableUploadProgress.State.FINALIZED) { + finalizedLatch.countDown(); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } finally { + concurrentExecutions.decrementAndGet(); + } + }, + multiThreadedExecutor); + + assertThat(future.get(10, TimeUnit.SECONDS)).isEqualTo("finished"); + assertThat(finalizedLatch.await(5, TimeUnit.SECONDS)).isTrue(); + assertThat(concurrencyDetected.get()).isFalse(); + + long lastBytes = 0; + for (ResumableUploadProgress s : events) { + assertThat(s.getBytesUploaded()).isAtLeast(lastBytes); + lastBytes = s.getBytesUploaded(); + } + assertThat(lastBytes).isEqualTo(40L); + } finally { + multiThreadedExecutor.shutdownNow(); + } + } + private static class HttpStatusStatusCode implements StatusCode { private final int httpStatus; private final StatusCode.Code code; @@ -1283,7 +1577,7 @@ private static void assertChunk( } private static class TrackableStream extends ByteArrayInputStream { - int closeCount = 0; + volatile int closeCount = 0; TrackableStream(String content) { super(content.getBytes(StandardCharsets.UTF_8)); diff --git a/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ResumableUploadChunkCoordinatorTest.java b/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ResumableUploadChunkCoordinatorTest.java index a527350e9775..cc0ba1672e33 100644 --- a/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ResumableUploadChunkCoordinatorTest.java +++ b/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ResumableUploadChunkCoordinatorTest.java @@ -82,6 +82,9 @@ void start_completesFutureOnFinalChunk() throws Exception { when(mockChunkCallable.futureCall(any(ChunkUploadRequest.class), any())) .thenReturn(chunkFuture); + UploadProgressTracker progressTracker = new UploadProgressTracker(); + progressTracker.onStarted("https://upload.url/coordinator-test"); + ResumableUploadChunkCoordinator coordinator = new ResumableUploadChunkCoordinator<>( mockChunkCallable, @@ -91,7 +94,8 @@ void start_completesFutureOnFinalChunk() throws Exception { settings, Long.MAX_VALUE, callContext, - clientContext); + clientContext, + progressTracker); ApiFuture result = coordinator.start(); assertThat(result.isDone()).isFalse(); @@ -100,6 +104,7 @@ void start_completesFutureOnFinalChunk() throws Exception { ChunkUploadResponse.create(ResumableUploadStatus.FINAL, "done")); assertThat(result.get()).isEqualTo("done"); + assertThat(progressTracker.getStatus().getBytesUploaded()).isEqualTo(64L); } @Test @@ -117,7 +122,8 @@ void cancel_propagatesToInFlightChunkFuture() { settings, Long.MAX_VALUE, callContext, - clientContext); + clientContext, + new UploadProgressTracker()); ApiFuture result = coordinator.start(); assertThat(result.cancel(true)).isTrue(); diff --git a/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ResumableUploadFutureImplTest.java b/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ResumableUploadFutureImplTest.java index 1ae303c1a67e..f5e15f283357 100644 --- a/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ResumableUploadFutureImplTest.java +++ b/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ResumableUploadFutureImplTest.java @@ -46,6 +46,9 @@ import com.google.common.util.concurrent.MoreExecutors; import java.io.ByteArrayInputStream; import java.io.IOException; +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; import java.util.concurrent.CountDownLatch; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; @@ -172,6 +175,17 @@ public void close() throws IOException { callContext, clientContext); + List terminalStatuses = + Collections.synchronizedList(new ArrayList<>()); + handle.addProgressListener( + status -> { + if (status.getState() == ResumableUploadProgress.State.FINALIZED + || status.getState() == ResumableUploadProgress.State.FAILED) { + terminalStatuses.add(status); + } + }, + MoreExecutors.directExecutor()); + AtomicInteger completionListenerCount = new AtomicInteger(0); handle.addListener(completionListenerCount::incrementAndGet, MoreExecutors.directExecutor()); @@ -209,5 +223,6 @@ public void close() throws IOException { assertThat(handle.isDone()).isTrue(); assertThat(closeCount.get()).isEqualTo(1); assertThat(completionListenerCount.get()).isEqualTo(1); + assertThat(terminalStatuses).hasSize(1); } }