Skip to content

Commit 554caca

Browse files
committed
feat(gax): implement MVP ResumableUploadCallable and ResumableUploadFuture
1 parent 1c643de commit 554caca

5 files changed

Lines changed: 691 additions & 3 deletions

File tree

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

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -47,8 +47,12 @@ protected ResumableUploadCallable() {}
4747
/**
4848
* Performs a new resumable upload asynchronously.
4949
*
50+
* <p>The provided {@code payload} stream is consumed asynchronously by the returned {@link
51+
* ResumableUploadFuture} and will be closed automatically upon completion, failure, or
52+
* cancellation.
53+
*
5054
* @param request the request message
51-
* @param payload the data payload input stream
55+
* @param payload the data payload input stream to upload and close
5256
* @param settings call settings overrides; may be {@code null}
5357
* @return future for tracking and controlling the upload
5458
*/
@@ -58,8 +62,12 @@ public abstract ResumableUploadFuture<ResponseT> futureCall(
5862
/**
5963
* Resumes an existing resumable upload session asynchronously using a saved session URL.
6064
*
65+
* <p>The provided {@code payload} stream is consumed asynchronously by the returned {@link
66+
* ResumableUploadFuture} and will be closed automatically upon completion, failure, or
67+
* cancellation.
68+
*
6169
* @param sessionUrl the upload session URL
62-
* @param payload the data payload input stream
70+
* @param payload the data payload input stream to upload and close
6371
* @param settings call settings overrides; may be {@code null}
6472
* @return future for tracking and controlling the upload
6573
*/
Lines changed: 88 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,88 @@
1+
/*
2+
* Copyright 2026 Google LLC
3+
*
4+
* Redistribution and use in source and binary forms, with or without
5+
* modification, are permitted provided that the following conditions are
6+
* met:
7+
*
8+
* * Redistributions of source code must retain the above copyright
9+
* notice, this list of conditions and the following disclaimer.
10+
* * Redistributions in binary form must reproduce the above
11+
* copyright notice, this list of conditions and the following disclaimer
12+
* in the documentation and/or other materials provided with the
13+
* distribution.
14+
* * Neither the name of Google LLC nor the names of its
15+
* contributors may be used to endorse or promote products derived from
16+
* this software without specific prior written permission.
17+
*
18+
* THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS
19+
* "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT
20+
* LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR
21+
* A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT
22+
* OWNER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL,
23+
* SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT
24+
* LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE,
25+
* DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY
26+
* THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT
27+
* (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE
28+
* OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
29+
*/
30+
package com.google.api.gax.rpc;
31+
32+
import static com.google.common.base.Preconditions.checkNotNull;
33+
34+
import com.google.api.core.BetaApi;
35+
import com.google.api.core.InternalApi;
36+
import com.google.api.gax.resumable.ResumableUploadClient;
37+
import java.io.InputStream;
38+
import java.util.concurrent.Executor;
39+
import org.jspecify.annotations.NullMarked;
40+
import org.jspecify.annotations.Nullable;
41+
42+
/**
43+
* Concrete implementation of {@link ResumableUploadCallable} that delegates the upload pipeline to
44+
* {@link ResumableUploadFutureImpl}.
45+
*
46+
* @param <RequestT> request type
47+
* @param <ResponseT> response type
48+
*/
49+
@BetaApi
50+
@InternalApi
51+
@NullMarked
52+
public class ResumableUploadCallableImpl<RequestT, ResponseT>
53+
extends ResumableUploadCallable<RequestT, ResponseT> {
54+
55+
private final ResumableUploadClient<RequestT, ResponseT> client;
56+
private final ResumableUploadCallSettings defaultCallSettings;
57+
private final @Nullable ApiCallContext defaultCallContext;
58+
private final Executor executor;
59+
60+
public ResumableUploadCallableImpl(
61+
ResumableUploadClient<RequestT, ResponseT> client,
62+
ResumableUploadCallSettings defaultCallSettings,
63+
@Nullable ApiCallContext defaultCallContext,
64+
Executor executor) {
65+
this.client = checkNotNull(client, "client must not be null");
66+
this.defaultCallSettings =
67+
checkNotNull(defaultCallSettings, "defaultCallSettings must not be null");
68+
this.defaultCallContext = defaultCallContext;
69+
this.executor = checkNotNull(executor, "executor must not be null");
70+
}
71+
72+
@Override
73+
public ResumableUploadFuture<ResponseT> futureCall(
74+
RequestT request, InputStream payload, @Nullable ResumableUploadCallSettings settings) {
75+
checkNotNull(request, "request must not be null");
76+
checkNotNull(payload, "payload must not be null");
77+
ResumableUploadCallSettings effectiveSettings = defaultCallSettings.merge(settings);
78+
79+
return ResumableUploadFutureImpl.create(
80+
client, request, payload, effectiveSettings.getChunkSize(), defaultCallContext, executor);
81+
}
82+
83+
@Override
84+
public ResumableUploadFuture<ResponseT> resumeCall(
85+
String sessionUrl, InputStream payload, @Nullable ResumableUploadCallSettings settings) {
86+
throw new UnsupportedOperationException("Session resumption is not yet implemented.");
87+
}
88+
}

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

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -31,15 +31,21 @@
3131

3232
import com.google.api.core.ApiFuture;
3333
import com.google.api.core.BetaApi;
34+
import org.jspecify.annotations.NullMarked;
35+
import org.jspecify.annotations.Nullable;
3436

3537
/**
3638
* A specialized {@link ApiFuture} for tracking and controlling an in-flight resumable upload.
3739
*
40+
* <p>The payload {@link java.io.InputStream} supplied when initiating the upload is managed by this
41+
* future and will be closed automatically upon completion, failure, or cancellation.
42+
*
3843
* @param <ResponseT> response type
3944
*/
4045
@BetaApi
46+
@NullMarked
4147
public interface ResumableUploadFuture<ResponseT> extends ApiFuture<ResponseT> {
4248

4349
/** Returns the upload session URL, or {@code null} if session initiation is in progress. */
44-
String getUploadSessionUrl();
50+
@Nullable String getUploadSessionUrl();
4551
}
Lines changed: 214 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,214 @@
1+
/*
2+
* Copyright 2026 Google LLC
3+
*
4+
* Redistribution and use in source and binary forms, with or without
5+
* modification, are permitted provided that the following conditions are
6+
* met:
7+
*
8+
* * Redistributions of source code must retain the above copyright
9+
* notice, this list of conditions and the following disclaimer.
10+
* * Redistributions in binary form must reproduce the above
11+
* copyright notice, this list of conditions and the following disclaimer
12+
* in the documentation and/or other materials provided with the
13+
* distribution.
14+
* * Neither the name of Google LLC nor the names of its
15+
* contributors may be used to endorse or promote products derived from
16+
* this software without specific prior written permission.
17+
*
18+
* THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS
19+
* "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT
20+
* LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR
21+
* A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT
22+
* OWNER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL,
23+
* SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT
24+
* LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE,
25+
* DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY
26+
* THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT
27+
* (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE
28+
* OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
29+
*/
30+
package com.google.api.gax.rpc;
31+
32+
import static com.google.common.base.Preconditions.checkArgument;
33+
import static com.google.common.base.Preconditions.checkNotNull;
34+
35+
import com.google.api.core.ApiFuture;
36+
import com.google.api.core.ApiFutures;
37+
import com.google.api.gax.resumable.ChunkUploadRequest;
38+
import com.google.api.gax.resumable.ChunkUploadResponse;
39+
import com.google.api.gax.resumable.ResumableUploadClient;
40+
import com.google.api.gax.resumable.ResumableUploadSession;
41+
import com.google.common.io.ByteStreams;
42+
import java.io.IOException;
43+
import java.io.InputStream;
44+
import java.util.Arrays;
45+
import java.util.concurrent.ExecutionException;
46+
import java.util.concurrent.Executor;
47+
import java.util.concurrent.TimeUnit;
48+
import java.util.concurrent.TimeoutException;
49+
import java.util.concurrent.atomic.AtomicReference;
50+
import org.jspecify.annotations.NullMarked;
51+
import org.jspecify.annotations.Nullable;
52+
53+
/**
54+
* Implementation of {@link ResumableUploadFuture} that coordinates session initiation and chunk
55+
* streaming.
56+
*
57+
* <p>The provided payload stream is automatically closed upon completion, failure, or cancellation.
58+
*
59+
* @param <ResponseT> response type
60+
*/
61+
@NullMarked
62+
final class ResumableUploadFutureImpl<ResponseT> implements ResumableUploadFuture<ResponseT> {
63+
64+
private final InputStream payload;
65+
private final AtomicReference<@Nullable String> uploadSessionUrl;
66+
private final ApiFuture<ResponseT> delegate;
67+
68+
static <RequestT, ResponseT> ResumableUploadFutureImpl<ResponseT> create(
69+
ResumableUploadClient<RequestT, ResponseT> client,
70+
RequestT initialRequest,
71+
InputStream payload,
72+
int chunkSize,
73+
@Nullable ApiCallContext callContext,
74+
Executor executor) {
75+
checkNotNull(client, "client must not be null");
76+
checkNotNull(initialRequest, "initialRequest must not be null");
77+
checkNotNull(payload, "payload must not be null");
78+
checkArgument(chunkSize > 0, "chunkSize must be > 0");
79+
checkNotNull(executor, "executor must not be null");
80+
81+
AtomicReference<@Nullable String> uploadSessionUrl = new AtomicReference<>();
82+
83+
ApiFuture<ResponseT> pipelineFuture;
84+
try {
85+
// Initiate upload session asynchronously, then chain into the chunk transmission loop.
86+
pipelineFuture =
87+
ApiFutures.<ResumableUploadSession, ResponseT>transformAsync(
88+
client.startUploadCallable().futureCall(initialRequest, callContext),
89+
session -> {
90+
String url = session.getUploadUrl();
91+
uploadSessionUrl.set(url);
92+
// Begin transmitting chunks starting at byte offset 0.
93+
return transmitChunks(client, payload, chunkSize, callContext, url, 0L, executor);
94+
},
95+
executor);
96+
} catch (Throwable t) {
97+
closePayload(payload);
98+
pipelineFuture = ApiFutures.immediateFailedFuture(t);
99+
}
100+
101+
pipelineFuture.addListener(() -> closePayload(payload), executor);
102+
103+
return new ResumableUploadFutureImpl<>(payload, uploadSessionUrl, pipelineFuture);
104+
}
105+
106+
private ResumableUploadFutureImpl(
107+
InputStream payload,
108+
AtomicReference<@Nullable String> uploadSessionUrl,
109+
ApiFuture<ResponseT> delegate) {
110+
this.payload = payload;
111+
this.uploadSessionUrl = uploadSessionUrl;
112+
this.delegate = delegate;
113+
}
114+
115+
private static <ResponseT> ApiFuture<ResponseT> transmitChunks(
116+
ResumableUploadClient<?, ResponseT> client,
117+
InputStream payload,
118+
int chunkSize,
119+
@Nullable ApiCallContext callContext,
120+
String uploadSessionUrl,
121+
long bytesUploaded,
122+
Executor executor) {
123+
byte[] buffer = new byte[chunkSize];
124+
int bytesRead;
125+
try {
126+
bytesRead = ByteStreams.read(payload, buffer, 0, chunkSize);
127+
} catch (IOException e) {
128+
return ApiFutures.immediateFailedFuture(e);
129+
}
130+
131+
boolean isFinal = bytesRead < chunkSize;
132+
byte[] chunkPayload =
133+
(isFinal && bytesRead == 0)
134+
? new byte[0]
135+
: (bytesRead == chunkSize ? buffer : Arrays.copyOf(buffer, bytesRead));
136+
137+
ChunkUploadRequest chunkRequest =
138+
ChunkUploadRequest.newBuilder()
139+
.setUploadUrl(uploadSessionUrl)
140+
.setPayload(chunkPayload)
141+
.setOffset(bytesUploaded)
142+
.setFinal(isFinal)
143+
.build();
144+
145+
long nextOffset = bytesUploaded + chunkPayload.length;
146+
return ApiFutures.<ChunkUploadResponse<ResponseT>, ResponseT>transformAsync(
147+
client.uploadChunkCallable().futureCall(chunkRequest, callContext),
148+
(ChunkUploadResponse<ResponseT> response) -> {
149+
// Terminal success: server finalized upload and returned the completion response.
150+
if (response.isComplete()) {
151+
ResponseT completedResponse =
152+
checkNotNull(
153+
response.getResponse(), "Upload is complete but server returned null response");
154+
return ApiFutures.immediateFuture(completedResponse);
155+
}
156+
// Protocol error: payload stream reached EOF but server did not finalize upload.
157+
if (isFinal) {
158+
return ApiFutures.immediateFailedFuture(
159+
new IllegalStateException(
160+
"Upload stream ended and final chunk was transmitted, but server returned"
161+
+ " incomplete status"));
162+
}
163+
// Continuation: asynchronously transmit subsequent chunk with updated offset.
164+
return transmitChunks(
165+
client, payload, chunkSize, callContext, uploadSessionUrl, nextOffset, executor);
166+
},
167+
executor);
168+
}
169+
170+
private static void closePayload(InputStream payload) {
171+
try {
172+
payload.close();
173+
} catch (IOException ignored) {
174+
// Suppressed during stream cleanup
175+
}
176+
}
177+
178+
@Override
179+
public @Nullable String getUploadSessionUrl() {
180+
return uploadSessionUrl.get();
181+
}
182+
183+
@Override
184+
public void addListener(Runnable listener, Executor executor) {
185+
delegate.addListener(listener, executor);
186+
}
187+
188+
@Override
189+
public boolean cancel(boolean mayInterruptIfRunning) {
190+
closePayload(payload);
191+
return delegate.cancel(mayInterruptIfRunning);
192+
}
193+
194+
@Override
195+
public boolean isCancelled() {
196+
return delegate.isCancelled();
197+
}
198+
199+
@Override
200+
public boolean isDone() {
201+
return delegate.isDone();
202+
}
203+
204+
@Override
205+
public ResponseT get() throws InterruptedException, ExecutionException {
206+
return delegate.get();
207+
}
208+
209+
@Override
210+
public ResponseT get(long timeout, TimeUnit unit)
211+
throws InterruptedException, ExecutionException, TimeoutException {
212+
return delegate.get(timeout, unit);
213+
}
214+
}

0 commit comments

Comments
 (0)