Skip to content

Commit 1c643de

Browse files
authored
feat(gax): implement uploadChunkCallable for resumable uploads (#14140)
This handles sending binary chunks and finalize commands over HTTP/JSON, extracting offset and status from response headers/codes.
1 parent 4723064 commit 1c643de

8 files changed

Lines changed: 659 additions & 0 deletions

File tree

‎sdk-platform-java/gax-java/gax-httpjson/src/main/java/com/google/api/gax/httpjson/HttpJsonResumableUploadClient.java‎

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,8 @@
3232
import com.google.api.client.http.HttpMethods;
3333
import com.google.api.core.BetaApi;
3434
import com.google.api.core.InternalApi;
35+
import com.google.api.gax.resumable.ChunkUploadRequest;
36+
import com.google.api.gax.resumable.ChunkUploadResponse;
3537
import com.google.api.gax.resumable.ResumableUploadClient;
3638
import com.google.api.gax.resumable.ResumableUploadSession;
3739
import com.google.api.gax.rpc.ClientContext;
@@ -54,6 +56,8 @@ public final class HttpJsonResumableUploadClient<RequestT, ResponseT>
5456
implements ResumableUploadClient<RequestT, ResponseT> {
5557

5658
private final UnaryCallable<RequestT, ResumableUploadSession> startUploadCallable;
59+
private final UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>>
60+
uploadChunkCallable;
5761

5862
public static <RequestT, ResponseT> HttpJsonResumableUploadClient<RequestT, ResponseT> create(
5963
ClientContext clientContext, ApiMethodDescriptor<RequestT, ResponseT> methodDescriptor) {
@@ -64,6 +68,8 @@ private HttpJsonResumableUploadClient(
6468
ClientContext clientContext, ApiMethodDescriptor<RequestT, ResponseT> methodDescriptor) {
6569
Preconditions.checkNotNull(clientContext);
6670
Preconditions.checkNotNull(methodDescriptor);
71+
HttpResponseParser<ResponseT> responseParser =
72+
Preconditions.checkNotNull(methodDescriptor.getResponseParser());
6773

6874
ApiMethodDescriptor<RequestT, String> startUploadDescriptor =
6975
ApiMethodDescriptor.<RequestT, String>newBuilder()
@@ -75,10 +81,16 @@ private HttpJsonResumableUploadClient(
7581
.build();
7682
this.startUploadCallable =
7783
ResumableUploadStartCallable.create(clientContext, startUploadDescriptor);
84+
this.uploadChunkCallable = ResumableUploadChunkCallable.create(clientContext, responseParser);
7885
}
7986

8087
@Override
8188
public UnaryCallable<RequestT, ResumableUploadSession> startUploadCallable() {
8289
return startUploadCallable;
8390
}
91+
92+
@Override
93+
public UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>> uploadChunkCallable() {
94+
return uploadChunkCallable;
95+
}
8496
}
Lines changed: 228 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,228 @@
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.httpjson;
31+
32+
import com.google.api.client.http.HttpMethods;
33+
import com.google.api.core.ApiFuture;
34+
import com.google.api.gax.resumable.ChunkUploadRequest;
35+
import com.google.api.gax.resumable.ChunkUploadResponse;
36+
import com.google.api.gax.rpc.ApiCallContext;
37+
import com.google.api.gax.rpc.ApiExceptionFactory;
38+
import com.google.api.gax.rpc.ClientContext;
39+
import com.google.api.gax.rpc.StatusCode;
40+
import com.google.api.gax.rpc.UnaryCallable;
41+
import com.google.api.pathtemplate.PathTemplate;
42+
import com.google.common.base.Preconditions;
43+
import com.google.common.collect.ImmutableList;
44+
import com.google.common.collect.ImmutableMap;
45+
import java.io.ByteArrayInputStream;
46+
import java.io.InputStream;
47+
import java.nio.charset.StandardCharsets;
48+
import java.util.Collections;
49+
import java.util.List;
50+
import java.util.Map;
51+
import org.jspecify.annotations.NullMarked;
52+
import org.jspecify.annotations.Nullable;
53+
54+
/** A {@link UnaryCallable} that transmits individual chunks in a resumable upload session. */
55+
@NullMarked
56+
class ResumableUploadChunkCallable<ResponseT>
57+
extends UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>> {
58+
59+
private static final String UPLOAD_COMMAND_HEADER = "X-Goog-Upload-Command";
60+
private static final String UPLOAD_OFFSET_HEADER = "X-Goog-Upload-Offset";
61+
private static final String UPLOAD_STATUS_HEADER = "X-Goog-Upload-Status";
62+
private static final String STATUS_FINAL = "final";
63+
64+
private static final String COMMAND_UPLOAD = "upload";
65+
private static final String COMMAND_FINALIZE = "finalize";
66+
private static final String COMMAND_UPLOAD_FINALIZE = "upload, finalize";
67+
68+
private static final PathTemplate PATH_TEMPLATE = PathTemplate.create("**");
69+
70+
private static final ApiMethodDescriptor<ChunkUploadRequest, String> UPLOAD_CHUNK_DESCRIPTOR =
71+
ApiMethodDescriptor.<ChunkUploadRequest, String>newBuilder()
72+
.setFullMethodName("ResumableUpload/UploadChunk")
73+
.setHttpMethod(HttpMethods.POST)
74+
.setType(ApiMethodDescriptor.MethodType.UNARY)
75+
.setRequestFormatter(
76+
new ResumableUploadChunkRequestFormatter<ChunkUploadRequest>() {
77+
@Override
78+
public Map<String, List<String>> getQueryParamNames(ChunkUploadRequest request) {
79+
return Collections.emptyMap();
80+
}
81+
82+
@Override
83+
public byte[] getBinaryRequestBody(ChunkUploadRequest request) {
84+
return request.getPayload();
85+
}
86+
87+
@Override
88+
public String getPath(ChunkUploadRequest request) {
89+
return request.getUploadUrl();
90+
}
91+
92+
@Override
93+
public PathTemplate getPathTemplate() {
94+
return PATH_TEMPLATE;
95+
}
96+
})
97+
.setResponseParser(ResumableUploadResponseParser.create())
98+
.build();
99+
100+
private final ClientContext clientContext;
101+
private final HttpResponseParser<ResponseT> responseParser;
102+
103+
private ResumableUploadChunkCallable(
104+
ClientContext clientContext, HttpResponseParser<ResponseT> responseParser) {
105+
this.clientContext = Preconditions.checkNotNull(clientContext);
106+
this.responseParser = Preconditions.checkNotNull(responseParser);
107+
}
108+
109+
@Override
110+
public ApiFuture<ChunkUploadResponse<ResponseT>> futureCall(
111+
ChunkUploadRequest request, @Nullable ApiCallContext inputContext) {
112+
Preconditions.checkNotNull(request);
113+
boolean isPayloadEmpty = request.getPayload().length == 0;
114+
String command;
115+
if (request.isFinal()) {
116+
command = !isPayloadEmpty ? COMMAND_UPLOAD_FINALIZE : COMMAND_FINALIZE;
117+
} else {
118+
command = COMMAND_UPLOAD;
119+
}
120+
ImmutableMap.Builder<String, List<String>> chunkHeadersBuilder =
121+
ImmutableMap.<String, List<String>>builder()
122+
.put(UPLOAD_COMMAND_HEADER, ImmutableList.of(command));
123+
if (!COMMAND_FINALIZE.equals(command)) {
124+
chunkHeadersBuilder.put(
125+
UPLOAD_OFFSET_HEADER, ImmutableList.of(String.valueOf(request.getOffset())));
126+
}
127+
Map<String, List<String>> chunkHeaders = chunkHeadersBuilder.build();
128+
129+
HttpJsonCallContext context =
130+
(HttpJsonCallContext)
131+
HttpJsonCallContext.createDefault()
132+
.nullToSelf(clientContext.getDefaultCallContext())
133+
.merge(inputContext)
134+
.withExtraHeaders(chunkHeaders);
135+
136+
HttpJsonClientCall<ChunkUploadRequest, String> clientCall =
137+
HttpJsonClientCalls.newCall(UPLOAD_CHUNK_DESCRIPTOR, context);
138+
139+
ResumableUploadHttpJsonFuture<ChunkUploadResponse<ResponseT>> future =
140+
new ResumableUploadHttpJsonFuture<>(clientCall);
141+
HttpJsonClientCalls.startUnaryCall(
142+
clientCall, request, context, new ChunkUploadResponseListener<>(future, responseParser));
143+
144+
return future;
145+
}
146+
147+
static <ResponseT> UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>> create(
148+
ClientContext clientContext, HttpResponseParser<ResponseT> responseParser) {
149+
UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>> rawCallable =
150+
new ResumableUploadChunkCallable<>(clientContext, responseParser);
151+
UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>> callable =
152+
new HttpJsonExceptionCallable<>(
153+
rawCallable,
154+
// Wire calls do not retry directly; retries are managed by ResumableUploadCallable.
155+
Collections.emptySet());
156+
return callable.withDefaultCallContext(clientContext.getDefaultCallContext());
157+
}
158+
159+
/**
160+
* A listener that processes chunk upload response headers and bodies to produce the {@link
161+
* ChunkUploadResponse}.
162+
*/
163+
private static class ChunkUploadResponseListener<ResponseT>
164+
extends HttpJsonClientCall.Listener<String> {
165+
166+
private final ResumableUploadHttpJsonFuture<ChunkUploadResponse<ResponseT>> future;
167+
private final HttpResponseParser<ResponseT> responseParser;
168+
@Nullable private String uploadStatus = null;
169+
private String responseBody = "";
170+
171+
private ChunkUploadResponseListener(
172+
ResumableUploadHttpJsonFuture<ChunkUploadResponse<ResponseT>> future,
173+
HttpResponseParser<ResponseT> responseParser) {
174+
this.future = future;
175+
this.responseParser = responseParser;
176+
}
177+
178+
@Override
179+
public void onHeaders(HttpJsonMetadata responseHeaders) {
180+
Map<String, Object> headers = responseHeaders.getHeaders();
181+
this.uploadStatus = HttpHeadersUtils.getSingleHeader(headers, UPLOAD_STATUS_HEADER);
182+
}
183+
184+
@Override
185+
public void onMessage(@Nullable String message) {
186+
if (message != null) {
187+
this.responseBody = message;
188+
}
189+
}
190+
191+
@Override
192+
public void onClose(int statusCode, HttpJsonMetadata trailers) {
193+
try {
194+
if (statusCode >= 200 && statusCode < 300) {
195+
if (uploadStatus == null) {
196+
future.setException(
197+
ApiExceptionFactory.createException(
198+
"Upload chunk response did not contain valid "
199+
+ UPLOAD_STATUS_HEADER
200+
+ " header",
201+
/* cause= */ null,
202+
HttpJsonStatusCode.of(StatusCode.Code.INTERNAL),
203+
/* retryable= */ false));
204+
return;
205+
}
206+
boolean isComplete = STATUS_FINAL.equalsIgnoreCase(uploadStatus);
207+
ChunkUploadResponse.Builder<ResponseT> chunkResponseBuilder =
208+
ChunkUploadResponse.<ResponseT>newBuilder().setComplete(isComplete);
209+
if (isComplete) {
210+
InputStream stream =
211+
new ByteArrayInputStream(responseBody.getBytes(StandardCharsets.UTF_8));
212+
chunkResponseBuilder.setResponse(responseParser.parse(stream));
213+
}
214+
future.set(chunkResponseBuilder.build());
215+
} else {
216+
Throwable cause = trailers.getException();
217+
future.setException(
218+
cause != null
219+
? cause
220+
: new HttpJsonStatusRuntimeException(
221+
statusCode, "Failed to upload chunk with status code: " + statusCode, null));
222+
}
223+
} catch (Throwable t) {
224+
future.setException(t);
225+
}
226+
}
227+
}
228+
}

0 commit comments

Comments
 (0)