Skip to content

Commit 864d8fd

Browse files
committed
feat(gax): enforce global timeout for resumable uploads
Enforces ResumableUploadCallSettings.getGlobalTimeout() (defaulting to 15 minutes when unset) in ResumableUploadFutureImpl across the entire upload session lifecycle, including session initiation and chunk streaming. Cancels in-flight RPCs, closes the payload stream, and completes the future with DeadlineExceededException when the deadline is exceeded.
1 parent ca7da1b commit 864d8fd

3 files changed

Lines changed: 264 additions & 8 deletions

File tree

‎sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ResumableUploadCallableImpl.java‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -120,7 +120,7 @@ public ResumableUploadFuture<ResponseT> futureCall(
120120
retryingQueryCallable,
121121
payload,
122122
effectiveSettings,
123-
clientContext.getDefaultCallContext());
123+
clientContext);
124124
}
125125

126126
@Override

‎sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ResumableUploadFutureImpl.java‎

Lines changed: 52 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@
2929
*/
3030
package com.google.api.gax.rpc;
3131

32+
import static com.google.common.base.MoreObjects.firstNonNull;
3233
import static com.google.common.base.Preconditions.checkArgument;
3334
import static com.google.common.base.Preconditions.checkNotNull;
3435

@@ -45,9 +46,12 @@
4546
import com.google.errorprone.annotations.concurrent.GuardedBy;
4647
import java.io.IOException;
4748
import java.io.InputStream;
49+
import java.time.Duration;
4850
import java.util.concurrent.CancellationException;
4951
import java.util.concurrent.ExecutionException;
5052
import java.util.concurrent.Executor;
53+
import java.util.concurrent.ScheduledExecutorService;
54+
import java.util.concurrent.ScheduledFuture;
5155
import java.util.concurrent.TimeUnit;
5256
import java.util.concurrent.TimeoutException;
5357
import org.jspecify.annotations.NullMarked;
@@ -62,6 +66,8 @@
6266
@NullMarked
6367
final class ResumableUploadFutureImpl<ResponseT> implements ResumableUploadFuture<ResponseT> {
6468

69+
private static final Duration DEFAULT_GLOBAL_TIMEOUT = Duration.ofMinutes(15);
70+
6571
private final Object lock = new Object();
6672

6773
private final ApiFuture<ResumableUploadSession> startFuture;
@@ -72,6 +78,7 @@ final class ResumableUploadFutureImpl<ResponseT> implements ResumableUploadFutur
7278
private final InputStream payload;
7379
private final ResumableUploadCallSettings settings;
7480
private final ApiCallContext callContext;
81+
private final ScheduledExecutorService executor;
7582
private final SettableApiFuture<ResponseT> resultFuture = SettableApiFuture.create();
7683

7784
private volatile @Nullable String uploadSessionUrl;
@@ -93,10 +100,15 @@ static <ResponseT> ResumableUploadFutureImpl<ResponseT> create(
93100
UnaryCallable<QueryStatusRequest, QueryStatusResponse<ResponseT>> queryStatusCallable,
94101
InputStream payload,
95102
ResumableUploadCallSettings settings,
96-
ApiCallContext callContext) {
103+
ClientContext clientContext) {
97104
ResumableUploadFutureImpl<ResponseT> future =
98105
new ResumableUploadFutureImpl<>(
99-
startFuture, uploadChunkCallable, queryStatusCallable, payload, settings, callContext);
106+
startFuture,
107+
uploadChunkCallable,
108+
queryStatusCallable,
109+
payload,
110+
settings,
111+
clientContext);
100112
try {
101113
future.start();
102114
} catch (Throwable t) {
@@ -111,7 +123,7 @@ private ResumableUploadFutureImpl(
111123
UnaryCallable<QueryStatusRequest, QueryStatusResponse<ResponseT>> queryStatusCallable,
112124
InputStream payload,
113125
ResumableUploadCallSettings settings,
114-
ApiCallContext callContext) {
126+
ClientContext clientContext) {
115127
this.startFuture = checkNotNull(startFuture, "startFuture must not be null");
116128
this.uploadChunkCallable =
117129
checkNotNull(uploadChunkCallable, "uploadChunkCallable must not be null");
@@ -120,11 +132,17 @@ private ResumableUploadFutureImpl(
120132
this.payload = checkNotNull(payload, "payload must not be null");
121133
this.settings = checkNotNull(settings, "settings must not be null");
122134
checkArgument(settings.getChunkSize() > 0, "chunkSize must be > 0");
123-
this.callContext = checkNotNull(callContext, "callContext must not be null");
135+
checkNotNull(clientContext, "clientContext must not be null");
136+
this.callContext = clientContext.getDefaultCallContext();
137+
this.executor = checkNotNull(clientContext.getExecutor(), "executor must not be null");
124138
this.inFlightFuture = startFuture;
125139
}
126140

127141
private void start() {
142+
Duration timeout = firstNonNull(settings.getGlobalTimeout(), DEFAULT_GLOBAL_TIMEOUT);
143+
ScheduledFuture<?> timeoutFuture =
144+
executor.schedule(this::onTimeout, timeout.toMillis(), TimeUnit.MILLISECONDS);
145+
resultFuture.addListener(() -> timeoutFuture.cancel(false), MoreExecutors.directExecutor());
128146
ApiFutures.addCallback(
129147
startFuture,
130148
new ApiFutureCallback<ResumableUploadSession>() {
@@ -178,6 +196,17 @@ public void onFailure(Throwable t) {
178196
MoreExecutors.directExecutor());
179197
}
180198

199+
private void onTimeout() {
200+
String sessionUrl = uploadSessionUrl;
201+
String message;
202+
if (sessionUrl != null) {
203+
message = "Resumable upload timed out for session: " + sessionUrl;
204+
} else {
205+
message = "Resumable upload timed out before session initiation completed";
206+
}
207+
fail(new DeadlineExceededException(message, null, TIMEOUT_STATUS_CODE, false));
208+
}
209+
181210
private void succeed(@Nullable ResponseT result) {
182211
synchronized (lock) {
183212
inFlightFuture = null;
@@ -187,8 +216,13 @@ private void succeed(@Nullable ResponseT result) {
187216
}
188217

189218
private void fail(Throwable t) {
219+
ApiFuture<?> inFlight;
190220
synchronized (lock) {
191-
inFlightFuture = null;
221+
inFlight = this.inFlightFuture;
222+
this.inFlightFuture = null;
223+
}
224+
if (inFlight != null) {
225+
inFlight.cancel(true);
192226
}
193227
closePayload();
194228
resultFuture.setException(t);
@@ -248,4 +282,17 @@ public ResponseT get(long timeout, TimeUnit unit)
248282
throws InterruptedException, ExecutionException, TimeoutException {
249283
return resultFuture.get(timeout, unit);
250284
}
285+
286+
private static final StatusCode TIMEOUT_STATUS_CODE =
287+
new StatusCode() {
288+
@Override
289+
public StatusCode.Code getCode() {
290+
return StatusCode.Code.DEADLINE_EXCEEDED;
291+
}
292+
293+
@Override
294+
public @Nullable Object getTransportCode() {
295+
return null;
296+
}
297+
};
251298
}

‎sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ResumableUploadCallableImplTest.java‎

Lines changed: 211 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,8 @@
3232
import static com.google.common.truth.Truth.assertThat;
3333
import static org.junit.jupiter.api.Assertions.assertThrows;
3434
import static org.mockito.ArgumentMatchers.any;
35+
import static org.mockito.ArgumentMatchers.anyLong;
36+
import static org.mockito.ArgumentMatchers.eq;
3537
import static org.mockito.Mockito.lenient;
3638
import static org.mockito.Mockito.mock;
3739
import static org.mockito.Mockito.times;
@@ -55,10 +57,13 @@
5557
import java.io.IOException;
5658
import java.io.InputStream;
5759
import java.nio.charset.StandardCharsets;
60+
import java.time.Duration;
5861
import java.util.List;
5962
import java.util.concurrent.CancellationException;
6063
import java.util.concurrent.CountDownLatch;
6164
import java.util.concurrent.ExecutionException;
65+
import java.util.concurrent.ScheduledExecutorService;
66+
import java.util.concurrent.ScheduledFuture;
6267
import java.util.concurrent.TimeUnit;
6368
import java.util.concurrent.atomic.AtomicReference;
6469
import org.jspecify.annotations.Nullable;
@@ -78,6 +83,7 @@ class ResumableUploadCallableImplTest {
7883

7984
private ResumableUploadCallSettings defaultSettings;
8085
private FakeCallContext callContext;
86+
private ClientContext clientContext;
8187
private ResumableUploadCallableImpl<String, String> callable;
8288

8389
@BeforeEach
@@ -94,8 +100,7 @@ void setUp() {
94100

95101
defaultSettings = ResumableUploadCallSettings.newBuilder().setChunkSize(8).build();
96102
callContext = FakeCallContext.createDefault();
97-
ClientContext clientContext =
98-
ClientContext.newBuilder().setDefaultCallContext(callContext).build();
103+
clientContext = ClientContext.newBuilder().setDefaultCallContext(callContext).build();
99104
callable = new ResumableUploadCallableImpl<>(mockClient, defaultSettings, clientContext);
100105
}
101106

@@ -792,6 +797,210 @@ void testRecovery_queryTransientError_retriesAndSucceeds() throws Exception {
792797
verify(mockChunkCallable, times(2)).futureCall(any(), any());
793798
}
794799

800+
@Test
801+
void testGlobalTimeout_firesAndFailsSessionWithDeadlineExceeded() throws Exception {
802+
stubStartSession("https://upload.url/timeout-fire");
803+
SettableApiFuture<ChunkUploadResponse<String>> hungChunk = SettableApiFuture.create();
804+
when(mockChunkCallable.futureCall(any(ChunkUploadRequest.class), any())).thenReturn(hungChunk);
805+
806+
ResumableUploadCallSettings timeoutSettings =
807+
defaultSettings.toBuilder().setGlobalTimeout(Duration.ofMillis(100)).build();
808+
809+
ResumableUploadFuture<String> future =
810+
callable.futureCall("resource-path", streamOf("hello"), timeoutSettings);
811+
812+
ExecutionException exception =
813+
assertThrows(ExecutionException.class, () -> future.get(5, TimeUnit.SECONDS));
814+
assertThat(exception.getCause()).isInstanceOf(DeadlineExceededException.class);
815+
DeadlineExceededException cause = (DeadlineExceededException) exception.getCause();
816+
assertThat(cause.getStatusCode().getCode()).isEqualTo(StatusCode.Code.DEADLINE_EXCEEDED);
817+
assertThat(cause.getMessage()).contains("https://upload.url/timeout-fire");
818+
assertThat(future.isDone()).isTrue();
819+
assertThat(future.isCancelled()).isFalse();
820+
}
821+
822+
@Test
823+
void testGlobalTimeout_cancelledCleanlyOnSuccess() throws Exception {
824+
stubStartSession("https://upload.url/timeout-success");
825+
when(mockChunkCallable.futureCall(any(ChunkUploadRequest.class), any()))
826+
.thenReturn(
827+
ApiFutures.immediateFuture(
828+
ChunkUploadResponse.create(ResumableUploadStatus.FINAL, "ok")));
829+
830+
ScheduledExecutorService mockExecutor = mock(ScheduledExecutorService.class);
831+
ScheduledFuture<?> mockScheduledFuture = mock(ScheduledFuture.class);
832+
when(mockExecutor.schedule(any(Runnable.class), anyLong(), any()))
833+
.thenAnswer(inv -> mockScheduledFuture);
834+
835+
ClientContext customClientContext = clientContext.toBuilder().setExecutor(mockExecutor).build();
836+
ResumableUploadCallableImpl<String, String> customCallable =
837+
new ResumableUploadCallableImpl<>(mockClient, defaultSettings, customClientContext);
838+
839+
ResumableUploadCallSettings timeoutSettings =
840+
defaultSettings.toBuilder().setGlobalTimeout(Duration.ofSeconds(60)).build();
841+
842+
ResumableUploadFuture<String> future =
843+
customCallable.futureCall("resource-path", streamOf("hello"), timeoutSettings);
844+
845+
assertThat(future.get()).isEqualTo("ok");
846+
verify(mockScheduledFuture).cancel(false);
847+
}
848+
849+
@Test
850+
void testGlobalTimeout_cancelledCleanlyOnFailure() throws Exception {
851+
when(mockStartCallable.futureCall(any(), any()))
852+
.thenReturn(
853+
ApiFutures.immediateFailedFuture(
854+
createApiException(401, StatusCode.Code.UNAUTHENTICATED)));
855+
856+
ScheduledExecutorService mockExecutor = mock(ScheduledExecutorService.class);
857+
ScheduledFuture<?> mockScheduledFuture = mock(ScheduledFuture.class);
858+
when(mockExecutor.schedule(any(Runnable.class), anyLong(), any()))
859+
.thenAnswer(inv -> mockScheduledFuture);
860+
861+
ClientContext customClientContext = clientContext.toBuilder().setExecutor(mockExecutor).build();
862+
ResumableUploadCallableImpl<String, String> customCallable =
863+
new ResumableUploadCallableImpl<>(mockClient, defaultSettings, customClientContext);
864+
865+
ResumableUploadCallSettings timeoutSettings =
866+
defaultSettings.toBuilder().setGlobalTimeout(Duration.ofSeconds(60)).build();
867+
868+
ResumableUploadFuture<String> future =
869+
customCallable.futureCall("resource-path", streamOf("hello"), timeoutSettings);
870+
871+
assertThrows(ExecutionException.class, future::get);
872+
verify(mockScheduledFuture).cancel(false);
873+
}
874+
875+
@Test
876+
void testGlobalTimeout_cancelledCleanlyOnUserCancel() throws Exception {
877+
stubStartSession("https://upload.url/timeout-cancel");
878+
SettableApiFuture<ChunkUploadResponse<String>> hungChunk = SettableApiFuture.create();
879+
when(mockChunkCallable.futureCall(any(ChunkUploadRequest.class), any())).thenReturn(hungChunk);
880+
881+
ScheduledExecutorService mockExecutor = mock(ScheduledExecutorService.class);
882+
ScheduledFuture<?> mockScheduledFuture = mock(ScheduledFuture.class);
883+
when(mockExecutor.schedule(any(Runnable.class), anyLong(), any()))
884+
.thenAnswer(inv -> mockScheduledFuture);
885+
886+
ClientContext customClientContext = clientContext.toBuilder().setExecutor(mockExecutor).build();
887+
ResumableUploadCallableImpl<String, String> customCallable =
888+
new ResumableUploadCallableImpl<>(mockClient, defaultSettings, customClientContext);
889+
890+
ResumableUploadCallSettings timeoutSettings =
891+
defaultSettings.toBuilder().setGlobalTimeout(Duration.ofSeconds(60)).build();
892+
893+
ResumableUploadFuture<String> future =
894+
customCallable.futureCall("resource-path", streamOf("hello"), timeoutSettings);
895+
896+
assertThat(future.cancel(true)).isTrue();
897+
verify(mockScheduledFuture).cancel(false);
898+
}
899+
900+
@Test
901+
void testGlobalTimeout_usesDefaultWhenUnset() throws Exception {
902+
stubStartSession("https://upload.url/default-timeout");
903+
when(mockChunkCallable.futureCall(any(ChunkUploadRequest.class), any()))
904+
.thenReturn(
905+
ApiFutures.immediateFuture(
906+
ChunkUploadResponse.create(ResumableUploadStatus.FINAL, "ok")));
907+
908+
ScheduledExecutorService mockExecutor = mock(ScheduledExecutorService.class);
909+
ScheduledFuture<?> mockScheduledFuture = mock(ScheduledFuture.class);
910+
when(mockExecutor.schedule(any(Runnable.class), anyLong(), any()))
911+
.thenAnswer(inv -> mockScheduledFuture);
912+
913+
ClientContext customClientContext = clientContext.toBuilder().setExecutor(mockExecutor).build();
914+
915+
// 1. Unset on both stub and per-request -> falls back to GAX default (15m)
916+
ResumableUploadCallableImpl<String, String> unsetStubCallable =
917+
new ResumableUploadCallableImpl<>(mockClient, defaultSettings, customClientContext);
918+
assertThat(unsetStubCallable.futureCall("resource-path", streamOf("hello"), null).get())
919+
.isEqualTo("ok");
920+
verify(mockExecutor)
921+
.schedule(
922+
any(Runnable.class), eq(Duration.ofMinutes(15).toMillis()), eq(TimeUnit.MILLISECONDS));
923+
924+
// 2. Stub-level timeout (30m, e.g. from generator/client settings) + null per-request -> 30m
925+
ResumableUploadCallSettings stubWith30m =
926+
defaultSettings.toBuilder().setGlobalTimeout(Duration.ofMinutes(30)).build();
927+
ResumableUploadCallableImpl<String, String> configuredStubCallable =
928+
new ResumableUploadCallableImpl<>(mockClient, stubWith30m, customClientContext);
929+
assertThat(configuredStubCallable.futureCall("resource-path", streamOf("hello"), null).get())
930+
.isEqualTo("ok");
931+
verify(mockExecutor)
932+
.schedule(
933+
any(Runnable.class), eq(Duration.ofMinutes(30).toMillis()), eq(TimeUnit.MILLISECONDS));
934+
935+
// 3. Stub-level timeout (30m) + per-request with only chunkSize set -> preserves 30m
936+
ResumableUploadCallSettings perRequestChunkSizeOnly =
937+
ResumableUploadCallSettings.newBuilder().setChunkSize(16).build();
938+
assertThat(
939+
configuredStubCallable
940+
.futureCall("resource-path", streamOf("hello"), perRequestChunkSizeOnly)
941+
.get())
942+
.isEqualTo("ok");
943+
verify(mockExecutor, times(2))
944+
.schedule(
945+
any(Runnable.class), eq(Duration.ofMinutes(30).toMillis()), eq(TimeUnit.MILLISECONDS));
946+
947+
// 4. Stub-level timeout (30m) + per-request globalTimeout (5m) -> per-request wins (5m)
948+
ResumableUploadCallSettings perRequestWith5m =
949+
ResumableUploadCallSettings.newBuilder().setGlobalTimeout(Duration.ofMinutes(5)).build();
950+
assertThat(
951+
configuredStubCallable
952+
.futureCall("resource-path", streamOf("hello"), perRequestWith5m)
953+
.get())
954+
.isEqualTo("ok");
955+
verify(mockExecutor)
956+
.schedule(
957+
any(Runnable.class), eq(Duration.ofMinutes(5).toMillis()), eq(TimeUnit.MILLISECONDS));
958+
}
959+
960+
@Test
961+
void testGlobalTimeout_timeoutWhileAttemptInFlight_cancelsInFlightFutureAndDoesNotCorruptBuffer()
962+
throws Exception {
963+
stubStartSession("https://upload.url/in-flight-timeout");
964+
SettableApiFuture<ChunkUploadResponse<String>> inFlightFuture = SettableApiFuture.create();
965+
when(mockChunkCallable.futureCall(any(ChunkUploadRequest.class), any()))
966+
.thenReturn(inFlightFuture);
967+
968+
ResumableUploadCallSettings timeoutSettings =
969+
defaultSettings.toBuilder().setGlobalTimeout(Duration.ofMillis(80)).build();
970+
971+
TrackableStream stream = new TrackableStream("01234567890123456789");
972+
ResumableUploadFuture<String> future =
973+
callable.futureCall("resource-path", stream, timeoutSettings);
974+
975+
ExecutionException exception =
976+
assertThrows(ExecutionException.class, () -> future.get(5, TimeUnit.SECONDS));
977+
assertThat(exception.getCause()).isInstanceOf(DeadlineExceededException.class);
978+
// In-flight attempt future must be cancelled
979+
assertThat(inFlightFuture.isCancelled()).isTrue();
980+
981+
// Stream should have been read only up to the first chunk (chunkSize = 8), not refilled or
982+
// advanced
983+
assertThat(stream.totalBytesRead).isEqualTo(8);
984+
}
985+
986+
@Test
987+
void testGlobalTimeout_coversStartSessionTimeout() throws Exception {
988+
SettableApiFuture<ResumableUploadSession> hungStartFuture = SettableApiFuture.create();
989+
when(mockStartCallable.futureCall(any(), any())).thenReturn(hungStartFuture);
990+
991+
ResumableUploadCallSettings timeoutSettings =
992+
defaultSettings.toBuilder().setGlobalTimeout(Duration.ofMillis(80)).build();
993+
994+
ResumableUploadFuture<String> future =
995+
callable.futureCall("resource-path", streamOf("hello"), timeoutSettings);
996+
997+
ExecutionException exception =
998+
assertThrows(ExecutionException.class, () -> future.get(5, TimeUnit.SECONDS));
999+
assertThat(exception.getCause()).isInstanceOf(DeadlineExceededException.class);
1000+
assertThat(exception.getCause().getMessage()).contains("before session initiation completed");
1001+
assertThat(hungStartFuture.isCancelled()).isTrue();
1002+
}
1003+
7951004
private static class HttpStatusStatusCode implements StatusCode {
7961005
private final int httpStatus;
7971006
private final StatusCode.Code code;

0 commit comments

Comments
 (0)