Skip to content

Commit 3472f61

Browse files
committed
feat(gax): implement MVP ResumableUploadCallable and ResumableUploadFuture
1 parent a03ee10 commit 3472f61

6 files changed

Lines changed: 876 additions & 8 deletions

File tree

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

Lines changed: 17 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -31,38 +31,49 @@
3131

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

3537
/**
3638
* A ResumableUploadCallable is an API-transport-independent wrapper for the Resumable Upload
3739
* protocol. Operates directly on the request object and input stream payload.
3840
*
39-
* @param <RequestT> request type
40-
* @param <ResponseT> response type
41+
* @param <RequestT> the type of the initial request message that initiates the upload session
42+
* @param <ResponseT> the type of the final response message returned once the upload completes
4143
*/
4244
@BetaApi
45+
@NullMarked
4346
public abstract class ResumableUploadCallable<RequestT, ResponseT> {
4447

4548
protected ResumableUploadCallable() {}
4649

4750
/**
4851
* Performs a new resumable upload asynchronously.
4952
*
53+
* <p>The provided {@code payload} stream is consumed asynchronously by the returned {@link
54+
* ResumableUploadFuture} and will be closed automatically upon completion, failure, or
55+
* cancellation.
56+
*
5057
* @param request the request message
51-
* @param payload the data payload input stream
58+
* @param payload the data payload input stream to upload and close
5259
* @param settings call settings overrides; may be {@code null}
5360
* @return future for tracking and controlling the upload
5461
*/
5562
public abstract ResumableUploadFuture<ResponseT> futureCall(
56-
RequestT request, InputStream payload, ResumableUploadCallSettings settings);
63+
RequestT request, InputStream payload, @Nullable ResumableUploadCallSettings settings);
5764

5865
/**
5966
* Resumes an existing resumable upload session asynchronously using a saved session URL.
6067
*
68+
* <p>The provided {@code payload} stream is consumed asynchronously by the returned {@link
69+
* ResumableUploadFuture} and will be closed automatically upon completion, failure, or
70+
* cancellation.
71+
*
6172
* @param sessionUrl the upload session URL
62-
* @param payload the data payload input stream
73+
* @param payload the data payload input stream to upload and close
6374
* @param settings call settings overrides; may be {@code null}
6475
* @return future for tracking and controlling the upload
6576
*/
6677
public abstract ResumableUploadFuture<ResponseT> resumeCall(
67-
String sessionUrl, InputStream payload, ResumableUploadCallSettings settings);
78+
String sessionUrl, InputStream payload, @Nullable ResumableUploadCallSettings settings);
6879
}
Lines changed: 99 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,99 @@
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.ApiFuture;
35+
import com.google.api.core.ApiFutures;
36+
import com.google.api.core.BetaApi;
37+
import com.google.api.core.InternalApi;
38+
import com.google.api.gax.resumable.ResumableUploadClient;
39+
import com.google.api.gax.resumable.ResumableUploadSession;
40+
import java.io.InputStream;
41+
import org.jspecify.annotations.NullMarked;
42+
import org.jspecify.annotations.Nullable;
43+
44+
/**
45+
* Concrete implementation of {@link ResumableUploadCallable} that delegates the end-to-end
46+
* management of a resumable upload session to {@link ResumableUploadFutureImpl}.
47+
*
48+
* @param <RequestT> the type of the initial request message that initiates the upload session
49+
* @param <ResponseT> the type of the final response message returned once the upload completes
50+
*/
51+
@BetaApi
52+
@InternalApi
53+
@NullMarked
54+
public class ResumableUploadCallableImpl<RequestT, ResponseT>
55+
extends ResumableUploadCallable<RequestT, ResponseT> {
56+
57+
private final ResumableUploadClient<RequestT, ResponseT> client;
58+
private final ResumableUploadCallSettings defaultCallSettings;
59+
private final ApiCallContext defaultCallContext;
60+
61+
public ResumableUploadCallableImpl(
62+
ResumableUploadClient<RequestT, ResponseT> client,
63+
ResumableUploadCallSettings defaultCallSettings,
64+
ApiCallContext defaultCallContext) {
65+
this.client = checkNotNull(client, "client must not be null");
66+
this.defaultCallSettings =
67+
checkNotNull(defaultCallSettings, "defaultCallSettings must not be null");
68+
this.defaultCallContext =
69+
checkNotNull(defaultCallContext, "defaultCallContext 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+
ApiFuture<ResumableUploadSession> startFuture;
80+
try {
81+
startFuture = client.startUploadCallable().futureCall(request, defaultCallContext);
82+
} catch (Throwable t) {
83+
startFuture = ApiFutures.immediateFailedFuture(t);
84+
}
85+
86+
return ResumableUploadFutureImpl.create(
87+
startFuture,
88+
client.uploadChunkCallable(),
89+
payload,
90+
effectiveSettings,
91+
defaultCallContext);
92+
}
93+
94+
@Override
95+
public ResumableUploadFuture<ResponseT> resumeCall(
96+
String sessionUrl, InputStream payload, @Nullable ResumableUploadCallSettings settings) {
97+
throw new UnsupportedOperationException("Session resumption is not yet implemented.");
98+
}
99+
}
Lines changed: 164 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,164 @@
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.ApiFuture;
35+
import com.google.api.core.ApiFutureCallback;
36+
import com.google.api.core.ApiFutures;
37+
import com.google.api.core.InternalApi;
38+
import com.google.api.gax.resumable.ChunkUploadRequest;
39+
import com.google.api.gax.resumable.ChunkUploadResponse;
40+
import com.google.common.io.ByteStreams;
41+
import com.google.common.util.concurrent.MoreExecutors;
42+
import java.io.IOException;
43+
import java.io.InputStream;
44+
import java.util.Arrays;
45+
import java.util.concurrent.CancellationException;
46+
import org.jspecify.annotations.NullMarked;
47+
48+
/**
49+
* Coordinates chunk transmission steps of a resumable upload session.
50+
*
51+
* @param <ResponseT> the type of the final response message returned once the upload completes
52+
*/
53+
@InternalApi
54+
@NullMarked
55+
final class ResumableUploadChunkCoordinator<ResponseT> {
56+
57+
private static final byte[] EMPTY_PAYLOAD = new byte[0];
58+
59+
private final UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>>
60+
uploadChunkCallable;
61+
private final String uploadUrl;
62+
private final InputStream payload;
63+
private final byte[] buffer;
64+
private final int chunkSize;
65+
private final ApiCallContext callContext;
66+
private final ResumableUploadFutureImpl<ResponseT> sessionFuture;
67+
68+
ResumableUploadChunkCoordinator(
69+
UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>> uploadChunkCallable,
70+
String uploadUrl,
71+
InputStream payload,
72+
int chunkSize,
73+
ApiCallContext callContext,
74+
ResumableUploadFutureImpl<ResponseT> sessionFuture) {
75+
this.uploadChunkCallable =
76+
checkNotNull(uploadChunkCallable, "uploadChunkCallable must not be null");
77+
this.uploadUrl = checkNotNull(uploadUrl, "uploadUrl must not be null");
78+
this.payload = checkNotNull(payload, "payload must not be null");
79+
this.chunkSize = chunkSize;
80+
this.callContext = checkNotNull(callContext, "callContext must not be null");
81+
this.sessionFuture = checkNotNull(sessionFuture, "sessionFuture must not be null");
82+
this.buffer = new byte[chunkSize];
83+
}
84+
85+
void start() {
86+
transmitChunk(0L);
87+
}
88+
89+
private void transmitChunk(long currentOffset) {
90+
// Abort if the session was already completed or canceled.
91+
if (sessionFuture.isDone()) {
92+
return;
93+
}
94+
95+
// Read the next chunk slice from the payload stream.
96+
int bytesRead;
97+
try {
98+
bytesRead = ByteStreams.read(payload, buffer, 0, chunkSize);
99+
} catch (IOException e) {
100+
sessionFuture.fail(e);
101+
return;
102+
}
103+
104+
// Determine if this is the final chunk and build the chunk request.
105+
boolean isFinal = bytesRead < chunkSize;
106+
byte[] chunkPayload;
107+
if (bytesRead == chunkSize) {
108+
chunkPayload = buffer;
109+
} else if (bytesRead == 0) {
110+
chunkPayload = EMPTY_PAYLOAD;
111+
} else {
112+
chunkPayload = Arrays.copyOf(buffer, bytesRead);
113+
}
114+
115+
ChunkUploadRequest chunkRequest =
116+
ChunkUploadRequest.newBuilder()
117+
.setUploadUrl(uploadUrl)
118+
.setPayload(chunkPayload)
119+
.setOffset(currentOffset)
120+
.setFinal(isFinal)
121+
.build();
122+
123+
// Dispatch the chunk upload call and register the in-flight future for cancellation.
124+
long chunkLength = chunkPayload.length;
125+
try {
126+
ApiFuture<ChunkUploadResponse<ResponseT>> chunkFuture =
127+
uploadChunkCallable.futureCall(chunkRequest, callContext);
128+
sessionFuture.setInFlightFuture(chunkFuture);
129+
130+
// Asynchronously handle the response: complete, fail, or chain the next chunk.
131+
ApiFutures.addCallback(
132+
chunkFuture,
133+
new ApiFutureCallback<ChunkUploadResponse<ResponseT>>() {
134+
@Override
135+
public void onSuccess(ChunkUploadResponse<ResponseT> response) {
136+
if (sessionFuture.isDone()) {
137+
return;
138+
}
139+
long nextOffset = currentOffset + chunkLength;
140+
if (response.isComplete()) {
141+
sessionFuture.succeed(response.getResponse());
142+
} else if (isFinal) {
143+
sessionFuture.fail(
144+
new IllegalStateException(
145+
"Upload stream ended and final chunk was transmitted, but server returned"
146+
+ " incomplete status"));
147+
} else {
148+
transmitChunk(nextOffset);
149+
}
150+
}
151+
152+
@Override
153+
public void onFailure(Throwable t) {
154+
if (t instanceof CancellationException || sessionFuture.isDone()) {
155+
return;
156+
}
157+
sessionFuture.fail(t);
158+
}
159+
}, MoreExecutors.directExecutor());
160+
} catch (Throwable t) {
161+
sessionFuture.fail(t);
162+
}
163+
}
164+
}

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

Lines changed: 8 additions & 2 deletions
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
*
38-
* @param <ResponseT> response type
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+
*
43+
* @param <ResponseT> the type of the final response message returned once the upload completes
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
}

0 commit comments

Comments
 (0)