3333import com .google .api .core .AbstractApiFuture ;
3434import com .google .api .core .ApiFuture ;
3535import com .google .api .core .InternalApi ;
36+ import com .google .api .gax .resumable .ChunkUploadRequest ;
37+ import com .google .api .gax .resumable .ChunkUploadResponse ;
3638import com .google .api .gax .resumable .ResumableUploadClient ;
3739import com .google .api .gax .resumable .ResumableUploadSession ;
3840import com .google .api .gax .rpc .ApiCallContext ;
3941import com .google .api .gax .rpc .ApiExceptionFactory ;
4042import com .google .api .gax .rpc .ClientContext ;
4143import com .google .api .gax .rpc .StatusCode ;
4244import com .google .api .gax .rpc .UnaryCallable ;
45+ import com .google .api .pathtemplate .PathTemplate ;
4346import com .google .common .base .Preconditions ;
4447import com .google .common .base .Strings ;
4548import com .google .common .collect .ImmutableList ;
4649import com .google .common .collect .ImmutableMap ;
50+ import java .io .ByteArrayInputStream ;
51+ import java .io .InputStream ;
52+ import java .nio .charset .StandardCharsets ;
4753import java .util .Collections ;
4854import java .util .List ;
4955import java .util .Map ;
@@ -65,16 +71,54 @@ public final class HttpJsonResumableUploadClient<RequestT, ResponseT>
6571
6672 private static final String UPLOAD_PROTOCOL_HEADER = "X-Goog-Upload-Protocol" ;
6773 private static final String UPLOAD_COMMAND_HEADER = "X-Goog-Upload-Command" ;
74+ private static final String UPLOAD_OFFSET_HEADER = "X-Goog-Upload-Offset" ;
6875 private static final String UPLOAD_URL_HEADER = "X-Goog-Upload-URL" ;
6976 private static final String UPLOAD_GRANULARITY_HEADER = "X-Goog-Upload-Chunk-Granularity" ;
77+ private static final String UPLOAD_STATUS_HEADER = "X-Goog-Upload-Status" ;
78+ private static final String UPLOAD_SIZE_RECEIVED_HEADER = "X-Goog-Upload-Size-Received" ;
79+ private static final String STATUS_FINAL = "final" ;
7080
7181 private static final Map <String , List <String >> START_UPLOAD_HEADERS =
7282 ImmutableMap .of (
7383 UPLOAD_PROTOCOL_HEADER , ImmutableList .of ("resumable" ),
7484 UPLOAD_COMMAND_HEADER , ImmutableList .of ("start" ));
7585
86+ private static final PathTemplate PATH_TEMPLATE = PathTemplate .create ("{+path}" );
87+
88+ private static final ApiMethodDescriptor <ChunkUploadRequest , String > UPLOAD_CHUNK_DESCRIPTOR =
89+ ApiMethodDescriptor .<ChunkUploadRequest , String >newBuilder ()
90+ .setFullMethodName ("ResumableUpload/UploadChunk" )
91+ .setHttpMethod (HttpMethods .POST )
92+ .setType (ApiMethodDescriptor .MethodType .UNARY )
93+ .setRequestFormatter (
94+ new ResumableUploadChunkRequestFormatter <ChunkUploadRequest >() {
95+ @ Override
96+ public Map <String , List <String >> getQueryParamNames (ChunkUploadRequest request ) {
97+ return Collections .emptyMap ();
98+ }
99+
100+ @ Override
101+ public byte [] getBinaryRequestBody (ChunkUploadRequest request ) {
102+ return request .getPayload ().toByteArray ();
103+ }
104+
105+ @ Override
106+ public String getPath (ChunkUploadRequest request ) {
107+ return request .getUploadUrl ();
108+ }
109+
110+ @ Override
111+ public PathTemplate getPathTemplate () {
112+ return PATH_TEMPLATE ;
113+ }
114+ })
115+ .setResponseParser (ResumableUploadResponseParser .create ())
116+ .build ();
117+
76118 private final ApiMethodDescriptor <RequestT , String > startUploadDescriptor ;
77119 private final UnaryCallable <RequestT , ResumableUploadSession > startUploadCallable ;
120+ private final UnaryCallable <ChunkUploadRequest , ChunkUploadResponse <ResponseT >>
121+ uploadChunkCallable ;
78122
79123 public static <RequestT , ResponseT > HttpJsonResumableUploadClient <RequestT , ResponseT > create (
80124 ClientContext clientContext , ApiMethodDescriptor <RequestT , ResponseT > methodDescriptor ) {
@@ -85,6 +129,8 @@ private HttpJsonResumableUploadClient(
85129 ClientContext clientContext , ApiMethodDescriptor <RequestT , ResponseT > methodDescriptor ) {
86130 Preconditions .checkNotNull (clientContext );
87131 Preconditions .checkNotNull (methodDescriptor );
132+ HttpResponseParser <ResponseT > responseParser =
133+ Preconditions .checkNotNull (methodDescriptor .getResponseParser ());
88134
89135 this .startUploadDescriptor =
90136 ApiMethodDescriptor .<RequestT , String >newBuilder ()
@@ -95,13 +141,19 @@ private HttpJsonResumableUploadClient(
95141 .setResponseParser (ResumableUploadResponseParser .create ())
96142 .build ();
97143 this .startUploadCallable = createStartUploadCallable (clientContext );
144+ this .uploadChunkCallable = createUploadChunkCallable (clientContext , responseParser );
98145 }
99146
100147 @ Override
101148 public UnaryCallable <RequestT , ResumableUploadSession > startUploadCallable () {
102149 return startUploadCallable ;
103150 }
104151
152+ @ Override
153+ public UnaryCallable <ChunkUploadRequest , ChunkUploadResponse <ResponseT >> uploadChunkCallable () {
154+ return uploadChunkCallable ;
155+ }
156+
105157 private UnaryCallable <RequestT , ResumableUploadSession > createStartUploadCallable (
106158 ClientContext clientContext ) {
107159 UnaryCallable <RequestT , ResumableUploadSession > rawCallable =
@@ -127,6 +179,49 @@ public ApiFuture<ResumableUploadSession> futureCall(
127179 return createClientCallable (rawCallable , clientContext );
128180 }
129181
182+ private UnaryCallable <ChunkUploadRequest , ChunkUploadResponse <ResponseT >>
183+ createUploadChunkCallable (
184+ ClientContext clientContext , HttpResponseParser <ResponseT > responseParser ) {
185+ UnaryCallable <ChunkUploadRequest , ChunkUploadResponse <ResponseT >> rawCallable =
186+ new UnaryCallable <ChunkUploadRequest , ChunkUploadResponse <ResponseT >>() {
187+ @ Override
188+ public ApiFuture <ChunkUploadResponse <ResponseT >> futureCall (
189+ ChunkUploadRequest request , @ Nullable ApiCallContext inputContext ) {
190+ Preconditions .checkNotNull (request );
191+ boolean isPayloadEmpty = request .getPayload ().isEmpty ();
192+ String command ;
193+ if (request .isFinal ()) {
194+ command = !isPayloadEmpty ? "upload, finalize" : "finalize" ;
195+ } else {
196+ command = "upload" ;
197+ }
198+ Map <String , List <String >> chunkHeaders =
199+ ImmutableMap .of (
200+ UPLOAD_COMMAND_HEADER ,
201+ ImmutableList .of (command ),
202+ UPLOAD_OFFSET_HEADER ,
203+ ImmutableList .of (String .valueOf (request .getOffset ())));
204+
205+ HttpJsonCallContext context =
206+ createCallContext (clientContext , inputContext , chunkHeaders );
207+
208+ HttpJsonClientCall <ChunkUploadRequest , String > clientCall =
209+ HttpJsonClientCalls .newCall (UPLOAD_CHUNK_DESCRIPTOR , context );
210+
211+ HttpJsonCallFuture <ChunkUploadResponse <ResponseT >> future =
212+ new HttpJsonCallFuture <>(clientCall );
213+ HttpJsonClientCalls .startUnaryCall (
214+ clientCall ,
215+ request ,
216+ context ,
217+ new ChunkUploadResponseListener <>(request , future , responseParser ));
218+
219+ return future ;
220+ }
221+ };
222+ return createClientCallable (rawCallable , clientContext );
223+ }
224+
130225 private static HttpJsonCallContext createCallContext (
131226 ClientContext clientContext ,
132227 @ Nullable ApiCallContext inputContext ,
@@ -148,6 +243,20 @@ private static <CallReqT, CallRespT> UnaryCallable<CallReqT, CallRespT> createCl
148243 return callable .withDefaultCallContext (clientContext .getDefaultCallContext ());
149244 }
150245
246+ @ Nullable
247+ private static Long parseSizeReceived (HttpJsonMetadata responseHeaders ) {
248+ String sizeReceivedStr =
249+ HttpHeadersUtils .getSingleHeader (responseHeaders .getHeaders (), UPLOAD_SIZE_RECEIVED_HEADER );
250+ if (!Strings .isNullOrEmpty (sizeReceivedStr )) {
251+ try {
252+ return Long .parseLong (sizeReceivedStr );
253+ } catch (NumberFormatException ignored ) {
254+ // Unparseable header; return null and let the listener decide how to handle it.
255+ }
256+ }
257+ return null ;
258+ }
259+
151260 /**
152261 * An {@link ApiFuture} that cancels the underlying {@link HttpJsonClientCall} to prevent
153262 * connection leaks.
@@ -165,12 +274,10 @@ protected void interruptTask() {
165274 call .cancel ("Call was cancelled" , null );
166275 }
167276
168- @ Override
169277 public boolean set (T value ) {
170278 return super .set (value );
171279 }
172280
173- @ Override
174281 public boolean setException (Throwable throwable ) {
175282 return super .setException (throwable );
176283 }
@@ -264,4 +371,88 @@ public void onClose(int statusCode, HttpJsonMetadata trailers) {
264371 }
265372 }
266373 }
374+
375+ /**
376+ * A listener that processes chunk upload response headers and bodies to produce the {@link
377+ * ChunkUploadResponse}.
378+ */
379+ private static class ChunkUploadResponseListener <ResponseT >
380+ extends HttpJsonClientCall .Listener <String > {
381+
382+ private final ChunkUploadRequest request ;
383+ private final HttpJsonCallFuture <ChunkUploadResponse <ResponseT >> future ;
384+ private final HttpResponseParser <ResponseT > responseParser ;
385+ private boolean hasUploadStatusHeader = false ;
386+ private boolean isComplete = false ;
387+ @ Nullable private Long committedOffset = null ;
388+ private String responseBody = "" ;
389+
390+ ChunkUploadResponseListener (
391+ ChunkUploadRequest request ,
392+ HttpJsonCallFuture <ChunkUploadResponse <ResponseT >> future ,
393+ HttpResponseParser <ResponseT > responseParser ) {
394+ this .request = request ;
395+ this .future = future ;
396+ this .responseParser = responseParser ;
397+ }
398+
399+ @ Override
400+ public void onHeaders (HttpJsonMetadata responseHeaders ) {
401+ Map <String , Object > headers = responseHeaders .getHeaders ();
402+
403+ String statusStr = HttpHeadersUtils .getSingleHeader (headers , UPLOAD_STATUS_HEADER );
404+ if (statusStr != null ) {
405+ this .hasUploadStatusHeader = true ;
406+ if (STATUS_FINAL .equalsIgnoreCase (statusStr )) {
407+ this .isComplete = true ;
408+ }
409+ }
410+
411+ this .committedOffset = parseSizeReceived (responseHeaders );
412+ }
413+
414+ @ Override
415+ public void onMessage (@ Nullable String message ) {
416+ if (message != null ) {
417+ this .responseBody = message ;
418+ }
419+ }
420+
421+ @ Override
422+ public void onClose (int statusCode , HttpJsonMetadata trailers ) {
423+ try {
424+ if (statusCode >= 200 && statusCode < 300 ) {
425+ if (!hasUploadStatusHeader ) {
426+ future .setException (
427+ ApiExceptionFactory .createException (
428+ "Upload chunk response did not contain valid X-Goog-Upload-Status header" ,
429+ /* cause= */ null ,
430+ HttpJsonStatusCode .of (StatusCode .Code .INTERNAL ),
431+ /* retryable= */ false ));
432+ return ;
433+ }
434+ long confirmedOffset =
435+ committedOffset != null
436+ ? committedOffset
437+ : request .getOffset () + request .getPayload ().size ();
438+ ResponseT response = null ;
439+ if (isComplete ) {
440+ InputStream stream =
441+ new ByteArrayInputStream (responseBody .getBytes (StandardCharsets .UTF_8 ));
442+ response = responseParser .parse (stream );
443+ }
444+ future .set (ChunkUploadResponse .create (confirmedOffset , isComplete , response ));
445+ } else {
446+ Throwable cause = trailers .getException ();
447+ future .setException (
448+ cause != null
449+ ? cause
450+ : new HttpJsonStatusRuntimeException (
451+ statusCode , "Failed to upload chunk with status code: " + statusCode , null ));
452+ }
453+ } catch (Throwable t ) {
454+ future .setException (t );
455+ }
456+ }
457+ }
267458}
0 commit comments