3434import com .google .api .core .ApiFuture ;
3535import com .google .api .core .BetaApi ;
3636import com .google .api .core .InternalApi ;
37+ import com .google .api .gax .resumable .ChunkUploadRequest ;
38+ import com .google .api .gax .resumable .ChunkUploadResponse ;
3739import com .google .api .gax .resumable .ResumableUploadClient ;
3840import com .google .api .gax .resumable .ResumableUploadSession ;
3941import com .google .api .gax .rpc .ApiCallContext ;
4042import com .google .api .gax .rpc .ApiExceptionFactory ;
4143import com .google .api .gax .rpc .ClientContext ;
4244import com .google .api .gax .rpc .StatusCode ;
4345import com .google .api .gax .rpc .UnaryCallable ;
46+ import com .google .api .pathtemplate .PathTemplate ;
4447import com .google .common .base .Preconditions ;
4548import com .google .common .base .Strings ;
4649import com .google .common .collect .ImmutableList ;
4750import com .google .common .collect .ImmutableMap ;
51+ import java .io .ByteArrayInputStream ;
52+ import java .io .InputStream ;
53+ import java .nio .charset .StandardCharsets ;
4854import java .util .Collections ;
4955import java .util .List ;
5056import java .util .Map ;
@@ -67,16 +73,54 @@ public final class HttpJsonResumableUploadClient<RequestT, ResponseT>
6773
6874 private static final String UPLOAD_PROTOCOL_HEADER = "X-Goog-Upload-Protocol" ;
6975 private static final String UPLOAD_COMMAND_HEADER = "X-Goog-Upload-Command" ;
76+ private static final String UPLOAD_OFFSET_HEADER = "X-Goog-Upload-Offset" ;
7077 private static final String UPLOAD_URL_HEADER = "X-Goog-Upload-URL" ;
7178 private static final String UPLOAD_GRANULARITY_HEADER = "X-Goog-Upload-Chunk-Granularity" ;
79+ private static final String UPLOAD_STATUS_HEADER = "X-Goog-Upload-Status" ;
80+ private static final String UPLOAD_SIZE_RECEIVED_HEADER = "X-Goog-Upload-Size-Received" ;
81+ private static final String STATUS_FINAL = "final" ;
7282
7383 private static final Map <String , List <String >> START_UPLOAD_HEADERS =
7484 ImmutableMap .of (
7585 UPLOAD_PROTOCOL_HEADER , ImmutableList .of ("resumable" ),
7686 UPLOAD_COMMAND_HEADER , ImmutableList .of ("start" ));
7787
88+ private static final PathTemplate PATH_TEMPLATE = PathTemplate .create ("{+path}" );
89+
90+ private static final ApiMethodDescriptor <ChunkUploadRequest , String > UPLOAD_CHUNK_DESCRIPTOR =
91+ ApiMethodDescriptor .<ChunkUploadRequest , String >newBuilder ()
92+ .setFullMethodName ("ResumableUpload/UploadChunk" )
93+ .setHttpMethod (HttpMethods .POST )
94+ .setType (ApiMethodDescriptor .MethodType .UNARY )
95+ .setRequestFormatter (
96+ new ResumableUploadChunkRequestFormatter <ChunkUploadRequest >() {
97+ @ Override
98+ public Map <String , List <String >> getQueryParamNames (ChunkUploadRequest request ) {
99+ return Collections .emptyMap ();
100+ }
101+
102+ @ Override
103+ public byte [] getBinaryRequestBody (ChunkUploadRequest request ) {
104+ return request .getPayload ().toByteArray ();
105+ }
106+
107+ @ Override
108+ public String getPath (ChunkUploadRequest request ) {
109+ return request .getUploadUrl ();
110+ }
111+
112+ @ Override
113+ public PathTemplate getPathTemplate () {
114+ return PATH_TEMPLATE ;
115+ }
116+ })
117+ .setResponseParser (ResumableUploadResponseParser .create ())
118+ .build ();
119+
78120 private final ApiMethodDescriptor <RequestT , String > startUploadDescriptor ;
79121 private final UnaryCallable <RequestT , ResumableUploadSession > startUploadCallable ;
122+ private final UnaryCallable <ChunkUploadRequest , ChunkUploadResponse <ResponseT >>
123+ uploadChunkCallable ;
80124
81125 public static <RequestT , ResponseT > HttpJsonResumableUploadClient <RequestT , ResponseT > create (
82126 ClientContext clientContext , ApiMethodDescriptor <RequestT , ResponseT > methodDescriptor ) {
@@ -87,6 +131,8 @@ private HttpJsonResumableUploadClient(
87131 ClientContext clientContext , ApiMethodDescriptor <RequestT , ResponseT > methodDescriptor ) {
88132 Preconditions .checkNotNull (clientContext );
89133 Preconditions .checkNotNull (methodDescriptor );
134+ HttpResponseParser <ResponseT > responseParser =
135+ Preconditions .checkNotNull (methodDescriptor .getResponseParser ());
90136
91137 this .startUploadDescriptor =
92138 ApiMethodDescriptor .<RequestT , String >newBuilder ()
@@ -97,13 +143,19 @@ private HttpJsonResumableUploadClient(
97143 .setResponseParser (ResumableUploadResponseParser .create ())
98144 .build ();
99145 this .startUploadCallable = createStartUploadCallable (clientContext );
146+ this .uploadChunkCallable = createUploadChunkCallable (clientContext , responseParser );
100147 }
101148
102149 @ Override
103150 public UnaryCallable <RequestT , ResumableUploadSession > startUploadCallable () {
104151 return startUploadCallable ;
105152 }
106153
154+ @ Override
155+ public UnaryCallable <ChunkUploadRequest , ChunkUploadResponse <ResponseT >> uploadChunkCallable () {
156+ return uploadChunkCallable ;
157+ }
158+
107159 private UnaryCallable <RequestT , ResumableUploadSession > createStartUploadCallable (
108160 ClientContext clientContext ) {
109161 UnaryCallable <RequestT , ResumableUploadSession > rawCallable =
@@ -129,6 +181,51 @@ public ApiFuture<ResumableUploadSession> futureCall(
129181 return createClientCallable (rawCallable , clientContext );
130182 }
131183
184+ private UnaryCallable <ChunkUploadRequest , ChunkUploadResponse <ResponseT >>
185+ createUploadChunkCallable (
186+ ClientContext clientContext , HttpResponseParser <ResponseT > responseParser ) {
187+ UnaryCallable <ChunkUploadRequest , ChunkUploadResponse <ResponseT >> rawCallable =
188+ new UnaryCallable <ChunkUploadRequest , ChunkUploadResponse <ResponseT >>() {
189+ @ Override
190+ public ApiFuture <ChunkUploadResponse <ResponseT >> futureCall (
191+ ChunkUploadRequest request , @ Nullable ApiCallContext inputContext ) {
192+ Preconditions .checkNotNull (request );
193+ boolean isPayloadEmpty = request .getPayload ().isEmpty ();
194+ String command ;
195+ if (request .isFinal ()) {
196+ command = !isPayloadEmpty ? "upload, finalize" : "finalize" ;
197+ } else {
198+ command = "upload" ;
199+ }
200+ ImmutableMap .Builder <String , List <String >> chunkHeadersBuilder =
201+ ImmutableMap .<String , List <String >>builder ()
202+ .put (UPLOAD_COMMAND_HEADER , ImmutableList .of (command ));
203+ if (!"finalize" .equals (command )) {
204+ chunkHeadersBuilder .put (
205+ UPLOAD_OFFSET_HEADER , ImmutableList .of (String .valueOf (request .getOffset ())));
206+ }
207+ Map <String , List <String >> chunkHeaders = chunkHeadersBuilder .build ();
208+
209+ HttpJsonCallContext context =
210+ createCallContext (clientContext , inputContext , chunkHeaders );
211+
212+ HttpJsonClientCall <ChunkUploadRequest , String > clientCall =
213+ HttpJsonClientCalls .newCall (UPLOAD_CHUNK_DESCRIPTOR , context );
214+
215+ HttpJsonCallFuture <ChunkUploadResponse <ResponseT >> future =
216+ new HttpJsonCallFuture <>(clientCall );
217+ HttpJsonClientCalls .startUnaryCall (
218+ clientCall ,
219+ request ,
220+ context ,
221+ new ChunkUploadResponseListener <>(request , future , responseParser ));
222+
223+ return future ;
224+ }
225+ };
226+ return createClientCallable (rawCallable , clientContext );
227+ }
228+
132229 private static HttpJsonCallContext createCallContext (
133230 ClientContext clientContext ,
134231 @ Nullable ApiCallContext inputContext ,
@@ -150,6 +247,20 @@ private static <CallReqT, CallRespT> UnaryCallable<CallReqT, CallRespT> createCl
150247 return callable .withDefaultCallContext (clientContext .getDefaultCallContext ());
151248 }
152249
250+ @ Nullable
251+ private static Long parseSizeReceived (HttpJsonMetadata responseHeaders ) {
252+ String sizeReceivedStr =
253+ HttpHeadersUtils .getSingleHeader (responseHeaders .getHeaders (), UPLOAD_SIZE_RECEIVED_HEADER );
254+ if (!Strings .isNullOrEmpty (sizeReceivedStr )) {
255+ try {
256+ return Long .parseLong (sizeReceivedStr );
257+ } catch (NumberFormatException ignored ) {
258+ // Unparseable header; return null and let the listener decide how to handle it.
259+ }
260+ }
261+ return null ;
262+ }
263+
153264 /**
154265 * An {@link ApiFuture} that cancels the underlying {@link HttpJsonClientCall} to prevent
155266 * connection leaks.
@@ -167,12 +278,10 @@ protected void interruptTask() {
167278 call .cancel ("Call was cancelled" , null );
168279 }
169280
170- @ Override
171281 public boolean set (T value ) {
172282 return super .set (value );
173283 }
174284
175- @ Override
176285 public boolean setException (Throwable throwable ) {
177286 return super .setException (throwable );
178287 }
@@ -268,4 +377,88 @@ public void onClose(int statusCode, HttpJsonMetadata trailers) {
268377 }
269378 }
270379 }
380+
381+ /**
382+ * A listener that processes chunk upload response headers and bodies to produce the {@link
383+ * ChunkUploadResponse}.
384+ */
385+ private static class ChunkUploadResponseListener <ResponseT >
386+ extends HttpJsonClientCall .Listener <String > {
387+
388+ private final ChunkUploadRequest request ;
389+ private final HttpJsonCallFuture <ChunkUploadResponse <ResponseT >> future ;
390+ private final HttpResponseParser <ResponseT > responseParser ;
391+ private boolean hasUploadStatusHeader = false ;
392+ private boolean isComplete = false ;
393+ @ Nullable private Long committedOffset = null ;
394+ private String responseBody = "" ;
395+
396+ ChunkUploadResponseListener (
397+ ChunkUploadRequest request ,
398+ HttpJsonCallFuture <ChunkUploadResponse <ResponseT >> future ,
399+ HttpResponseParser <ResponseT > responseParser ) {
400+ this .request = request ;
401+ this .future = future ;
402+ this .responseParser = responseParser ;
403+ }
404+
405+ @ Override
406+ public void onHeaders (HttpJsonMetadata responseHeaders ) {
407+ Map <String , Object > headers = responseHeaders .getHeaders ();
408+
409+ String statusStr = HttpHeadersUtils .getSingleHeader (headers , UPLOAD_STATUS_HEADER );
410+ if (statusStr != null ) {
411+ this .hasUploadStatusHeader = true ;
412+ if (STATUS_FINAL .equalsIgnoreCase (statusStr )) {
413+ this .isComplete = true ;
414+ }
415+ }
416+
417+ this .committedOffset = parseSizeReceived (responseHeaders );
418+ }
419+
420+ @ Override
421+ public void onMessage (@ Nullable String message ) {
422+ if (message != null ) {
423+ this .responseBody = message ;
424+ }
425+ }
426+
427+ @ Override
428+ public void onClose (int statusCode , HttpJsonMetadata trailers ) {
429+ try {
430+ if (statusCode >= 200 && statusCode < 300 ) {
431+ if (!hasUploadStatusHeader ) {
432+ future .setException (
433+ ApiExceptionFactory .createException (
434+ "Upload chunk response did not contain valid X-Goog-Upload-Status header" ,
435+ /* cause= */ null ,
436+ HttpJsonStatusCode .of (StatusCode .Code .INTERNAL ),
437+ /* retryable= */ false ));
438+ return ;
439+ }
440+ long confirmedOffset =
441+ committedOffset != null
442+ ? committedOffset
443+ : request .getOffset () + request .getPayload ().size ();
444+ ResponseT response = null ;
445+ if (isComplete ) {
446+ InputStream stream =
447+ new ByteArrayInputStream (responseBody .getBytes (StandardCharsets .UTF_8 ));
448+ response = responseParser .parse (stream );
449+ }
450+ future .set (ChunkUploadResponse .create (confirmedOffset , isComplete , response ));
451+ } else {
452+ Throwable cause = trailers .getException ();
453+ future .setException (
454+ cause != null
455+ ? cause
456+ : new HttpJsonStatusRuntimeException (
457+ statusCode , "Failed to upload chunk with status code: " + statusCode , null ));
458+ }
459+ } catch (Throwable t ) {
460+ future .setException (t );
461+ }
462+ }
463+ }
271464}
0 commit comments