Skip to content

Commit 44dc342

Browse files
committed
feat(gax): implement queryStatus in HttpJsonResumableUploadClient
1 parent d1a4963 commit 44dc342

5 files changed

Lines changed: 483 additions & 16 deletions

File tree

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

Lines changed: 147 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,8 @@
3838
import com.google.api.core.SettableApiFuture;
3939
import com.google.api.gax.resumable.ChunkUploadRequest;
4040
import com.google.api.gax.resumable.ChunkUploadResponse;
41+
import com.google.api.gax.resumable.QueryStatusRequest;
42+
import com.google.api.gax.resumable.QueryStatusResponse;
4143
import com.google.api.gax.resumable.ResumableUploadClient;
4244
import com.google.api.gax.resumable.ResumableUploadSession;
4345
import com.google.api.gax.resumable.StartUploadRequest;
@@ -84,6 +86,9 @@ public final class HttpJsonResumableUploadClient implements ResumableUploadClien
8486
UPLOAD_PROTOCOL_HEADER, ImmutableList.of("resumable"),
8587
UPLOAD_COMMAND_HEADER, ImmutableList.of("start"));
8688

89+
private static final Map<String, List<String>> QUERY_STATUS_HEADERS =
90+
ImmutableMap.of(UPLOAD_COMMAND_HEADER, ImmutableList.of("query"));
91+
8792
private static final ApiMethodDescriptor<StartUploadRequest, String> START_UPLOAD_DESCRIPTOR =
8893
ApiMethodDescriptor.<StartUploadRequest, String>newBuilder()
8994
.setFullMethodName("ResumableUpload/StartUpload")
@@ -153,6 +158,41 @@ public PathTemplate getPathTemplate() {
153158
.setResponseParser(StringHttpResponseParser.create())
154159
.build();
155160

161+
private static final ApiMethodDescriptor<QueryStatusRequest, String> QUERY_STATUS_DESCRIPTOR =
162+
ApiMethodDescriptor.<QueryStatusRequest, String>newBuilder()
163+
.setFullMethodName("ResumableUpload/QueryStatus")
164+
.setHttpMethod(HttpMethods.POST)
165+
.setType(ApiMethodDescriptor.MethodType.UNARY)
166+
.setRequestFormatter(
167+
new HttpRequestFormatter<QueryStatusRequest>() {
168+
@Override
169+
public Map<String, List<String>> getQueryParamNames(QueryStatusRequest request) {
170+
return Collections.emptyMap();
171+
}
172+
173+
@Override
174+
public String getRequestBody(QueryStatusRequest request) {
175+
return "";
176+
}
177+
178+
@Override
179+
public HttpContent getHttpContent(QueryStatusRequest request) {
180+
return new EmptyContent();
181+
}
182+
183+
@Override
184+
public String getPath(QueryStatusRequest request) {
185+
return request.getUploadUrl();
186+
}
187+
188+
@Override
189+
public PathTemplate getPathTemplate() {
190+
return PathTemplate.create("{+path}");
191+
}
192+
})
193+
.setResponseParser(StringHttpResponseParser.create())
194+
.build();
195+
156196
private final ClientContext clientContext;
157197

158198
public static HttpJsonResumableUploadClient create(ClientContext clientContext) {
@@ -228,6 +268,32 @@ public ApiFuture<ChunkUploadResponse> futureCall(
228268
};
229269
}
230270

271+
@Override
272+
public UnaryCallable<QueryStatusRequest, QueryStatusResponse> queryStatusCallable() {
273+
return new UnaryCallable<QueryStatusRequest, QueryStatusResponse>() {
274+
@Override
275+
public ApiFuture<QueryStatusResponse> futureCall(
276+
QueryStatusRequest request, @Nullable ApiCallContext inputContext) {
277+
Preconditions.checkNotNull(request);
278+
HttpJsonCallContext context =
279+
(HttpJsonCallContext)
280+
HttpJsonCallContext.createDefault()
281+
.nullToSelf(clientContext.getDefaultCallContext())
282+
.merge(inputContext)
283+
.withExtraHeaders(QUERY_STATUS_HEADERS);
284+
285+
HttpJsonClientCall<QueryStatusRequest, String> clientCall =
286+
HttpJsonClientCalls.newCall(QUERY_STATUS_DESCRIPTOR, context);
287+
288+
SettableApiFuture<QueryStatusResponse> future = SettableApiFuture.create();
289+
HttpJsonClientCalls.startUnaryCall(
290+
clientCall, request, context, new QueryStatusResponseListener(future));
291+
292+
return future;
293+
}
294+
};
295+
}
296+
231297
private static class StartUploadResponseListener extends HttpJsonClientCall.Listener<String> {
232298

233299
private final SettableApiFuture<ResumableUploadSession> future;
@@ -306,22 +372,10 @@ private static class ChunkUploadResponseListener extends HttpJsonClientCall.List
306372

307373
@Override
308374
public void onHeaders(HttpJsonMetadata responseHeaders) {
309-
Map<String, Object> headers = responseHeaders.getHeaders();
310-
311-
String statusStr = HttpHeadersUtils.getFirstHeader(headers, UPLOAD_STATUS_HEADER);
312-
if (STATUS_FINAL.equalsIgnoreCase(statusStr)) {
313-
this.isComplete = true;
314-
}
315-
316-
String sizeReceivedStr =
317-
HttpHeadersUtils.getFirstHeader(headers, UPLOAD_SIZE_RECEIVED_HEADER);
318-
if (!Strings.isNullOrEmpty(sizeReceivedStr)) {
319-
try {
320-
this.committedOffset = Long.parseLong(sizeReceivedStr);
321-
} catch (NumberFormatException ignored) {
322-
// Ignore invalid/malformed size received header and fall back to local offset
323-
// calculation.
324-
}
375+
this.isComplete = isUploadFinal(responseHeaders);
376+
Long sizeReceived = parseSizeReceived(responseHeaders);
377+
if (sizeReceived != null) {
378+
this.committedOffset = sizeReceived;
325379
}
326380
}
327381

@@ -357,6 +411,83 @@ public void onClose(int statusCode, HttpJsonMetadata trailers) {
357411
}
358412
}
359413

414+
private static class QueryStatusResponseListener extends HttpJsonClientCall.Listener<String> {
415+
416+
private final SettableApiFuture<QueryStatusResponse> future;
417+
private boolean isComplete = false;
418+
@Nullable private Long committedOffset = null;
419+
private String responseBody = "";
420+
421+
QueryStatusResponseListener(SettableApiFuture<QueryStatusResponse> future) {
422+
this.future = future;
423+
}
424+
425+
@Override
426+
public void onHeaders(HttpJsonMetadata responseHeaders) {
427+
this.isComplete = isUploadFinal(responseHeaders);
428+
this.committedOffset = parseSizeReceived(responseHeaders);
429+
}
430+
431+
@Override
432+
public void onMessage(@Nullable String message) {
433+
if (message != null) {
434+
this.responseBody = message;
435+
}
436+
}
437+
438+
@Override
439+
public void onClose(int statusCode, HttpJsonMetadata trailers) {
440+
try {
441+
if (statusCode >= 200 && statusCode < 300) {
442+
if (isComplete || committedOffset != null) {
443+
future.set(
444+
QueryStatusResponse.create(
445+
committedOffset != null ? committedOffset : 0L,
446+
isComplete,
447+
isComplete ? responseBody : ""));
448+
} else {
449+
future.setException(
450+
ApiExceptionFactory.createException(
451+
"Query status response did not contain valid X-Goog-Upload-Size-Received header",
452+
/* cause= */ null,
453+
HttpJsonStatusCode.of(StatusCode.Code.INTERNAL),
454+
/* retryable= */ false));
455+
}
456+
} else {
457+
future.setException(
458+
createApiException(statusCode, trailers, "Failed to query upload status"));
459+
}
460+
} catch (Throwable t) {
461+
future.setException(
462+
ApiExceptionFactory.createException(
463+
"Internal error processing query status response",
464+
t,
465+
HttpJsonStatusCode.of(StatusCode.Code.INTERNAL),
466+
/* retryable= */ false));
467+
}
468+
}
469+
}
470+
471+
private static boolean isUploadFinal(HttpJsonMetadata responseHeaders) {
472+
String statusStr =
473+
HttpHeadersUtils.getFirstHeader(responseHeaders.getHeaders(), UPLOAD_STATUS_HEADER);
474+
return STATUS_FINAL.equalsIgnoreCase(statusStr);
475+
}
476+
477+
@Nullable
478+
private static Long parseSizeReceived(HttpJsonMetadata responseHeaders) {
479+
String sizeReceivedStr =
480+
HttpHeadersUtils.getFirstHeader(responseHeaders.getHeaders(), UPLOAD_SIZE_RECEIVED_HEADER);
481+
if (!Strings.isNullOrEmpty(sizeReceivedStr)) {
482+
try {
483+
return Long.parseLong(sizeReceivedStr);
484+
} catch (NumberFormatException ignored) {
485+
// Unparseable header; return null and let the listener decide how to handle it.
486+
}
487+
}
488+
return null;
489+
}
490+
360491
private static ApiException createApiException(
361492
int statusCode, @Nullable HttpJsonMetadata trailers, String actionDescription) {
362493
Throwable cause = trailers != null ? trailers.getException() : null;

0 commit comments

Comments
 (0)