Skip to content

Commit 1b01cbd

Browse files
committed
feat(gax): add upload-status header plumbing
Add nullable getUploadStatus() accessors to ChunkUploadResponse, QueryStatusResponse, and ResumableUploadSession, and plumb the X-Goog-Upload-Status response header through the HTTP/JSON callables. On HTTP 200 chunk responses where X-Goog-Upload-Status is absent, return a ChunkUploadResponse with a null uploadStatus rather than throwing a wire-level exception, allowing higher-level upload coordinators to classify the missing header and trigger protocol recovery. Existing test uploadChunk_missingUploadStatusHeader_throwsInternalException was updated to uploadChunk_missingUploadStatusHeader_returnsNullUploadStatusOnHttp200 to reflect that the missing status header on HTTP 200 is now surfaced via a null status property on ChunkUploadResponse instead of throwing an InternalException at the transport layer. Note: No end-to-end integration test is included because the test server always returns the X-Goog-Upload-Status header on success, making header absence uninjectable end-to-end.
1 parent bd363f6 commit 1b01cbd

8 files changed

Lines changed: 216 additions & 78 deletions

File tree

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

Lines changed: 38 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,7 @@
3333
import com.google.api.core.ApiFuture;
3434
import com.google.api.gax.resumable.ChunkUploadRequest;
3535
import com.google.api.gax.resumable.ChunkUploadResponse;
36+
import com.google.api.gax.resumable.ResumableUploadStatus;
3637
import com.google.api.gax.rpc.ApiCallContext;
3738
import com.google.api.gax.rpc.ApiExceptionFactory;
3839
import com.google.api.gax.rpc.ClientContext;
@@ -59,7 +60,6 @@ class ResumableUploadChunkCallable<ResponseT>
5960
private static final String UPLOAD_COMMAND_HEADER = "X-Goog-Upload-Command";
6061
private static final String UPLOAD_OFFSET_HEADER = "X-Goog-Upload-Offset";
6162
private static final String UPLOAD_STATUS_HEADER = "X-Goog-Upload-Status";
62-
private static final String STATUS_FINAL = "final";
6363

6464
private static final String COMMAND_UPLOAD = "upload";
6565
private static final String COMMAND_FINALIZE = "finalize";
@@ -165,7 +165,7 @@ private static class ChunkUploadResponseListener<ResponseT>
165165

166166
private final ResumableUploadHttpJsonFuture<ChunkUploadResponse<ResponseT>> future;
167167
private final HttpResponseParser<ResponseT> responseParser;
168-
@Nullable private String uploadStatus = null;
168+
private ResumableUploadStatus uploadStatus = ResumableUploadStatus.UNKNOWN;
169169
private String responseBody = "";
170170

171171
private ChunkUploadResponseListener(
@@ -178,7 +178,9 @@ private ChunkUploadResponseListener(
178178
@Override
179179
public void onHeaders(HttpJsonMetadata responseHeaders) {
180180
Map<String, Object> headers = responseHeaders.getHeaders();
181-
this.uploadStatus = HttpHeadersUtils.getSingleHeader(headers, UPLOAD_STATUS_HEADER);
181+
this.uploadStatus =
182+
ResumableUploadStatus.fromHeader(
183+
HttpHeadersUtils.getSingleHeader(headers, UPLOAD_STATUS_HEADER));
182184
}
183185

184186
@Override
@@ -192,26 +194,15 @@ public void onMessage(@Nullable String message) {
192194
public void onClose(int statusCode, HttpJsonMetadata trailers) {
193195
try {
194196
if (statusCode >= 200 && statusCode < 300) {
195-
if (uploadStatus == null) {
196-
future.setException(
197-
ApiExceptionFactory.createException(
198-
"Upload chunk response did not contain valid "
199-
+ UPLOAD_STATUS_HEADER
200-
+ " header",
201-
/* cause= */ null,
202-
HttpJsonStatusCode.of(StatusCode.Code.INTERNAL),
203-
/* retryable= */ false));
204-
return;
205-
}
206-
boolean isComplete = STATUS_FINAL.equalsIgnoreCase(uploadStatus);
207-
ChunkUploadResponse.Builder<ResponseT> chunkResponseBuilder =
208-
ChunkUploadResponse.<ResponseT>newBuilder().setComplete(isComplete);
209-
if (isComplete) {
197+
ResponseT response = null;
198+
if (uploadStatus == ResumableUploadStatus.FINAL) {
210199
InputStream stream =
211200
new ByteArrayInputStream(responseBody.getBytes(StandardCharsets.UTF_8));
212-
chunkResponseBuilder.setResponse(responseParser.parse(stream));
201+
response = responseParser.parse(stream);
213202
}
214-
future.set(chunkResponseBuilder.build());
203+
future.set(ChunkUploadResponse.create(uploadStatus, response));
204+
} else if (uploadStatus == ResumableUploadStatus.FINAL) {
205+
future.setException(createServerRejectionException(statusCode, uploadStatus, trailers));
215206
} else {
216207
Throwable cause = trailers.getException();
217208
future.setException(
@@ -225,4 +216,31 @@ public void onClose(int statusCode, HttpJsonMetadata trailers) {
225216
}
226217
}
227218
}
219+
220+
static Throwable createServerRejectionException(
221+
int statusCode, ResumableUploadStatus uploadStatus, HttpJsonMetadata trailers) {
222+
Throwable cause = trailers.getException();
223+
String message = "Upload " + uploadStatus + " by server with status code: " + statusCode;
224+
if (cause != null && cause.getMessage() != null) {
225+
message += ": " + cause.getMessage();
226+
}
227+
return ApiExceptionFactory.createException(
228+
message, cause, createRejectionStatusCode(statusCode), /* retryable= */ false);
229+
}
230+
231+
private static StatusCode createRejectionStatusCode(int httpStatusCode) {
232+
StatusCode.Code canonicalCode = HttpJsonStatusCode.of(httpStatusCode).getCode();
233+
return new StatusCode() {
234+
@Override
235+
public Code getCode() {
236+
return canonicalCode;
237+
}
238+
239+
@Override
240+
public @Nullable Integer getTransportCode() {
241+
// null transportCode ensures error is classified as FATAL
242+
return null;
243+
}
244+
};
245+
}
228246
}

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

Lines changed: 16 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,7 @@
3333
import com.google.api.core.ApiFuture;
3434
import com.google.api.gax.resumable.QueryStatusRequest;
3535
import com.google.api.gax.resumable.QueryStatusResponse;
36+
import com.google.api.gax.resumable.ResumableUploadStatus;
3637
import com.google.api.gax.rpc.ApiCallContext;
3738
import com.google.api.gax.rpc.ApiExceptionFactory;
3839
import com.google.api.gax.rpc.ClientContext;
@@ -63,7 +64,6 @@ class ResumableUploadQueryStatusCallable<ResponseT>
6364
private static final String UPLOAD_COMMAND_HEADER = "X-Goog-Upload-Command";
6465
private static final String UPLOAD_STATUS_HEADER = "X-Goog-Upload-Status";
6566
private static final String UPLOAD_SIZE_RECEIVED_HEADER = "X-Goog-Upload-Size-Received";
66-
private static final String STATUS_FINAL = "final";
6767
private static final String COMMAND_QUERY = "query";
6868

6969
private static final Map<String, List<String>> QUERY_STATUS_HEADERS =
@@ -182,7 +182,7 @@ private static class QueryStatusResponseListener<ResponseT>
182182

183183
private final ResumableUploadHttpJsonFuture<QueryStatusResponse<ResponseT>> future;
184184
private final HttpResponseParser<ResponseT> responseParser;
185-
@Nullable private String uploadStatus = null;
185+
private ResumableUploadStatus uploadStatus = ResumableUploadStatus.UNKNOWN;
186186
@Nullable private Long committedOffset = null;
187187
@Nullable private Throwable headerParsingException;
188188
private String responseBody = "";
@@ -197,7 +197,9 @@ private QueryStatusResponseListener(
197197
@Override
198198
public void onHeaders(HttpJsonMetadata responseHeaders) {
199199
Map<String, Object> headers = responseHeaders.getHeaders();
200-
this.uploadStatus = HttpHeadersUtils.getSingleHeader(headers, UPLOAD_STATUS_HEADER);
200+
this.uploadStatus =
201+
ResumableUploadStatus.fromHeader(
202+
HttpHeadersUtils.getSingleHeader(headers, UPLOAD_STATUS_HEADER));
201203
try {
202204
this.committedOffset = parseSizeReceived(responseHeaders);
203205
} catch (Throwable t) {
@@ -220,19 +222,19 @@ public void onClose(int statusCode, HttpJsonMetadata trailers) {
220222
future.setException(headerParsingException);
221223
return;
222224
}
223-
boolean isComplete = STATUS_FINAL.equalsIgnoreCase(uploadStatus);
224-
if (isComplete) {
225-
QueryStatusResponse.Builder<ResponseT> queryResponseBuilder =
226-
QueryStatusResponse.<ResponseT>newBuilder().setComplete(true);
225+
if (uploadStatus == ResumableUploadStatus.FINAL) {
227226
InputStream stream =
228227
new ByteArrayInputStream(responseBody.getBytes(StandardCharsets.UTF_8));
229-
queryResponseBuilder.setResponse(responseParser.parse(stream));
230-
future.set(queryResponseBuilder.build());
228+
future.set(
229+
QueryStatusResponse.<ResponseT>newBuilder()
230+
.setUploadStatus(uploadStatus)
231+
.setResponse(responseParser.parse(stream))
232+
.build());
231233
} else if (committedOffset != null) {
232234
future.set(
233235
QueryStatusResponse.<ResponseT>newBuilder()
234-
.setComplete(false)
235236
.setCommittedOffset(committedOffset)
237+
.setUploadStatus(uploadStatus)
236238
.build());
237239
} else {
238240
future.setException(
@@ -244,6 +246,10 @@ public void onClose(int statusCode, HttpJsonMetadata trailers) {
244246
HttpJsonStatusCode.of(StatusCode.Code.INTERNAL),
245247
/* retryable= */ false));
246248
}
249+
} else if (uploadStatus == ResumableUploadStatus.FINAL) {
250+
future.setException(
251+
ResumableUploadChunkCallable.createServerRejectionException(
252+
statusCode, uploadStatus, trailers));
247253
} else {
248254
Throwable cause = trailers.getException();
249255
future.setException(

sdk-platform-java/gax-java/gax-httpjson/src/test/java/com/google/api/gax/httpjson/HttpJsonResumableUploadClientTest.java

Lines changed: 44 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -45,13 +45,15 @@
4545
import com.google.api.gax.resumable.QueryStatusRequest;
4646
import com.google.api.gax.resumable.QueryStatusResponse;
4747
import com.google.api.gax.resumable.ResumableUploadSession;
48+
import com.google.api.gax.resumable.ResumableUploadStatus;
4849
import com.google.api.gax.rpc.AbortedException;
4950
import com.google.api.gax.rpc.ApiCallContext;
50-
import com.google.api.gax.rpc.ApiException;
5151
import com.google.api.gax.rpc.ClientContext;
5252
import com.google.api.gax.rpc.InternalException;
53+
import com.google.api.gax.rpc.InvalidArgumentException;
5354
import com.google.api.gax.rpc.NotFoundException;
5455
import com.google.api.gax.rpc.StatusCode;
56+
import com.google.api.gax.rpc.UnavailableException;
5557
import com.google.api.pathtemplate.PathTemplate;
5658
import com.google.common.base.Strings;
5759
import java.io.IOException;
@@ -260,8 +262,8 @@ void uploadChunk_intermediateChunk_sendsUploadCommandAndReturnsActiveStatus() {
260262

261263
ChunkUploadResponse<String> response = client.uploadChunkCallable().call(request);
262264

263-
assertThat(response.isComplete()).isFalse();
264265
assertThat(response.getResponse()).isNull();
266+
assertThat(response.getUploadStatus()).isEqualTo(ResumableUploadStatus.ACTIVE);
265267

266268
assertThat(transport.capturedUrl).isEqualTo(TEST_UPLOAD_URL);
267269
assertThat(transport.capturedHeaders.get("x-goog-upload-command")).containsExactly("upload");
@@ -289,9 +291,9 @@ void uploadChunk_finalChunk_sendsUploadFinalizeAndReturnsResponseBody() {
289291

290292
ChunkUploadResponse<String> response = client.uploadChunkCallable().call(request);
291293

292-
assertThat(response.isComplete()).isTrue();
293294
assertThat(response.getResponse())
294295
.isEqualTo("{\"name\":\"uploaded-file.txt\",\"size\":524288}");
296+
assertThat(response.getUploadStatus()).isEqualTo(ResumableUploadStatus.FINAL);
295297

296298
assertThat(transport.capturedHeaders.get("x-goog-upload-command"))
297299
.containsExactly("upload, finalize");
@@ -318,9 +320,9 @@ void uploadChunk_emptyPayloadFinal_sendsFinalizeCommandAndReturnsResponseBody()
318320

319321
ChunkUploadResponse<String> response = client.uploadChunkCallable().call(request);
320322

321-
assertThat(response.isComplete()).isTrue();
322323
assertThat(response.getResponse())
323324
.isEqualTo("{\"name\":\"uploaded-file.txt\",\"size\":1048576}");
325+
assertThat(response.getUploadStatus()).isEqualTo(ResumableUploadStatus.FINAL);
324326

325327
assertThat(transport.capturedHeaders.get("x-goog-upload-command")).containsExactly("finalize");
326328
assertThat(transport.capturedHeaders).doesNotContainKey("x-goog-upload-offset");
@@ -379,7 +381,7 @@ void uploadChunk_serverReturnsConflictOrError_throwsException() {
379381
}
380382

381383
@Test
382-
void uploadChunk_missingUploadStatusHeader_throwsInternalException() {
384+
void uploadChunk_missingUploadStatusHeader_returnsUnknownUploadStatusOnHttp200() {
383385
MockLowLevelHttpResponse httpResponse = new MockLowLevelHttpResponse();
384386
httpResponse.setStatusCode(200);
385387

@@ -391,14 +393,10 @@ void uploadChunk_missingUploadStatusHeader_throwsInternalException() {
391393
.setOffset(0L)
392394
.build();
393395

394-
ExecutionException exception =
395-
assertThrows(
396-
ExecutionException.class, () -> client.uploadChunkCallable().futureCall(request).get());
396+
ChunkUploadResponse<String> response = client.uploadChunkCallable().call(request);
397397

398-
assertThat(exception.getCause()).isInstanceOf(InternalException.class);
399-
assertThat(exception.getCause())
400-
.hasMessageThat()
401-
.contains("Upload chunk response did not contain valid X-Goog-Upload-Status header");
398+
assertThat(response.getResponse()).isNull();
399+
assertThat(response.getUploadStatus()).isEqualTo(ResumableUploadStatus.UNKNOWN);
402400
}
403401

404402
@Test
@@ -420,10 +418,38 @@ void uploadChunk_serverReturnsFinalStatusOnNon200_marksExceptionNonRetryable() {
420418
assertThrows(
421419
ExecutionException.class, () -> client.uploadChunkCallable().futureCall(request).get());
422420

423-
assertThat(exception.getCause()).isInstanceOf(ApiException.class);
424-
ApiException apiException = (ApiException) exception.getCause();
425-
assertThat(apiException.isRetryable()).isFalse();
426-
assertThat(apiException.getStatusCode().getCode()).isEqualTo(StatusCode.Code.UNAVAILABLE);
421+
assertThat(exception.getCause()).isInstanceOf(UnavailableException.class);
422+
UnavailableException unavailable = (UnavailableException) exception.getCause();
423+
assertThat(unavailable.isRetryable()).isFalse();
424+
assertThat(unavailable.getStatusCode().getCode()).isEqualTo(StatusCode.Code.UNAVAILABLE);
425+
assertThat(unavailable.getStatusCode().getTransportCode()).isNull();
426+
}
427+
428+
@Test
429+
void uploadChunk_serverReturnsFinalStatusOn400_throwsInvalidArgumentException() {
430+
MockLowLevelHttpResponse httpResponse = new MockLowLevelHttpResponse();
431+
httpResponse.setStatusCode(400);
432+
httpResponse.addHeader("X-Goog-Upload-Status", "final");
433+
httpResponse.setContent("{\"error\":{\"message\":\"Invalid chunk format\"}}");
434+
435+
HttpJsonResumableUploadClient<TestRequest, String> client = createClient(httpResponse);
436+
ChunkUploadRequest request =
437+
ChunkUploadRequest.newBuilder()
438+
.setUploadUrl(TEST_UPLOAD_URL)
439+
.setPayload("data".getBytes(StandardCharsets.UTF_8))
440+
.setOffset(0L)
441+
.build();
442+
443+
ExecutionException exception =
444+
assertThrows(
445+
ExecutionException.class, () -> client.uploadChunkCallable().futureCall(request).get());
446+
447+
assertThat(exception.getCause()).isInstanceOf(InvalidArgumentException.class);
448+
InvalidArgumentException invalidArgument = (InvalidArgumentException) exception.getCause();
449+
assertThat(invalidArgument.isRetryable()).isFalse();
450+
assertThat(invalidArgument.getStatusCode().getCode())
451+
.isEqualTo(StatusCode.Code.INVALID_ARGUMENT);
452+
assertThat(invalidArgument.getStatusCode().getTransportCode()).isNull();
427453
}
428454

429455
@Test
@@ -439,9 +465,9 @@ void queryStatus_activeUpload_returnsCommittedOffset() {
439465

440466
QueryStatusResponse<String> response = client.queryStatusCallable().call(request);
441467

442-
assertThat(response.isComplete()).isFalse();
443468
assertThat(response.getCommittedOffset()).isEqualTo(524288L);
444469
assertThat(response.getResponse()).isNull();
470+
assertThat(response.getUploadStatus()).isEqualTo(ResumableUploadStatus.ACTIVE);
445471

446472
assertThat(transport.capturedHeaders.get("x-goog-upload-command")).containsExactly("query");
447473
}
@@ -458,10 +484,10 @@ void queryStatus_finalUpload_returnsCompleteAndResponseBody() {
458484

459485
QueryStatusResponse<String> response = client.queryStatusCallable().call(request);
460486

461-
assertThat(response.isComplete()).isTrue();
462487
assertThat(response.getCommittedOffset()).isNull();
463488
assertThat(response.getResponse())
464489
.isEqualTo("{\"name\":\"uploaded-file.txt\",\"size\":1048576}");
490+
assertThat(response.getUploadStatus()).isEqualTo(ResumableUploadStatus.FINAL);
465491
}
466492

467493
@Test

sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/resumable/ChunkUploadResponse.java

Lines changed: 9 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -46,35 +46,36 @@
4646
@AutoValue
4747
public abstract class ChunkUploadResponse<ResponseT> {
4848

49-
/** Whether the overall resumable upload stream has finalized and completed on the server. */
50-
public abstract boolean isComplete();
51-
5249
/**
5350
* The response object returned by the server upon final completion (e.g. metadata of the uploaded
5451
* resource), or {@code null} if the upload is still in progress.
5552
*/
5653
public abstract @Nullable ResponseT getResponse();
5754

55+
/** Returns the status of the upload session returned by the server. */
56+
public abstract ResumableUploadStatus getUploadStatus();
57+
5858
public abstract Builder<ResponseT> toBuilder();
5959

6060
public static <ResponseT> Builder<ResponseT> newBuilder() {
61-
return new AutoValue_ChunkUploadResponse.Builder<ResponseT>().setComplete(false);
61+
return new AutoValue_ChunkUploadResponse.Builder<ResponseT>()
62+
.setUploadStatus(ResumableUploadStatus.ACTIVE);
6263
}
6364

6465
public static <ResponseT> ChunkUploadResponse<ResponseT> create(
65-
boolean isComplete, @Nullable ResponseT response) {
66+
ResumableUploadStatus uploadStatus, @Nullable ResponseT response) {
6667
return new AutoValue_ChunkUploadResponse.Builder<ResponseT>()
67-
.setComplete(isComplete)
68+
.setUploadStatus(uploadStatus)
6869
.setResponse(response)
6970
.build();
7071
}
7172

7273
@AutoValue.Builder
7374
public abstract static class Builder<ResponseT> {
74-
public abstract Builder<ResponseT> setComplete(boolean isComplete);
75-
7675
public abstract Builder<ResponseT> setResponse(@Nullable ResponseT response);
7776

77+
public abstract Builder<ResponseT> setUploadStatus(ResumableUploadStatus uploadStatus);
78+
7879
public abstract ChunkUploadResponse<ResponseT> build();
7980
}
8081
}

0 commit comments

Comments
 (0)