3636import com .google .api .core .InternalApi ;
3737import com .google .api .gax .resumable .ChunkUploadRequest ;
3838import com .google .api .gax .resumable .ChunkUploadResponse ;
39+ import com .google .api .gax .resumable .QueryStatusRequest ;
40+ import com .google .api .gax .resumable .QueryStatusResponse ;
3941import com .google .api .gax .resumable .ResumableUploadClient ;
4042import com .google .api .gax .resumable .ResumableUploadSession ;
4143import com .google .api .gax .rpc .ApiCallContext ;
@@ -87,6 +89,9 @@ public final class HttpJsonResumableUploadClient<RequestT, ResponseT>
8789
8890 private static final PathTemplate PATH_TEMPLATE = PathTemplate .create ("{+path}" );
8991
92+ private static final Map <String , List <String >> QUERY_STATUS_HEADERS =
93+ ImmutableMap .of (UPLOAD_COMMAND_HEADER , ImmutableList .of ("query" ));
94+
9095 private static final ApiMethodDescriptor <ChunkUploadRequest , String > UPLOAD_CHUNK_DESCRIPTOR =
9196 ApiMethodDescriptor .<ChunkUploadRequest , String >newBuilder ()
9297 .setFullMethodName ("ResumableUpload/UploadChunk" )
@@ -117,10 +122,42 @@ public PathTemplate getPathTemplate() {
117122 .setResponseParser (ResumableUploadResponseParser .create ())
118123 .build ();
119124
125+ private static final ApiMethodDescriptor <QueryStatusRequest , String > QUERY_STATUS_DESCRIPTOR =
126+ ApiMethodDescriptor .<QueryStatusRequest , String >newBuilder ()
127+ .setFullMethodName ("ResumableUpload/QueryStatus" )
128+ .setHttpMethod (HttpMethods .POST )
129+ .setType (ApiMethodDescriptor .MethodType .UNARY )
130+ .setRequestFormatter (
131+ new HttpRequestFormatter <QueryStatusRequest >() {
132+ @ Override
133+ public Map <String , List <String >> getQueryParamNames (QueryStatusRequest request ) {
134+ return Collections .emptyMap ();
135+ }
136+
137+ @ Override
138+ public String getRequestBody (QueryStatusRequest request ) {
139+ return "" ;
140+ }
141+
142+ @ Override
143+ public String getPath (QueryStatusRequest request ) {
144+ return request .getUploadUrl ();
145+ }
146+
147+ @ Override
148+ public PathTemplate getPathTemplate () {
149+ return PATH_TEMPLATE ;
150+ }
151+ })
152+ .setResponseParser (ResumableUploadResponseParser .create ())
153+ .build ();
154+
120155 private final ApiMethodDescriptor <RequestT , String > startUploadDescriptor ;
121156 private final UnaryCallable <RequestT , ResumableUploadSession > startUploadCallable ;
122157 private final UnaryCallable <ChunkUploadRequest , ChunkUploadResponse <ResponseT >>
123158 uploadChunkCallable ;
159+ private final UnaryCallable <QueryStatusRequest , QueryStatusResponse <ResponseT >>
160+ queryStatusCallable ;
124161
125162 public static <RequestT , ResponseT > HttpJsonResumableUploadClient <RequestT , ResponseT > create (
126163 ClientContext clientContext , ApiMethodDescriptor <RequestT , ResponseT > methodDescriptor ) {
@@ -144,6 +181,7 @@ private HttpJsonResumableUploadClient(
144181 .build ();
145182 this .startUploadCallable = createStartUploadCallable (clientContext );
146183 this .uploadChunkCallable = createUploadChunkCallable (clientContext , responseParser );
184+ this .queryStatusCallable = createQueryStatusCallable (clientContext , responseParser );
147185 }
148186
149187 @ Override
@@ -156,6 +194,11 @@ public UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>> uploadC
156194 return uploadChunkCallable ;
157195 }
158196
197+ @ Override
198+ public UnaryCallable <QueryStatusRequest , QueryStatusResponse <ResponseT >> queryStatusCallable () {
199+ return queryStatusCallable ;
200+ }
201+
159202 private UnaryCallable <RequestT , ResumableUploadSession > createStartUploadCallable (
160203 ClientContext clientContext ) {
161204 UnaryCallable <RequestT , ResumableUploadSession > rawCallable =
@@ -224,6 +267,35 @@ public ApiFuture<ChunkUploadResponse<ResponseT>> futureCall(
224267 return createClientCallable (rawCallable , clientContext );
225268 }
226269
270+ private UnaryCallable <QueryStatusRequest , QueryStatusResponse <ResponseT >>
271+ createQueryStatusCallable (
272+ ClientContext clientContext , HttpResponseParser <ResponseT > responseParser ) {
273+ UnaryCallable <QueryStatusRequest , QueryStatusResponse <ResponseT >> rawCallable =
274+ new UnaryCallable <QueryStatusRequest , QueryStatusResponse <ResponseT >>() {
275+ @ Override
276+ public ApiFuture <QueryStatusResponse <ResponseT >> futureCall (
277+ QueryStatusRequest request , @ Nullable ApiCallContext inputContext ) {
278+ Preconditions .checkNotNull (request );
279+ HttpJsonCallContext context =
280+ createCallContext (clientContext , inputContext , QUERY_STATUS_HEADERS );
281+
282+ HttpJsonClientCall <QueryStatusRequest , String > clientCall =
283+ HttpJsonClientCalls .newCall (QUERY_STATUS_DESCRIPTOR , context );
284+
285+ HttpJsonCallFuture <QueryStatusResponse <ResponseT >> future =
286+ new HttpJsonCallFuture <>(clientCall );
287+ HttpJsonClientCalls .startUnaryCall (
288+ clientCall ,
289+ request ,
290+ context ,
291+ new QueryStatusResponseListener <>(future , responseParser ));
292+
293+ return future ;
294+ }
295+ };
296+ return createClientCallable (rawCallable , clientContext );
297+ }
298+
227299 private static HttpJsonCallContext createCallContext (
228300 ClientContext clientContext ,
229301 @ Nullable ApiCallContext inputContext ,
@@ -245,6 +317,22 @@ private static <CallReqT, CallRespT> UnaryCallable<CallReqT, CallRespT> createCl
245317 return callable .withDefaultCallContext (clientContext .getDefaultCallContext ());
246318 }
247319
320+ @ Nullable
321+ private static <ResponseT > ResponseT parseResponseBody (
322+ String responseBody , HttpResponseParser <ResponseT > responseParser ) {
323+ InputStream stream = new ByteArrayInputStream (responseBody .getBytes (StandardCharsets .UTF_8 ));
324+ return responseParser .parse (stream );
325+ }
326+
327+ @ Nullable
328+ private static String getUploadStatus (HttpJsonMetadata responseHeaders ) {
329+ return HttpHeadersUtils .getSingleHeader (responseHeaders .getHeaders (), UPLOAD_STATUS_HEADER );
330+ }
331+
332+ private static boolean isUploadFinal (HttpJsonMetadata responseHeaders ) {
333+ return STATUS_FINAL .equalsIgnoreCase (getUploadStatus (responseHeaders ));
334+ }
335+
248336 @ Nullable
249337 private static Long parseSizeReceived (HttpJsonMetadata responseHeaders ) {
250338 String sizeReceivedStr =
@@ -259,6 +347,15 @@ private static Long parseSizeReceived(HttpJsonMetadata responseHeaders) {
259347 return null ;
260348 }
261349
350+ private static Throwable createStatusException (
351+ int statusCode , HttpJsonMetadata trailers , String actionMessage ) {
352+ Throwable cause = trailers .getException ();
353+ return cause != null
354+ ? cause
355+ : new HttpJsonStatusRuntimeException (
356+ statusCode , actionMessage + " with status code: " + statusCode , null );
357+ }
358+
262359 /**
263360 * An {@link ApiFuture} that cancels the underlying {@link HttpJsonClientCall} to prevent
264361 * connection leaks.
@@ -402,16 +499,11 @@ private static class ChunkUploadResponseListener<ResponseT>
402499
403500 @ Override
404501 public void onHeaders (HttpJsonMetadata responseHeaders ) {
405- Map <String , Object > headers = responseHeaders .getHeaders ();
406-
407- String statusStr = HttpHeadersUtils .getSingleHeader (headers , UPLOAD_STATUS_HEADER );
502+ String statusStr = getUploadStatus (responseHeaders );
408503 if (statusStr != null ) {
409504 this .hasUploadStatusHeader = true ;
410- if (STATUS_FINAL .equalsIgnoreCase (statusStr )) {
411- this .isComplete = true ;
412- }
505+ this .isComplete = STATUS_FINAL .equalsIgnoreCase (statusStr );
413506 }
414-
415507 this .committedOffset = parseSizeReceived (responseHeaders );
416508 }
417509
@@ -439,20 +531,70 @@ public void onClose(int statusCode, HttpJsonMetadata trailers) {
439531 committedOffset != null
440532 ? committedOffset
441533 : request .getOffset () + request .getPayload ().size ();
442- ResponseT response = null ;
443- if (isComplete ) {
444- InputStream stream =
445- new ByteArrayInputStream (responseBody .getBytes (StandardCharsets .UTF_8 ));
446- response = responseParser .parse (stream );
447- }
534+ ResponseT response = isComplete ? parseResponseBody (responseBody , responseParser ) : null ;
448535 future .set (ChunkUploadResponse .create (confirmedOffset , isComplete , response ));
449536 } else {
450- Throwable cause = trailers .getException ();
451537 future .setException (
452- cause != null
453- ? cause
454- : new HttpJsonStatusRuntimeException (
455- statusCode , "Failed to upload chunk with status code: " + statusCode , null ));
538+ createStatusException (statusCode , trailers , "Failed to upload chunk" ));
539+ }
540+ } catch (Throwable t ) {
541+ future .setException (t );
542+ }
543+ }
544+ }
545+
546+ /** A listener that parses query response headers to produce the {@link QueryStatusResponse}. */
547+ private static class QueryStatusResponseListener <ResponseT >
548+ extends HttpJsonClientCall .Listener <String > {
549+
550+ private final HttpJsonCallFuture <QueryStatusResponse <ResponseT >> future ;
551+ private final HttpResponseParser <ResponseT > responseParser ;
552+ private boolean isComplete = false ;
553+ @ Nullable private Long committedOffset = null ;
554+ private String responseBody = "" ;
555+
556+ QueryStatusResponseListener (
557+ HttpJsonCallFuture <QueryStatusResponse <ResponseT >> future ,
558+ HttpResponseParser <ResponseT > responseParser ) {
559+ this .future = future ;
560+ this .responseParser = responseParser ;
561+ }
562+
563+ @ Override
564+ public void onHeaders (HttpJsonMetadata responseHeaders ) {
565+ this .isComplete = isUploadFinal (responseHeaders );
566+ this .committedOffset = parseSizeReceived (responseHeaders );
567+ }
568+
569+ @ Override
570+ public void onMessage (@ Nullable String message ) {
571+ if (message != null ) {
572+ this .responseBody = message ;
573+ }
574+ }
575+
576+ @ Override
577+ public void onClose (int statusCode , HttpJsonMetadata trailers ) {
578+ try {
579+ if (statusCode >= 200 && statusCode < 300 ) {
580+ if (isComplete || committedOffset != null ) {
581+ ResponseT response =
582+ isComplete ? parseResponseBody (responseBody , responseParser ) : null ;
583+ future .set (
584+ QueryStatusResponse .create (
585+ committedOffset != null ? committedOffset : 0L , isComplete , response ));
586+ } else {
587+ future .setException (
588+ ApiExceptionFactory .createException (
589+ "Query status response did not contain valid X-Goog-Upload-Size-Received"
590+ + " header" ,
591+ /* cause= */ null ,
592+ HttpJsonStatusCode .of (StatusCode .Code .INTERNAL ),
593+ /* retryable= */ false ));
594+ }
595+ } else {
596+ future .setException (
597+ createStatusException (statusCode , trailers , "Failed to query upload status" ));
456598 }
457599 } catch (Throwable t ) {
458600 future .setException (t );
0 commit comments