Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@
*/
package com.google.api.gax.rpc;

import static com.google.common.base.MoreObjects.firstNonNull;
import static com.google.common.base.Preconditions.checkArgument;
import static com.google.common.base.Preconditions.checkNotNull;

Expand All @@ -46,10 +47,12 @@
import com.google.errorprone.annotations.concurrent.GuardedBy;
import java.io.IOException;
import java.io.InputStream;
import java.time.Duration;
import java.util.concurrent.CancellationException;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.Executor;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import org.jspecify.annotations.NullMarked;
Expand All @@ -64,6 +67,8 @@
@NullMarked
final class ResumableUploadFutureImpl<ResponseT> implements ResumableUploadFuture<ResponseT> {

private static final Duration DEFAULT_GLOBAL_TIMEOUT = Duration.ofMinutes(15);

private final Object lock = new Object();

private final ApiFuture<ResumableUploadSession> startFuture;
Expand All @@ -80,7 +85,8 @@

private volatile @Nullable String uploadSessionUrl;

// Tracks the current operation's Future (start, chunk upload) to propagate cancellation.
// Tracks the current operation's Future (start, chunk upload) to propagate cancellation, or
// null once a terminal transition (succeed, fail, cancel) has claimed completion.
@GuardedBy("lock")
private @Nullable ApiFuture<?> inFlightFuture;

Expand Down Expand Up @@ -140,6 +146,10 @@
}

private void start() {
Duration timeout = firstNonNull(options.getGlobalTimeout(), DEFAULT_GLOBAL_TIMEOUT);
ScheduledFuture<?> timeoutFuture =
executor.schedule(this::onTimeout, timeout.toMillis(), TimeUnit.MILLISECONDS);
resultFuture.addListener(() -> timeoutFuture.cancel(false), MoreExecutors.directExecutor());
ApiFutures.addCallback(
startFuture,
new ApiFutureCallback<ResumableUploadSession>() {
Expand All @@ -158,7 +168,7 @@
executor);
ApiFuture<ResponseT> uploadFuture = coordinator.getFuture();
synchronized (lock) {
if (resultFuture.isDone()) {
if (inFlightFuture == null) {
return;
}
uploadSessionUrl = sessionUrl;
Expand Down Expand Up @@ -195,18 +205,38 @@
MoreExecutors.directExecutor());
}

private void onTimeout() {
String sessionUrl = uploadSessionUrl;
String message;
if (sessionUrl != null) {
message = "Resumable upload timed out for session: " + sessionUrl;
} else {
message = "Resumable upload timed out before session initiation completed";
}
fail(ApiExceptionFactory.createException(message, null, TIMEOUT_STATUS_CODE, false));
Comment thread
blakeli0 marked this conversation as resolved.
}

private void succeed(@Nullable ResponseT result) {
synchronized (lock) {
if (inFlightFuture == null) {
return;
}
inFlightFuture = null;
}
closePayload();
resultFuture.set(result);
}

private void fail(Throwable t) {
ApiFuture<?> inFlight;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

There could be a race condition per AI:

  1. onTimeout() → fail() enters synchronized (lock), reads inFlight = startFuture, sets this.inFlightFuture = null, and exits synchronized (lock). At this moment, resultFuture.isDone() is still false.
  2. startFuture.onSuccess() enters synchronized (lock), checks if (resultFuture.isDone()) (which is still false!), sets inFlightFuture = uploadFuture, and exits synchronized (lock).
  3. fail() calls closePayload() and resultFuture.setException(DeadlineExceededException).
  4. startFuture.onSuccess() calls coordinator.start(), launching the background chunk coordinator—and uploadFuture is never cancelled because fail() already cleared inFlightFuture before step 2.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

All terminal transitions (succeed, fail, cancel) now return early if inFlightFuture == null; otherwise they null it and complete resultFuture. So only the first one "wins".

startFuture.onSuccess() also checks inFlightFuture == null (instead of resultFuture.isDone()) before installing uploadFuture, so it can't start the coordinator after fail(). If onSuccess() gets there first, fail() now cancels uploadFuture.

synchronized (lock) {
if (inFlightFuture == null) {
return;
}
inFlight = inFlightFuture;
inFlightFuture = null;
}
inFlight.cancel(true);
closePayload();
resultFuture.setException(t);
}
Expand Down Expand Up @@ -234,13 +264,14 @@
boolean cancelled;
ApiFuture<?> inFlight;
synchronized (lock) {
if (inFlightFuture == null) {
return false;
}
cancelled = resultFuture.cancel(mayInterruptIfRunning);
inFlight = this.inFlightFuture;
this.inFlightFuture = null;
}
if (inFlight != null) {
inFlight.cancel(mayInterruptIfRunning);
inFlight = inFlightFuture;
inFlightFuture = null;
}
inFlight.cancel(mayInterruptIfRunning);
closePayload();
return cancelled;
}
Expand All @@ -265,4 +296,17 @@
throws InterruptedException, ExecutionException, TimeoutException {
return resultFuture.get(timeout, unit);
}

private static final StatusCode TIMEOUT_STATUS_CODE =
new StatusCode() {
@Override
public StatusCode.Code getCode() {
return StatusCode.Code.DEADLINE_EXCEEDED;
}

@Override
public @Nullable Object getTransportCode() {

Check failure on line 308 in sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ResumableUploadFutureImpl.java

View check run for this annotation

SonarQubeCloud / [gapic-generator-java-root] SonarCloud Code Analysis

Fix the incompatibility of the annotation @Nullable to honor @NullMarked at class level of the overridden method.

See more on https://sonarcloud.io/project/issues?id=googleapis_google-cloud-java_showcase&issues=AaC91aRmz6wgCEivc2NM&open=AaC91aRmz6wgCEivc2NM&pullRequest=14425
return null;
}
};
}
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,8 @@
import static com.google.common.truth.Truth.assertThat;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.lenient;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.times;
Expand All @@ -56,10 +58,13 @@
import java.io.IOException;
import java.io.InputStream;
import java.nio.charset.StandardCharsets;
import java.time.Duration;
import java.util.List;
import java.util.concurrent.CancellationException;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import org.jspecify.annotations.Nullable;
Expand All @@ -79,6 +84,7 @@

private ResumableUploadOptions defaultSettings;
private FakeCallContext callContext;
private ClientContext clientContext;
private ResumableUploadCallableImpl<String, String> callable;

@BeforeEach
Expand All @@ -95,8 +101,7 @@

defaultSettings = ResumableUploadOptions.newBuilder().setChunkSize(8).build();
callContext = FakeCallContext.createDefault();
ClientContext clientContext =
ClientContext.newBuilder().setDefaultCallContext(callContext).build();
clientContext = ClientContext.newBuilder().setDefaultCallContext(callContext).build();
callable = new ResumableUploadCallableImpl<>(mockClient, defaultSettings, clientContext);
}

Expand Down Expand Up @@ -860,6 +865,210 @@
verify(mockChunkCallable, times(10)).futureCall(any(), any());
}

@Test
void testGlobalTimeout_firesAndFailsSessionWithDeadlineExceeded() throws Exception {

Check warning on line 869 in sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ResumableUploadCallableImplTest.java

View check run for this annotation

SonarQubeCloud / [gapic-generator-java-root] SonarCloud Code Analysis

Remove the declaration of thrown exception 'java.lang.Exception', as it cannot be thrown from method's body.

See more on https://sonarcloud.io/project/issues?id=googleapis_google-cloud-java_showcase&issues=AaCyXlS_NxZs2PJzo0AY&open=AaCyXlS_NxZs2PJzo0AY&pullRequest=14425
stubStartSession("https://upload.url/timeout-fire");
SettableApiFuture<ChunkUploadResponse<String>> hungChunk = SettableApiFuture.create();
when(mockChunkCallable.futureCall(any(ChunkUploadRequest.class), any())).thenReturn(hungChunk);

ResumableUploadOptions timeoutSettings =
defaultSettings.toBuilder().setGlobalTimeout(Duration.ofMillis(100)).build();

ResumableUploadFuture<String> future =
callable.futureCall("resource-path", streamOf("hello"), timeoutSettings);

ExecutionException exception =
assertThrows(ExecutionException.class, () -> future.get(5, TimeUnit.SECONDS));
assertThat(exception.getCause()).isInstanceOf(DeadlineExceededException.class);
DeadlineExceededException cause = (DeadlineExceededException) exception.getCause();
assertThat(cause.getStatusCode().getCode()).isEqualTo(StatusCode.Code.DEADLINE_EXCEEDED);
assertThat(cause.getMessage()).contains("https://upload.url/timeout-fire");
assertThat(future.isDone()).isTrue();
assertThat(future.isCancelled()).isFalse();
}

@Test
void testGlobalTimeout_cancelledCleanlyOnSuccess() throws Exception {
stubStartSession("https://upload.url/timeout-success");
when(mockChunkCallable.futureCall(any(ChunkUploadRequest.class), any()))
.thenReturn(
ApiFutures.immediateFuture(
ChunkUploadResponse.create(ResumableUploadStatus.FINAL, "ok")));

ScheduledExecutorService mockExecutor = mock(ScheduledExecutorService.class);
ScheduledFuture<?> mockScheduledFuture = mock(ScheduledFuture.class);
when(mockExecutor.schedule(any(Runnable.class), anyLong(), any()))
.thenAnswer(inv -> mockScheduledFuture);

ClientContext customClientContext = clientContext.toBuilder().setExecutor(mockExecutor).build();
ResumableUploadCallableImpl<String, String> customCallable =
new ResumableUploadCallableImpl<>(mockClient, defaultSettings, customClientContext);

ResumableUploadOptions timeoutSettings =
defaultSettings.toBuilder().setGlobalTimeout(Duration.ofSeconds(60)).build();

ResumableUploadFuture<String> future =
customCallable.futureCall("resource-path", streamOf("hello"), timeoutSettings);

assertThat(future.get()).isEqualTo("ok");
verify(mockScheduledFuture).cancel(false);
}

@Test
void testGlobalTimeout_cancelledCleanlyOnFailure() throws Exception {

Check warning on line 918 in sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ResumableUploadCallableImplTest.java

View check run for this annotation

SonarQubeCloud / [gapic-generator-java-root] SonarCloud Code Analysis

Remove the declaration of thrown exception 'java.lang.Exception', as it cannot be thrown from method's body.

See more on https://sonarcloud.io/project/issues?id=googleapis_google-cloud-java_showcase&issues=AaCyXlS_NxZs2PJzo0AZ&open=AaCyXlS_NxZs2PJzo0AZ&pullRequest=14425
when(mockStartCallable.futureCall(any(), any()))
.thenReturn(
ApiFutures.immediateFailedFuture(
createApiException(401, StatusCode.Code.UNAUTHENTICATED)));

ScheduledExecutorService mockExecutor = mock(ScheduledExecutorService.class);
ScheduledFuture<?> mockScheduledFuture = mock(ScheduledFuture.class);
when(mockExecutor.schedule(any(Runnable.class), anyLong(), any()))
.thenAnswer(inv -> mockScheduledFuture);

ClientContext customClientContext = clientContext.toBuilder().setExecutor(mockExecutor).build();
ResumableUploadCallableImpl<String, String> customCallable =
new ResumableUploadCallableImpl<>(mockClient, defaultSettings, customClientContext);

ResumableUploadOptions timeoutSettings =
defaultSettings.toBuilder().setGlobalTimeout(Duration.ofSeconds(60)).build();

ResumableUploadFuture<String> future =
customCallable.futureCall("resource-path", streamOf("hello"), timeoutSettings);

assertThrows(ExecutionException.class, future::get);
verify(mockScheduledFuture).cancel(false);
}

@Test
void testGlobalTimeout_cancelledCleanlyOnUserCancel() throws Exception {

Check warning on line 944 in sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ResumableUploadCallableImplTest.java

View check run for this annotation

SonarQubeCloud / [gapic-generator-java-root] SonarCloud Code Analysis

Remove the declaration of thrown exception 'java.lang.Exception', as it cannot be thrown from method's body.

See more on https://sonarcloud.io/project/issues?id=googleapis_google-cloud-java_showcase&issues=AaCyXlS_NxZs2PJzo0Aa&open=AaCyXlS_NxZs2PJzo0Aa&pullRequest=14425
stubStartSession("https://upload.url/timeout-cancel");
SettableApiFuture<ChunkUploadResponse<String>> hungChunk = SettableApiFuture.create();
when(mockChunkCallable.futureCall(any(ChunkUploadRequest.class), any())).thenReturn(hungChunk);

ScheduledExecutorService mockExecutor = mock(ScheduledExecutorService.class);
ScheduledFuture<?> mockScheduledFuture = mock(ScheduledFuture.class);
when(mockExecutor.schedule(any(Runnable.class), anyLong(), any()))
.thenAnswer(inv -> mockScheduledFuture);

ClientContext customClientContext = clientContext.toBuilder().setExecutor(mockExecutor).build();
ResumableUploadCallableImpl<String, String> customCallable =
new ResumableUploadCallableImpl<>(mockClient, defaultSettings, customClientContext);

ResumableUploadOptions timeoutSettings =
defaultSettings.toBuilder().setGlobalTimeout(Duration.ofSeconds(60)).build();

ResumableUploadFuture<String> future =
customCallable.futureCall("resource-path", streamOf("hello"), timeoutSettings);

assertThat(future.cancel(true)).isTrue();
verify(mockScheduledFuture).cancel(false);
}

@Test
void testGlobalTimeout_usesDefaultWhenUnset() throws Exception {
stubStartSession("https://upload.url/default-timeout");
when(mockChunkCallable.futureCall(any(ChunkUploadRequest.class), any()))
.thenReturn(
ApiFutures.immediateFuture(
ChunkUploadResponse.create(ResumableUploadStatus.FINAL, "ok")));

ScheduledExecutorService mockExecutor = mock(ScheduledExecutorService.class);
ScheduledFuture<?> mockScheduledFuture = mock(ScheduledFuture.class);
when(mockExecutor.schedule(any(Runnable.class), anyLong(), any()))
.thenAnswer(inv -> mockScheduledFuture);

ClientContext customClientContext = clientContext.toBuilder().setExecutor(mockExecutor).build();

// 1. Unset on both stub and per-request -> falls back to GAX default (15m)
ResumableUploadCallableImpl<String, String> unsetStubCallable =
new ResumableUploadCallableImpl<>(mockClient, defaultSettings, customClientContext);
assertThat(unsetStubCallable.futureCall("resource-path", streamOf("hello"), null).get())
.isEqualTo("ok");
verify(mockExecutor)
.schedule(
any(Runnable.class), eq(Duration.ofMinutes(15).toMillis()), eq(TimeUnit.MILLISECONDS));

// 2. Stub-level timeout (30m, e.g. from generator/client settings) + null per-request -> 30m
ResumableUploadOptions stubWith30m =
defaultSettings.toBuilder().setGlobalTimeout(Duration.ofMinutes(30)).build();
ResumableUploadCallableImpl<String, String> configuredStubCallable =
new ResumableUploadCallableImpl<>(mockClient, stubWith30m, customClientContext);
assertThat(configuredStubCallable.futureCall("resource-path", streamOf("hello"), null).get())
.isEqualTo("ok");
verify(mockExecutor)
.schedule(
any(Runnable.class), eq(Duration.ofMinutes(30).toMillis()), eq(TimeUnit.MILLISECONDS));

// 3. Stub-level timeout (30m) + per-request with only chunkSize set -> preserves 30m
ResumableUploadOptions perRequestChunkSizeOnly =
ResumableUploadOptions.newBuilder().setChunkSize(16).build();
assertThat(
configuredStubCallable
.futureCall("resource-path", streamOf("hello"), perRequestChunkSizeOnly)
.get())
.isEqualTo("ok");
verify(mockExecutor, times(2))
.schedule(
any(Runnable.class), eq(Duration.ofMinutes(30).toMillis()), eq(TimeUnit.MILLISECONDS));

// 4. Stub-level timeout (30m) + per-request globalTimeout (5m) -> per-request wins (5m)
ResumableUploadOptions perRequestWith5m =
ResumableUploadOptions.newBuilder().setGlobalTimeout(Duration.ofMinutes(5)).build();
assertThat(
configuredStubCallable
.futureCall("resource-path", streamOf("hello"), perRequestWith5m)
.get())
.isEqualTo("ok");
verify(mockExecutor)
.schedule(
any(Runnable.class), eq(Duration.ofMinutes(5).toMillis()), eq(TimeUnit.MILLISECONDS));
}

@Test
void testGlobalTimeout_timeoutWhileAttemptInFlight_cancelsInFlightFutureAndDoesNotCorruptBuffer()
throws Exception {

Check warning on line 1030 in sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ResumableUploadCallableImplTest.java

View check run for this annotation

SonarQubeCloud / [gapic-generator-java-root] SonarCloud Code Analysis

Remove the declaration of thrown exception 'java.lang.Exception', as it cannot be thrown from method's body.

See more on https://sonarcloud.io/project/issues?id=googleapis_google-cloud-java_showcase&issues=AaCyXlS_NxZs2PJzo0Ab&open=AaCyXlS_NxZs2PJzo0Ab&pullRequest=14425
stubStartSession("https://upload.url/in-flight-timeout");
SettableApiFuture<ChunkUploadResponse<String>> inFlightFuture = SettableApiFuture.create();
when(mockChunkCallable.futureCall(any(ChunkUploadRequest.class), any()))
.thenReturn(inFlightFuture);

ResumableUploadOptions timeoutSettings =
defaultSettings.toBuilder().setGlobalTimeout(Duration.ofMillis(80)).build();

TrackableStream stream = new TrackableStream("01234567890123456789");
ResumableUploadFuture<String> future =
callable.futureCall("resource-path", stream, timeoutSettings);

ExecutionException exception =
assertThrows(ExecutionException.class, () -> future.get(5, TimeUnit.SECONDS));
assertThat(exception.getCause()).isInstanceOf(DeadlineExceededException.class);
// In-flight attempt future must be cancelled
assertThat(inFlightFuture.isCancelled()).isTrue();

// Stream should have been read only up to the first chunk (chunkSize = 8), not refilled or
// advanced
assertThat(stream.totalBytesRead).isEqualTo(8);
}

@Test
void testGlobalTimeout_coversStartSessionTimeout() throws Exception {

Check warning on line 1055 in sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ResumableUploadCallableImplTest.java

View check run for this annotation

SonarQubeCloud / [gapic-generator-java-root] SonarCloud Code Analysis

Remove the declaration of thrown exception 'java.lang.Exception', as it cannot be thrown from method's body.

See more on https://sonarcloud.io/project/issues?id=googleapis_google-cloud-java_showcase&issues=AaCyXlS_NxZs2PJzo0Ac&open=AaCyXlS_NxZs2PJzo0Ac&pullRequest=14425
SettableApiFuture<ResumableUploadSession> hungStartFuture = SettableApiFuture.create();
when(mockStartCallable.futureCall(any(), any())).thenReturn(hungStartFuture);

ResumableUploadOptions timeoutSettings =
defaultSettings.toBuilder().setGlobalTimeout(Duration.ofMillis(80)).build();

ResumableUploadFuture<String> future =
callable.futureCall("resource-path", streamOf("hello"), timeoutSettings);

ExecutionException exception =
assertThrows(ExecutionException.class, () -> future.get(5, TimeUnit.SECONDS));
assertThat(exception.getCause()).isInstanceOf(DeadlineExceededException.class);
assertThat(exception.getCause().getMessage()).contains("before session initiation completed");
assertThat(hungStartFuture.isCancelled()).isTrue();
}

private static class HttpStatusStatusCode implements StatusCode {
private final int httpStatus;
private final StatusCode.Code code;
Expand Down
Loading