Skip to content

Commit 1e2a630

Browse files
committed
feat(gax): add recoverable error query and buffer realignment loop
Introduces a query-and-realign loop when encountering recoverable protocol errors during chunk uploads. Queries the server for the committed offset and adjusts the buffer window before resuming chunk transmission.
1 parent 074e557 commit 1e2a630

5 files changed

Lines changed: 435 additions & 34 deletions

File tree

sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ResumableUploadCallableImpl.java

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,8 @@
3737
import com.google.api.core.InternalApi;
3838
import com.google.api.gax.resumable.ChunkUploadRequest;
3939
import com.google.api.gax.resumable.ChunkUploadResponse;
40+
import com.google.api.gax.resumable.QueryStatusRequest;
41+
import com.google.api.gax.resumable.QueryStatusResponse;
4042
import com.google.api.gax.resumable.ResumableUploadClient;
4143
import com.google.api.gax.resumable.ResumableUploadSession;
4244
import com.google.api.gax.retrying.ExponentialRetryAlgorithm;
@@ -74,6 +76,8 @@ public class ResumableUploadCallableImpl<RequestT, ResponseT>
7476
private final ClientContext clientContext;
7577
private final UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>>
7678
retryingUploadChunkCallable;
79+
private final UnaryCallable<QueryStatusRequest, QueryStatusResponse<ResponseT>>
80+
retryingQueryCallable;
7781

7882
public ResumableUploadCallableImpl(
7983
ResumableUploadClient<RequestT, ResponseT> client,
@@ -86,6 +90,9 @@ public ResumableUploadCallableImpl(
8690
this.retryingUploadChunkCallable =
8791
createRetryingCallable(
8892
client.uploadChunkCallable(), ResumableUploadCommand.UPLOAD, clientContext);
93+
this.retryingQueryCallable =
94+
createRetryingCallable(
95+
client.queryStatusCallable(), ResumableUploadCommand.QUERY, clientContext);
8996
}
9097

9198
@Override
@@ -109,6 +116,7 @@ public ResumableUploadFuture<ResponseT> futureCall(
109116
return ResumableUploadFutureImpl.create(
110117
startFuture,
111118
retryingUploadChunkCallable,
119+
retryingQueryCallable,
112120
payload,
113121
effectiveSettings,
114122
clientContext.getDefaultCallContext());

sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ResumableUploadChunkCoordinator.java

Lines changed: 140 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -38,7 +38,10 @@
3838
import com.google.api.core.SettableApiFuture;
3939
import com.google.api.gax.resumable.ChunkUploadRequest;
4040
import com.google.api.gax.resumable.ChunkUploadResponse;
41+
import com.google.api.gax.resumable.QueryStatusRequest;
42+
import com.google.api.gax.resumable.QueryStatusResponse;
4143
import com.google.api.gax.resumable.ResumableUploadStatus;
44+
import com.google.api.gax.rpc.ResumableUploadErrorClassifier.Category;
4245
import com.google.common.util.concurrent.MoreExecutors;
4346
import java.io.IOException;
4447
import java.io.InputStream;
@@ -62,6 +65,8 @@ final class ResumableUploadChunkCoordinator<ResponseT> {
6265

6366
private final UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>>
6467
uploadChunkCallable;
68+
private final UnaryCallable<QueryStatusRequest, QueryStatusResponse<ResponseT>>
69+
queryStatusCallable;
6570
private final String uploadUrl;
6671
private final RewindableStreamBuffer buffer;
6772
private final ApiCallContext callContext;
@@ -70,12 +75,15 @@ final class ResumableUploadChunkCoordinator<ResponseT> {
7075

7176
ResumableUploadChunkCoordinator(
7277
UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>> uploadChunkCallable,
78+
UnaryCallable<QueryStatusRequest, QueryStatusResponse<ResponseT>> queryStatusCallable,
7379
String uploadUrl,
7480
InputStream payload,
7581
int chunkSize,
7682
ApiCallContext callContext) {
7783
this.uploadChunkCallable =
7884
checkNotNull(uploadChunkCallable, "uploadChunkCallable must not be null");
85+
this.queryStatusCallable =
86+
checkNotNull(queryStatusCallable, "queryStatusCallable must not be null");
7987
this.uploadUrl = checkNotNull(uploadUrl, "uploadUrl must not be null");
8088
checkNotNull(payload, "payload must not be null");
8189
this.callContext = checkNotNull(callContext, "callContext must not be null");
@@ -100,37 +108,25 @@ private void transmitChunk(long currentOffset) {
100108
if (result.isDone()) {
101109
return;
102110
}
103-
104111
// Read the next chunk slice from the payload stream.
105112
try {
106113
buffer.fill(currentOffset);
107-
} catch (IOException e) {
108-
result.setException(e);
109-
return;
114+
dispatchCurrentChunk();
115+
} catch (Throwable t) {
116+
result.setException(t);
110117
}
118+
}
111119

112-
// Determine if this is the final chunk and build the chunk request.
113-
ChunkUploadRequest chunkRequest =
114-
ChunkUploadRequest.newBuilder()
115-
.setUploadUrl(uploadUrl)
116-
.setPayload(buffer.getBuffer())
117-
.setPayloadLength(buffer.getPayloadLength())
118-
.setOffset(buffer.getBufferBaseOffset())
119-
.setFinal(buffer.isFinal())
120-
.build();
121-
122-
// Dispatch the chunk upload call and register the in-flight future for cancellation.
123-
long chunkLength = chunkRequest.getPayloadLength();
124-
boolean isFinal = chunkRequest.isFinal();
120+
private void dispatchCurrentChunk() {
125121
try {
122+
ChunkUploadRequest chunkRequest = buildCurrentChunkRequest();
123+
124+
// Dispatch the chunk upload call and register the in-flight future for cancellation.
126125
ApiFuture<ChunkUploadResponse<ResponseT>> chunkFuture =
127126
uploadChunkCallable.futureCall(chunkRequest, callContext);
128-
this.currentChunkFuture = chunkFuture;
129-
if (result.isCancelled()) {
130-
chunkFuture.cancel(true);
127+
if (!tryRegisterInFlight(chunkFuture)) {
131128
return;
132129
}
133-
134130
ApiFutures.addCallback(
135131
chunkFuture,
136132
new ApiFutureCallback<ChunkUploadResponse<ResponseT>>() {
@@ -139,17 +135,55 @@ public void onSuccess(ChunkUploadResponse<ResponseT> response) {
139135
if (result.isDone()) {
140136
return;
141137
}
142-
long nextOffset = currentOffset + chunkLength;
143-
if (response.getUploadStatus() == ResumableUploadStatus.FINAL) {
144-
result.set(response.getResponse());
145-
} else if (isFinal) {
146-
result.setException(
147-
new IllegalStateException(
148-
"Upload stream ended and final chunk was transmitted, but server returned"
149-
+ " incomplete status for upload URL: "
150-
+ uploadUrl));
138+
if (response.getUploadStatus() == ResumableUploadStatus.UNKNOWN) {
139+
recover();
151140
} else {
152-
chunkExecutor.execute(() -> transmitChunk(nextOffset));
141+
onChunkUploaded(response);
142+
}
143+
}
144+
145+
@Override
146+
public void onFailure(Throwable t) {
147+
if (t instanceof CancellationException || result.isDone()) {
148+
return;
149+
}
150+
Category category =
151+
ResumableUploadErrorClassifier.classify(t, ResumableUploadCommand.UPLOAD);
152+
if (category == Category.RECOVERABLE) {
153+
recover();
154+
} else {
155+
// Category.TRANSIENT errors reaching here have already exhausted their retry budget
156+
// in the underlying RetryingCallable and become fatal per protocol specification.
157+
result.setException(t);
158+
}
159+
}
160+
},
161+
chunkExecutor);
162+
} catch (Throwable t) {
163+
result.setException(t);
164+
}
165+
}
166+
167+
private void recover() {
168+
try {
169+
// Dispatch the query status call and register the in-flight future for cancellation.
170+
ApiFuture<QueryStatusResponse<ResponseT>> queryFuture =
171+
queryStatusCallable.futureCall(QueryStatusRequest.create(uploadUrl), callContext);
172+
if (!tryRegisterInFlight(queryFuture)) {
173+
return;
174+
}
175+
ApiFutures.addCallback(
176+
queryFuture,
177+
new ApiFutureCallback<QueryStatusResponse<ResponseT>>() {
178+
@Override
179+
public void onSuccess(QueryStatusResponse<ResponseT> queryResponse) {
180+
if (result.isDone()) {
181+
return;
182+
}
183+
try {
184+
handleQueryResponse(queryResponse);
185+
} catch (Throwable t) {
186+
result.setException(t);
153187
}
154188
}
155189

@@ -161,10 +195,83 @@ public void onFailure(Throwable t) {
161195
result.setException(t);
162196
}
163197
},
164-
MoreExecutors.directExecutor());
198+
chunkExecutor);
165199
} catch (Throwable t) {
166200
result.setException(t);
167201
}
168202
}
169-
}
170203

204+
/** Registers the in-flight future for cancellation, returning false if already cancelled. */
205+
private boolean tryRegisterInFlight(ApiFuture<?> future) {
206+
this.currentChunkFuture = future;
207+
if (result.isCancelled()) {
208+
future.cancel(true);
209+
return false;
210+
}
211+
return true;
212+
}
213+
214+
private void handleQueryResponse(QueryStatusResponse<ResponseT> queryResponse)
215+
throws IOException {
216+
if (queryResponse.getUploadStatus() == ResumableUploadStatus.UNKNOWN) {
217+
throw protocolViolation(
218+
"Query status response missing X-Goog-Upload-Status header for upload URL: " + uploadUrl);
219+
}
220+
if (queryResponse.getUploadStatus() == ResumableUploadStatus.FINAL) {
221+
onChunkUploaded(
222+
ChunkUploadResponse.create(ResumableUploadStatus.FINAL, queryResponse.getResponse()));
223+
return;
224+
}
225+
Long committedOffset = queryResponse.getCommittedOffset();
226+
if (committedOffset == null) {
227+
throw protocolViolation(
228+
"Incomplete query status response did not include a committed offset for upload URL: "
229+
+ uploadUrl);
230+
}
231+
buffer.realignTo(committedOffset);
232+
dispatchCurrentChunk();
233+
}
234+
235+
private void onChunkUploaded(ChunkUploadResponse<ResponseT> response) {
236+
long nextOffset = buffer.getBufferBaseOffset() + buffer.getPayloadLength();
237+
if (response.getUploadStatus() == ResumableUploadStatus.FINAL) {
238+
result.set(response.getResponse());
239+
} else if (buffer.isFinal()) {
240+
result.setException(
241+
new IllegalStateException(
242+
"Upload stream ended and final chunk was transmitted, but server returned"
243+
+ " incomplete status for upload URL: "
244+
+ uploadUrl));
245+
} else {
246+
chunkExecutor.execute(() -> transmitChunk(nextOffset));
247+
}
248+
}
249+
250+
private ChunkUploadRequest buildCurrentChunkRequest() {
251+
// Determine if this is the final chunk and build the chunk request.
252+
return ChunkUploadRequest.newBuilder()
253+
.setUploadUrl(uploadUrl)
254+
.setPayload(buffer.getBuffer())
255+
.setPayloadLength(buffer.getPayloadLength())
256+
.setOffset(buffer.getBufferBaseOffset())
257+
.setFinal(buffer.isFinal())
258+
.build();
259+
}
260+
261+
private static final StatusCode FAILED_PRECONDITION_STATUS_CODE =
262+
new StatusCode() {
263+
@Override
264+
public StatusCode.Code getCode() {
265+
return StatusCode.Code.FAILED_PRECONDITION;
266+
}
267+
268+
@Override
269+
public @Nullable Object getTransportCode() {
270+
return null;
271+
}
272+
};
273+
274+
private static FailedPreconditionException protocolViolation(String message) {
275+
return new FailedPreconditionException(message, null, FAILED_PRECONDITION_STATUS_CODE, false);
276+
}
277+
}

sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ResumableUploadFutureImpl.java

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,8 @@
3838
import com.google.api.core.SettableApiFuture;
3939
import com.google.api.gax.resumable.ChunkUploadRequest;
4040
import com.google.api.gax.resumable.ChunkUploadResponse;
41+
import com.google.api.gax.resumable.QueryStatusRequest;
42+
import com.google.api.gax.resumable.QueryStatusResponse;
4143
import com.google.api.gax.resumable.ResumableUploadSession;
4244
import com.google.common.util.concurrent.MoreExecutors;
4345
import com.google.errorprone.annotations.concurrent.GuardedBy;
@@ -65,6 +67,8 @@ final class ResumableUploadFutureImpl<ResponseT> implements ResumableUploadFutur
6567
private final ApiFuture<ResumableUploadSession> startFuture;
6668
private final UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>>
6769
uploadChunkCallable;
70+
private final UnaryCallable<QueryStatusRequest, QueryStatusResponse<ResponseT>>
71+
queryStatusCallable;
6872
private final InputStream payload;
6973
private final ResumableUploadCallSettings settings;
7074
private final ApiCallContext callContext;
@@ -85,12 +89,13 @@ final class ResumableUploadFutureImpl<ResponseT> implements ResumableUploadFutur
8589
static <ResponseT> ResumableUploadFutureImpl<ResponseT> create(
8690
ApiFuture<ResumableUploadSession> startFuture,
8791
UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>> uploadChunkCallable,
92+
UnaryCallable<QueryStatusRequest, QueryStatusResponse<ResponseT>> queryStatusCallable,
8893
InputStream payload,
8994
ResumableUploadCallSettings settings,
9095
ApiCallContext callContext) {
9196
ResumableUploadFutureImpl<ResponseT> future =
9297
new ResumableUploadFutureImpl<>(
93-
startFuture, uploadChunkCallable, payload, settings, callContext);
98+
startFuture, uploadChunkCallable, queryStatusCallable, payload, settings, callContext);
9499
try {
95100
future.start();
96101
} catch (Throwable t) {
@@ -102,12 +107,15 @@ static <ResponseT> ResumableUploadFutureImpl<ResponseT> create(
102107
private ResumableUploadFutureImpl(
103108
ApiFuture<ResumableUploadSession> startFuture,
104109
UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>> uploadChunkCallable,
110+
UnaryCallable<QueryStatusRequest, QueryStatusResponse<ResponseT>> queryStatusCallable,
105111
InputStream payload,
106112
ResumableUploadCallSettings settings,
107113
ApiCallContext callContext) {
108114
this.startFuture = checkNotNull(startFuture, "startFuture must not be null");
109115
this.uploadChunkCallable =
110116
checkNotNull(uploadChunkCallable, "uploadChunkCallable must not be null");
117+
this.queryStatusCallable =
118+
checkNotNull(queryStatusCallable, "queryStatusCallable must not be null");
111119
this.payload = checkNotNull(payload, "payload must not be null");
112120
this.settings = checkNotNull(settings, "settings must not be null");
113121
checkArgument(settings.getChunkSize() > 0, "chunkSize must be > 0");
@@ -128,6 +136,7 @@ public void onSuccess(ResumableUploadSession session) {
128136
ResumableUploadChunkCoordinator<ResponseT> coordinator =
129137
new ResumableUploadChunkCoordinator<>(
130138
uploadChunkCallable,
139+
queryStatusCallable,
131140
uploadSessionUrl,
132141
payload,
133142
settings.getChunkSize(),

sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/CallableTest.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -216,6 +216,8 @@ void testResumableUploadCallable() {
216216
mock(ResumableUploadClient.class, Mockito.withSettings().withoutAnnotations());
217217
when(uploadClient.uploadChunkCallable())
218218
.thenReturn(mock(UnaryCallable.class, Mockito.withSettings().withoutAnnotations()));
219+
when(uploadClient.queryStatusCallable())
220+
.thenReturn(mock(UnaryCallable.class, Mockito.withSettings().withoutAnnotations()));
219221
ResumableUploadCallSettings settings =
220222
ResumableUploadCallSettings.newBuilder().setChunkSize(1024).build();
221223

0 commit comments

Comments
 (0)