3333import com .google .api .core .ApiFuture ;
3434import com .google .api .core .InternalApi ;
3535import com .google .api .core .SettableApiFuture ;
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 ()
@@ -120,13 +166,61 @@ public ApiFuture<ResumableUploadSession> futureCall(
120166 };
121167 this .startUploadCallable =
122168 new HttpJsonExceptionCallable <>(rawStartUploadCallable , Collections .emptySet ());
169+
170+ UnaryCallable <ChunkUploadRequest , ChunkUploadResponse <ResponseT >> rawUploadChunkCallable =
171+ new UnaryCallable <ChunkUploadRequest , ChunkUploadResponse <ResponseT >>() {
172+ @ Override
173+ public ApiFuture <ChunkUploadResponse <ResponseT >> futureCall (
174+ ChunkUploadRequest request , @ Nullable ApiCallContext inputContext ) {
175+ Preconditions .checkNotNull (request );
176+ boolean isPayloadEmpty = request .getPayload ().isEmpty ();
177+ String command ;
178+ if (request .isFinal ()) {
179+ command = !isPayloadEmpty ? "upload, finalize" : "finalize" ;
180+ } else {
181+ command = "upload" ;
182+ }
183+ Map <String , List <String >> chunkHeaders =
184+ ImmutableMap .of (
185+ UPLOAD_COMMAND_HEADER ,
186+ ImmutableList .of (command ),
187+ UPLOAD_OFFSET_HEADER ,
188+ ImmutableList .of (String .valueOf (request .getOffset ())));
189+
190+ HttpJsonCallContext context =
191+ (HttpJsonCallContext )
192+ HttpJsonCallContext .createDefault ()
193+ .nullToSelf (clientContext .getDefaultCallContext ())
194+ .merge (inputContext )
195+ .withExtraHeaders (chunkHeaders );
196+
197+ HttpJsonClientCall <ChunkUploadRequest , String > clientCall =
198+ HttpJsonClientCalls .newCall (UPLOAD_CHUNK_DESCRIPTOR , context );
199+
200+ SettableApiFuture <ChunkUploadResponse <ResponseT >> future = SettableApiFuture .create ();
201+ HttpJsonClientCalls .startUnaryCall (
202+ clientCall ,
203+ request ,
204+ context ,
205+ new ChunkUploadResponseListener <>(request , future , responseParser ));
206+
207+ return future ;
208+ }
209+ };
210+ this .uploadChunkCallable =
211+ new HttpJsonExceptionCallable <>(rawUploadChunkCallable , Collections .emptySet ());
123212 }
124213
125214 @ Override
126215 public UnaryCallable <RequestT , ResumableUploadSession > startUploadCallable () {
127216 return startUploadCallable ;
128217 }
129218
219+ @ Override
220+ public UnaryCallable <ChunkUploadRequest , ChunkUploadResponse <ResponseT >> uploadChunkCallable () {
221+ return uploadChunkCallable ;
222+ }
223+
130224 private static class StartUploadResponseListener extends HttpJsonClientCall .Listener <String > {
131225
132226 private final SettableApiFuture <ResumableUploadSession > future ;
@@ -214,4 +308,93 @@ public void onClose(int statusCode, HttpJsonMetadata trailers) {
214308 }
215309 }
216310 }
311+
312+ private static class ChunkUploadResponseListener <ResponseT >
313+ extends HttpJsonClientCall .Listener <String > {
314+
315+ private final ChunkUploadRequest request ;
316+ private final SettableApiFuture <ChunkUploadResponse <ResponseT >> future ;
317+ private final HttpResponseParser <ResponseT > responseParser ;
318+ private boolean hasUploadStatusHeader = false ;
319+ private boolean isComplete = false ;
320+ private long committedOffset = -1L ;
321+ private String responseBody = "" ;
322+
323+ ChunkUploadResponseListener (
324+ ChunkUploadRequest request ,
325+ SettableApiFuture <ChunkUploadResponse <ResponseT >> future ,
326+ HttpResponseParser <ResponseT > responseParser ) {
327+ this .request = request ;
328+ this .future = future ;
329+ this .responseParser = responseParser ;
330+ }
331+
332+ @ Override
333+ public void onHeaders (HttpJsonMetadata responseHeaders ) {
334+ Map <String , Object > headers = responseHeaders .getHeaders ();
335+
336+ String statusStr = HttpHeadersUtils .getSingleHeader (headers , UPLOAD_STATUS_HEADER );
337+ if (statusStr != null ) {
338+ this .hasUploadStatusHeader = true ;
339+ if (STATUS_FINAL .equalsIgnoreCase (statusStr )) {
340+ this .isComplete = true ;
341+ }
342+ }
343+
344+ String sizeReceivedStr =
345+ HttpHeadersUtils .getSingleHeader (headers , UPLOAD_SIZE_RECEIVED_HEADER );
346+ if (!Strings .isNullOrEmpty (sizeReceivedStr )) {
347+ try {
348+ this .committedOffset = Long .parseLong (sizeReceivedStr );
349+ } catch (NumberFormatException ignored ) {
350+ // Ignore invalid/malformed size received header and fall back to local offset
351+ // calculation.
352+ }
353+ }
354+ }
355+
356+ @ Override
357+ public void onMessage (@ Nullable String message ) {
358+ if (message != null ) {
359+ this .responseBody = message ;
360+ }
361+ }
362+
363+ @ Override
364+ public void onClose (int statusCode , HttpJsonMetadata trailers ) {
365+ try {
366+ if (statusCode >= 200 && statusCode < 300 ) {
367+ if (!hasUploadStatusHeader ) {
368+ future .setException (
369+ ApiExceptionFactory .createException (
370+ "Upload chunk response did not contain valid X-Goog-Upload-Status header" ,
371+ /* cause= */ null ,
372+ HttpJsonStatusCode .of (StatusCode .Code .INTERNAL ),
373+ /* retryable= */ false ));
374+ return ;
375+ }
376+ long confirmedOffset =
377+ committedOffset >= 0
378+ ? committedOffset
379+ : request .getOffset () + request .getPayload ().size ();
380+ ResponseT response = null ;
381+ if (isComplete ) {
382+ InputStream stream =
383+ new ByteArrayInputStream (responseBody .getBytes (StandardCharsets .UTF_8 ));
384+ response = responseParser .parse (stream );
385+ }
386+ future .set (ChunkUploadResponse .create (confirmedOffset , isComplete , response ));
387+ } else {
388+ Throwable cause = trailers .getException ();
389+ future .setException (
390+ cause != null
391+ ? cause
392+ : new HttpJsonStatusRuntimeException (
393+ statusCode , "Failed to upload chunk with status code: " + statusCode , null ));
394+ }
395+ } catch (Throwable t ) {
396+ future .setException (t );
397+ }
398+ }
399+ }
217400}
0 commit comments