Skip to content

Commit 387f05d

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 3ae0f6e commit 387f05d

5 files changed

Lines changed: 437 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: 137 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;
@@ -50,6 +53,9 @@
5053
/**
5154
* Coordinates chunk transmission steps of a resumable upload session.
5255
*
56+
* <p>Expects {@code uploadChunkCallable} and {@code queryStatusCallable} to be pre-wrapped in
57+
* retrying callables that handle transient errors.
58+
*
5359
* @param <ResponseT> the type of the final response message returned once the upload completes
5460
*/
5561
@InternalApi
@@ -62,6 +68,8 @@ final class ResumableUploadChunkCoordinator<ResponseT> {
6268

6369
private final UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>>
6470
uploadChunkCallable;
71+
private final UnaryCallable<QueryStatusRequest, QueryStatusResponse<ResponseT>>
72+
queryStatusCallable;
6573
private final String uploadUrl;
6674
private final RewindableStreamBuffer buffer;
6775
private final ApiCallContext callContext;
@@ -70,12 +78,15 @@ final class ResumableUploadChunkCoordinator<ResponseT> {
7078

7179
ResumableUploadChunkCoordinator(
7280
UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>> uploadChunkCallable,
81+
UnaryCallable<QueryStatusRequest, QueryStatusResponse<ResponseT>> queryStatusCallable,
7382
String uploadUrl,
7483
InputStream payload,
7584
int chunkSize,
7685
ApiCallContext callContext) {
7786
this.uploadChunkCallable =
7887
checkNotNull(uploadChunkCallable, "uploadChunkCallable must not be null");
88+
this.queryStatusCallable =
89+
checkNotNull(queryStatusCallable, "queryStatusCallable must not be null");
7990
this.uploadUrl = checkNotNull(uploadUrl, "uploadUrl must not be null");
8091
checkNotNull(payload, "payload must not be null");
8192
this.callContext = checkNotNull(callContext, "callContext must not be null");
@@ -96,60 +107,85 @@ ApiFuture<ResponseT> start() {
96107
}
97108

98109
private void transmitChunk(long currentOffset) {
99-
// Abort if the session was already completed or canceled.
100110
if (result.isDone()) {
101111
return;
102112
}
103-
104-
// Read the next chunk slice from the payload stream.
105113
try {
106114
buffer.fill(currentOffset);
107-
} catch (IOException e) {
108-
result.setException(e);
109-
return;
115+
dispatchCurrentChunk();
116+
} catch (Throwable t) {
117+
result.setException(t);
110118
}
119+
}
111120

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();
121+
private void dispatchCurrentChunk() {
122+
if (result.isDone()) {
123+
return;
124+
}
125125
try {
126+
ChunkUploadRequest chunkRequest = buildCurrentChunkRequest();
126127
ApiFuture<ChunkUploadResponse<ResponseT>> chunkFuture =
127128
uploadChunkCallable.futureCall(chunkRequest, callContext);
128129
this.currentChunkFuture = chunkFuture;
129130
if (result.isCancelled()) {
130131
chunkFuture.cancel(true);
131132
return;
132133
}
133-
134134
ApiFutures.addCallback(
135135
chunkFuture,
136136
new ApiFutureCallback<ChunkUploadResponse<ResponseT>>() {
137137
@Override
138138
public void onSuccess(ChunkUploadResponse<ResponseT> response) {
139-
if (result.isDone()) {
139+
if (response.getUploadStatus() == ResumableUploadStatus.UNKNOWN) {
140+
recover();
141+
} else {
142+
onChunkUploaded(response);
143+
}
144+
}
145+
146+
@Override
147+
public void onFailure(Throwable t) {
148+
if (t instanceof CancellationException || result.isDone()) {
140149
return;
141150
}
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));
151+
Category category =
152+
ResumableUploadErrorClassifier.classify(t, ResumableUploadCommand.UPLOAD);
153+
if (category == Category.RECOVERABLE) {
154+
recover();
151155
} else {
152-
chunkExecutor.execute(() -> transmitChunk(nextOffset));
156+
// Category.TRANSIENT errors reaching here have already exhausted their retry budget
157+
// in the underlying RetryingCallable and become fatal per protocol specification.
158+
result.setException(t);
159+
}
160+
}
161+
},
162+
chunkExecutor);
163+
} catch (Throwable t) {
164+
result.setException(t);
165+
}
166+
}
167+
168+
private void recover() {
169+
if (result.isDone()) {
170+
return;
171+
}
172+
try {
173+
ApiFuture<QueryStatusResponse<ResponseT>> queryFuture =
174+
queryStatusCallable.futureCall(QueryStatusRequest.create(uploadUrl), callContext);
175+
this.currentChunkFuture = queryFuture;
176+
if (result.isCancelled()) {
177+
queryFuture.cancel(true);
178+
return;
179+
}
180+
ApiFutures.addCallback(
181+
queryFuture,
182+
new ApiFutureCallback<QueryStatusResponse<ResponseT>>() {
183+
@Override
184+
public void onSuccess(QueryStatusResponse<ResponseT> queryResponse) {
185+
try {
186+
handleQueryResponse(queryResponse);
187+
} catch (Throwable t) {
188+
result.setException(t);
153189
}
154190
}
155191

@@ -161,10 +197,78 @@ public void onFailure(Throwable t) {
161197
result.setException(t);
162198
}
163199
},
164-
MoreExecutors.directExecutor());
200+
chunkExecutor);
165201
} catch (Throwable t) {
166202
result.setException(t);
167203
}
168204
}
169-
}
170205

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

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

Lines changed: 15 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,18 @@ 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,
99+
uploadChunkCallable,
100+
queryStatusCallable,
101+
payload,
102+
settings,
103+
callContext);
94104
try {
95105
future.start();
96106
} catch (Throwable t) {
@@ -102,12 +112,15 @@ static <ResponseT> ResumableUploadFutureImpl<ResponseT> create(
102112
private ResumableUploadFutureImpl(
103113
ApiFuture<ResumableUploadSession> startFuture,
104114
UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>> uploadChunkCallable,
115+
UnaryCallable<QueryStatusRequest, QueryStatusResponse<ResponseT>> queryStatusCallable,
105116
InputStream payload,
106117
ResumableUploadCallSettings settings,
107118
ApiCallContext callContext) {
108119
this.startFuture = checkNotNull(startFuture, "startFuture must not be null");
109120
this.uploadChunkCallable =
110121
checkNotNull(uploadChunkCallable, "uploadChunkCallable must not be null");
122+
this.queryStatusCallable =
123+
checkNotNull(queryStatusCallable, "queryStatusCallable must not be null");
111124
this.payload = checkNotNull(payload, "payload must not be null");
112125
this.settings = checkNotNull(settings, "settings must not be null");
113126
checkArgument(settings.getChunkSize() > 0, "chunkSize must be > 0");
@@ -128,6 +141,7 @@ public void onSuccess(ResumableUploadSession session) {
128141
ResumableUploadChunkCoordinator<ResponseT> coordinator =
129142
new ResumableUploadChunkCoordinator<>(
130143
uploadChunkCallable,
144+
queryStatusCallable,
131145
uploadSessionUrl,
132146
payload,
133147
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)