Skip to content

Commit d91c682

Browse files
committed
feat(gax): implement queryStatus in HttpJsonResumableUploadClient
1 parent 07159cd commit d91c682

5 files changed

Lines changed: 487 additions & 16 deletions

File tree

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

Lines changed: 151 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,8 @@
3535
import com.google.api.core.SettableApiFuture;
3636
import com.google.api.gax.resumable.ChunkUploadRequest;
3737
import com.google.api.gax.resumable.ChunkUploadResponse;
38+
import com.google.api.gax.resumable.QueryStatusRequest;
39+
import com.google.api.gax.resumable.QueryStatusResponse;
3840
import com.google.api.gax.resumable.ResumableUploadClient;
3941
import com.google.api.gax.resumable.ResumableUploadSession;
4042
import com.google.api.gax.rpc.ApiCallContext;
@@ -89,6 +91,9 @@ public final class HttpJsonResumableUploadClient<RequestT, ResponseT>
8991

9092
private static final PathTemplate PATH_TEMPLATE = PathTemplate.create("{+path}");
9193

94+
private static final Map<String, List<String>> QUERY_STATUS_HEADERS =
95+
ImmutableMap.of(UPLOAD_COMMAND_HEADER, ImmutableList.of("query"));
96+
9297
private static final ApiMethodDescriptor<ChunkUploadRequest, String> UPLOAD_CHUNK_DESCRIPTOR =
9398
ApiMethodDescriptor.<ChunkUploadRequest, String>newBuilder()
9499
.setFullMethodName("ResumableUpload/UploadChunk")
@@ -118,6 +123,36 @@ public PathTemplate getPathTemplate() {
118123
})
119124
.setResponseParser(ResumableUploadResponseParser.create())
120125
.build();
126+
127+
private static final ApiMethodDescriptor<QueryStatusRequest, String> QUERY_STATUS_DESCRIPTOR =
128+
ApiMethodDescriptor.<QueryStatusRequest, String>newBuilder()
129+
.setFullMethodName("ResumableUpload/QueryStatus")
130+
.setHttpMethod(HttpMethods.POST)
131+
.setType(ApiMethodDescriptor.MethodType.UNARY)
132+
.setRequestFormatter(
133+
new HttpRequestFormatter<QueryStatusRequest>() {
134+
@Override
135+
public Map<String, List<String>> getQueryParamNames(QueryStatusRequest request) {
136+
return Collections.emptyMap();
137+
}
138+
139+
@Override
140+
public String getRequestBody(QueryStatusRequest request) {
141+
return "";
142+
}
143+
144+
@Override
145+
public String getPath(QueryStatusRequest request) {
146+
return request.getUploadUrl();
147+
}
148+
149+
@Override
150+
public PathTemplate getPathTemplate() {
151+
return PATH_TEMPLATE;
152+
}
153+
})
154+
.setResponseParser(ResumableUploadResponseParser.create())
155+
.build();
121156
private final ClientContext clientContext;
122157
private final ApiMethodDescriptor<RequestT, String> startUploadDescriptor;
123158
private final HttpResponseParser<ResponseT> responseParser;
@@ -212,6 +247,35 @@ public ApiFuture<ChunkUploadResponse<ResponseT>> futureCall(
212247
};
213248
}
214249

250+
@Override
251+
public UnaryCallable<QueryStatusRequest, QueryStatusResponse<ResponseT>> queryStatusCallable() {
252+
return new UnaryCallable<QueryStatusRequest, QueryStatusResponse<ResponseT>>() {
253+
@Override
254+
public ApiFuture<QueryStatusResponse<ResponseT>> futureCall(
255+
QueryStatusRequest request, @Nullable ApiCallContext inputContext) {
256+
Preconditions.checkNotNull(request);
257+
HttpJsonCallContext context =
258+
(HttpJsonCallContext)
259+
HttpJsonCallContext.createDefault()
260+
.nullToSelf(clientContext.getDefaultCallContext())
261+
.merge(inputContext)
262+
.withExtraHeaders(QUERY_STATUS_HEADERS);
263+
264+
HttpJsonClientCall<QueryStatusRequest, String> clientCall =
265+
HttpJsonClientCalls.newCall(QUERY_STATUS_DESCRIPTOR, context);
266+
267+
SettableApiFuture<QueryStatusResponse<ResponseT>> future = SettableApiFuture.create();
268+
HttpJsonClientCalls.startUnaryCall(
269+
clientCall,
270+
request,
271+
context,
272+
new QueryStatusResponseListener<>(future, responseParser));
273+
274+
return future;
275+
}
276+
};
277+
}
278+
215279
private static class StartUploadResponseListener extends HttpJsonClientCall.Listener<String> {
216280

217281
private final SettableApiFuture<ResumableUploadSession> future;
@@ -300,25 +364,15 @@ private static class ChunkUploadResponseListener<ResponseT>
300364

301365
@Override
302366
public void onHeaders(HttpJsonMetadata responseHeaders) {
303-
Map<String, Object> headers = responseHeaders.getHeaders();
304-
305-
String statusStr = HttpHeadersUtils.getSingleHeader(headers, UPLOAD_STATUS_HEADER);
367+
String statusStr =
368+
HttpHeadersUtils.getSingleHeader(responseHeaders.getHeaders(), UPLOAD_STATUS_HEADER);
306369
if (statusStr != null) {
307370
this.hasUploadStatusHeader = true;
308-
if (STATUS_FINAL.equalsIgnoreCase(statusStr)) {
309-
this.isComplete = true;
310-
}
371+
this.isComplete = STATUS_FINAL.equalsIgnoreCase(statusStr);
311372
}
312-
313-
String sizeReceivedStr =
314-
HttpHeadersUtils.getSingleHeader(headers, UPLOAD_SIZE_RECEIVED_HEADER);
315-
if (!Strings.isNullOrEmpty(sizeReceivedStr)) {
316-
try {
317-
this.committedOffset = Long.parseLong(sizeReceivedStr);
318-
} catch (NumberFormatException ignored) {
319-
// Ignore invalid/malformed size received header and fall back to local offset
320-
// calculation.
321-
}
373+
Long sizeReceived = parseSizeReceived(responseHeaders);
374+
if (sizeReceived != null) {
375+
this.committedOffset = sizeReceived;
322376
}
323377
}
324378

@@ -367,12 +421,93 @@ public void onClose(int statusCode, HttpJsonMetadata trailers) {
367421
}
368422
}
369423

424+
private static class QueryStatusResponseListener<ResponseT>
425+
extends HttpJsonClientCall.Listener<String> {
426+
427+
private final SettableApiFuture<QueryStatusResponse<ResponseT>> future;
428+
private final HttpResponseParser<ResponseT> responseParser;
429+
private boolean isComplete = false;
430+
@Nullable private Long committedOffset = null;
431+
private String responseBody = "";
432+
433+
QueryStatusResponseListener(
434+
SettableApiFuture<QueryStatusResponse<ResponseT>> future,
435+
HttpResponseParser<ResponseT> responseParser) {
436+
this.future = future;
437+
this.responseParser = responseParser;
438+
}
439+
440+
@Override
441+
public void onHeaders(HttpJsonMetadata responseHeaders) {
442+
this.isComplete = isUploadFinal(responseHeaders);
443+
this.committedOffset = parseSizeReceived(responseHeaders);
444+
}
445+
446+
@Override
447+
public void onMessage(@Nullable String message) {
448+
if (message != null) {
449+
this.responseBody = message;
450+
}
451+
}
452+
453+
@Override
454+
public void onClose(int statusCode, HttpJsonMetadata trailers) {
455+
try {
456+
if (statusCode >= 200 && statusCode < 300) {
457+
if (isComplete || committedOffset != null) {
458+
ResponseT response = null;
459+
if (isComplete) {
460+
InputStream stream =
461+
new ByteArrayInputStream(responseBody.getBytes(StandardCharsets.UTF_8));
462+
response = responseParser.parse(stream);
463+
}
464+
future.set(
465+
QueryStatusResponse.create(
466+
committedOffset != null ? committedOffset : 0L, isComplete, response));
467+
} else {
468+
future.setException(
469+
ApiExceptionFactory.createException(
470+
"Query status response did not contain valid X-Goog-Upload-Size-Received"
471+
+ " header",
472+
/* cause= */ null,
473+
HttpJsonStatusCode.of(StatusCode.Code.INTERNAL),
474+
/* retryable= */ false));
475+
}
476+
} else {
477+
future.setException(
478+
createApiException(statusCode, trailers, "Failed to query upload status"));
479+
}
480+
} catch (Exception e) {
481+
future.setException(
482+
ApiExceptionFactory.createException(
483+
"Internal error processing query status response",
484+
e,
485+
HttpJsonStatusCode.of(StatusCode.Code.INTERNAL),
486+
/* retryable= */ false));
487+
}
488+
}
489+
}
490+
370491
private static boolean isUploadFinal(HttpJsonMetadata responseHeaders) {
371492
String statusStr =
372493
HttpHeadersUtils.getSingleHeader(responseHeaders.getHeaders(), UPLOAD_STATUS_HEADER);
373494
return STATUS_FINAL.equalsIgnoreCase(statusStr);
374495
}
375496

497+
@Nullable
498+
private static Long parseSizeReceived(HttpJsonMetadata responseHeaders) {
499+
String sizeReceivedStr =
500+
HttpHeadersUtils.getSingleHeader(responseHeaders.getHeaders(), UPLOAD_SIZE_RECEIVED_HEADER);
501+
if (!Strings.isNullOrEmpty(sizeReceivedStr)) {
502+
try {
503+
return Long.parseLong(sizeReceivedStr);
504+
} catch (NumberFormatException ignored) {
505+
// Unparseable header; return null and let the listener decide how to handle it.
506+
}
507+
}
508+
return null;
509+
}
510+
376511
private static ApiException createApiException(
377512
int statusCode, @Nullable HttpJsonMetadata trailers, String actionDescription) {
378513
Throwable cause = trailers != null ? trailers.getException() : null;

0 commit comments

Comments
 (0)