Skip to content

Commit f8a42a5

Browse files
committed
feat(gax): implement uploadChunk in HttpJsonResumableUploadClient
1 parent 3655e2e commit f8a42a5

7 files changed

Lines changed: 699 additions & 0 deletions

File tree

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

Lines changed: 155 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,10 +29,15 @@
2929
*/
3030
package com.google.api.gax.httpjson;
3131

32+
import com.google.api.client.http.ByteArrayContent;
33+
import com.google.api.client.http.EmptyContent;
34+
import com.google.api.client.http.HttpContent;
3235
import com.google.api.client.http.HttpMethods;
3336
import com.google.api.core.ApiFuture;
3437
import com.google.api.core.InternalApi;
3538
import com.google.api.core.SettableApiFuture;
39+
import com.google.api.gax.resumable.ChunkUploadRequest;
40+
import com.google.api.gax.resumable.ChunkUploadResponse;
3641
import com.google.api.gax.resumable.ResumableUploadClient;
3742
import com.google.api.gax.resumable.ResumableUploadSession;
3843
import com.google.api.gax.resumable.StartUploadRequest;
@@ -67,8 +72,12 @@ public final class HttpJsonResumableUploadClient implements ResumableUploadClien
6772

6873
private static final String UPLOAD_PROTOCOL_HEADER = "X-Goog-Upload-Protocol";
6974
private static final String UPLOAD_COMMAND_HEADER = "X-Goog-Upload-Command";
75+
private static final String UPLOAD_OFFSET_HEADER = "X-Goog-Upload-Offset";
7076
private static final String UPLOAD_URL_HEADER = "X-Goog-Upload-URL";
7177
private static final String UPLOAD_GRANULARITY_HEADER = "X-Goog-Upload-Chunk-Granularity";
78+
private static final String UPLOAD_STATUS_HEADER = "X-Goog-Upload-Status";
79+
private static final String UPLOAD_SIZE_RECEIVED_HEADER = "X-Goog-Upload-Size-Received";
80+
private static final String STATUS_FINAL = "final";
7281

7382
private static final Map<String, List<String>> START_UPLOAD_HEADERS =
7483
ImmutableMap.of(
@@ -105,6 +114,45 @@ public PathTemplate getPathTemplate() {
105114
.setResponseParser(StringHttpResponseParser.create())
106115
.build();
107116

117+
private static final ApiMethodDescriptor<ChunkUploadRequest, String> UPLOAD_CHUNK_DESCRIPTOR =
118+
ApiMethodDescriptor.<ChunkUploadRequest, String>newBuilder()
119+
.setFullMethodName("ResumableUpload/UploadChunk")
120+
.setHttpMethod(HttpMethods.POST)
121+
.setType(ApiMethodDescriptor.MethodType.UNARY)
122+
.setRequestFormatter(
123+
new HttpRequestFormatter<ChunkUploadRequest>() {
124+
@Override
125+
public Map<String, List<String>> getQueryParamNames(ChunkUploadRequest request) {
126+
return Collections.emptyMap();
127+
}
128+
129+
@Override
130+
public String getRequestBody(ChunkUploadRequest request) {
131+
return "";
132+
}
133+
134+
@Override
135+
public HttpContent getHttpContent(ChunkUploadRequest request) {
136+
if (request.getPayload() != null && !request.getPayload().isEmpty()) {
137+
return new ByteArrayContent(
138+
"application/octet-stream", request.getPayload().toByteArray());
139+
}
140+
return new EmptyContent();
141+
}
142+
143+
@Override
144+
public String getPath(ChunkUploadRequest request) {
145+
return request.getUploadUrl();
146+
}
147+
148+
@Override
149+
public PathTemplate getPathTemplate() {
150+
return PathTemplate.create("{+path}");
151+
}
152+
})
153+
.setResponseParser(StringHttpResponseParser.create())
154+
.build();
155+
108156
private final ClientContext clientContext;
109157

110158
public static HttpJsonResumableUploadClient create(ClientContext clientContext) {
@@ -141,6 +189,48 @@ public ApiFuture<ResumableUploadSession> futureCall(
141189
};
142190
}
143191

192+
@Override
193+
public UnaryCallable<ChunkUploadRequest, ChunkUploadResponse> uploadChunkCallable() {
194+
return new UnaryCallable<ChunkUploadRequest, ChunkUploadResponse>() {
195+
@Override
196+
public ApiFuture<ChunkUploadResponse> futureCall(
197+
ChunkUploadRequest request, @Nullable ApiCallContext inputContext) {
198+
Preconditions.checkNotNull(request);
199+
String command;
200+
if (request.isFinal()) {
201+
command =
202+
(request.getPayload() != null && !request.getPayload().isEmpty())
203+
? "upload, finalize"
204+
: "finalize";
205+
} else {
206+
command = "upload";
207+
}
208+
Map<String, List<String>> chunkHeaders =
209+
ImmutableMap.of(
210+
UPLOAD_COMMAND_HEADER,
211+
ImmutableList.of(command),
212+
UPLOAD_OFFSET_HEADER,
213+
ImmutableList.of(String.valueOf(request.getOffset())));
214+
215+
HttpJsonCallContext context =
216+
(HttpJsonCallContext)
217+
HttpJsonCallContext.createDefault()
218+
.nullToSelf(clientContext.getDefaultCallContext())
219+
.merge(inputContext)
220+
.withExtraHeaders(chunkHeaders);
221+
222+
HttpJsonClientCall<ChunkUploadRequest, String> clientCall =
223+
HttpJsonClientCalls.newCall(UPLOAD_CHUNK_DESCRIPTOR, context);
224+
225+
SettableApiFuture<ChunkUploadResponse> future = SettableApiFuture.create();
226+
HttpJsonClientCalls.startUnaryCall(
227+
clientCall, request, context, new ChunkUploadResponseListener(request, future));
228+
229+
return future;
230+
}
231+
};
232+
}
233+
144234
private static class StartUploadResponseListener extends HttpJsonClientCall.Listener<String> {
145235

146236
private final SettableApiFuture<ResumableUploadSession> future;
@@ -205,4 +295,69 @@ public void onClose(int statusCode, HttpJsonMetadata trailers) {
205295
}
206296
}
207297
}
298+
299+
private static class ChunkUploadResponseListener extends HttpJsonClientCall.Listener<String> {
300+
301+
private final ChunkUploadRequest request;
302+
private final SettableApiFuture<ChunkUploadResponse> future;
303+
private boolean isComplete = false;
304+
private long committedOffset = -1L;
305+
@Nullable private String responseBody;
306+
307+
ChunkUploadResponseListener(
308+
ChunkUploadRequest request, SettableApiFuture<ChunkUploadResponse> future) {
309+
this.request = request;
310+
this.future = future;
311+
}
312+
313+
@Override
314+
public void onHeaders(HttpJsonMetadata responseHeaders) {
315+
Map<String, Object> headers = responseHeaders.getHeaders();
316+
317+
String statusStr = HttpHeadersUtils.getFirstHeader(headers, UPLOAD_STATUS_HEADER);
318+
if (STATUS_FINAL.equalsIgnoreCase(statusStr)) {
319+
this.isComplete = true;
320+
}
321+
322+
String sizeReceivedStr =
323+
HttpHeadersUtils.getFirstHeader(headers, UPLOAD_SIZE_RECEIVED_HEADER);
324+
if (!Strings.isNullOrEmpty(sizeReceivedStr)) {
325+
try {
326+
this.committedOffset = Long.parseLong(sizeReceivedStr);
327+
} catch (NumberFormatException ignored) {
328+
}
329+
}
330+
}
331+
332+
@Override
333+
public void onMessage(@Nullable String message) {
334+
if (message != null && !message.isEmpty()) {
335+
this.responseBody = message;
336+
}
337+
}
338+
339+
@Override
340+
public void onClose(int statusCode, HttpJsonMetadata trailers) {
341+
if ((statusCode >= 200 && statusCode < 300) || statusCode == 308) {
342+
long confirmedOffset =
343+
committedOffset >= 0
344+
? committedOffset
345+
: request.getOffset() + request.getPayload().size();
346+
future.set(
347+
ChunkUploadResponse.create(
348+
confirmedOffset, isComplete, isComplete ? responseBody : null));
349+
} else {
350+
Throwable cause = trailers.getException();
351+
ApiException apiException =
352+
cause != null
353+
? API_EXCEPTION_FACTORY.create(cause)
354+
: ApiExceptionFactory.createException(
355+
"Failed to upload chunk with status code: " + statusCode,
356+
/* cause= */ null,
357+
HttpJsonStatusCode.of(statusCode),
358+
/* retryable= */ false);
359+
future.setException(apiException);
360+
}
361+
}
362+
}
208363
}

0 commit comments

Comments
 (0)