Skip to content

Commit cf10342

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

5 files changed

Lines changed: 531 additions & 19 deletions

File tree

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

Lines changed: 12 additions & 1 deletion
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,
@@ -89,6 +93,8 @@ public ResumableUploadCallableImpl(
8993
.build();
9094
this.retryingUploadChunkCallable =
9195
createRetryingCallable(client.uploadChunkCallable(), ResumableUploadCommand.UPLOAD);
96+
this.retryingQueryCallable =
97+
createRetryingCallable(client.queryStatusCallable(), ResumableUploadCommand.QUERY);
9298
}
9399

94100
@Override
@@ -110,7 +116,12 @@ public ResumableUploadFuture<ResponseT> futureCall(
110116
}
111117

112118
return ResumableUploadFutureImpl.create(
113-
startFuture, retryingUploadChunkCallable, payload, effectiveSettings, clientContext);
119+
startFuture,
120+
retryingUploadChunkCallable,
121+
retryingQueryCallable,
122+
payload,
123+
effectiveSettings,
124+
clientContext);
114125
}
115126

116127
@Override

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

Lines changed: 126 additions & 9 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;
@@ -56,25 +59,44 @@
5659
@NullMarked
5760
final class ResumableUploadChunkCoordinator<ResponseT> {
5861

62+
private static final StatusCode FAILED_PRECONDITION_STATUS_CODE =
63+
new StatusCode() {
64+
@Override
65+
public StatusCode.Code getCode() {
66+
return StatusCode.Code.FAILED_PRECONDITION;
67+
}
68+
69+
@Override
70+
public @Nullable Object getTransportCode() {
71+
return null;
72+
}
73+
};
74+
5975
private final Executor chunkExecutor =
6076
MoreExecutors.newSequentialExecutor(MoreExecutors.directExecutor());
6177

6278
private final UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>>
6379
uploadChunkCallable;
80+
private final UnaryCallable<QueryStatusRequest, QueryStatusResponse<ResponseT>>
81+
queryStatusCallable;
6482
private final String uploadUrl;
6583
private final RewindableStreamBuffer buffer;
6684
private final ApiCallContext callContext;
6785
private final SettableApiFuture<ResponseT> result = SettableApiFuture.create();
6886
private volatile @Nullable ApiFuture<?> currentChunkFuture;
87+
private int recoveryAttempts;
6988

7089
ResumableUploadChunkCoordinator(
7190
UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>> uploadChunkCallable,
91+
UnaryCallable<QueryStatusRequest, QueryStatusResponse<ResponseT>> queryStatusCallable,
7292
String uploadUrl,
7393
InputStream payload,
7494
int chunkSize,
7595
ClientContext clientContext) {
7696
this.uploadChunkCallable =
7797
checkNotNull(uploadChunkCallable, "uploadChunkCallable must not be null");
98+
this.queryStatusCallable =
99+
checkNotNull(queryStatusCallable, "queryStatusCallable must not be null");
78100
this.uploadUrl = checkNotNull(uploadUrl, "uploadUrl must not be null");
79101
checkNotNull(payload, "payload must not be null");
80102
checkNotNull(clientContext, "clientContext must not be null");
@@ -96,20 +118,26 @@ ApiFuture<ResponseT> start() {
96118
}
97119

98120
private void transmitChunk(long currentOffset) {
99-
// Abort if the session was already completed or canceled.
100121
if (result.isDone()) {
101122
return;
102123
}
103124

104-
// Read the next chunk slice from the payload stream.
105125
try {
106126
buffer.fill(currentOffset);
107127
} catch (IOException e) {
108128
result.setException(e);
109129
return;
110130
}
111131

112-
// Determine if this is the final chunk and build the chunk request.
132+
recoveryAttempts = 0;
133+
sendCurrentBuffer();
134+
}
135+
136+
private void sendCurrentBuffer() {
137+
if (result.isDone()) {
138+
return;
139+
}
140+
113141
ChunkUploadRequest chunkRequest =
114142
ChunkUploadRequest.newBuilder()
115143
.setUploadUrl(uploadUrl)
@@ -119,9 +147,6 @@ private void transmitChunk(long currentOffset) {
119147
.setFinal(buffer.isFinal())
120148
.build();
121149

122-
// Dispatch the chunk upload call and register the in-flight future for cancellation.
123-
long chunkLength = chunkRequest.getPayloadLength();
124-
boolean isFinal = chunkRequest.isFinal();
125150
try {
126151
ApiFuture<ChunkUploadResponse<ResponseT>> chunkFuture =
127152
uploadChunkCallable.futureCall(chunkRequest, callContext);
@@ -139,16 +164,18 @@ public void onSuccess(ChunkUploadResponse<ResponseT> response) {
139164
if (result.isDone()) {
140165
return;
141166
}
142-
long nextOffset = currentOffset + chunkLength;
143-
if (response.getUploadStatus() == ResumableUploadStatus.FINAL) {
167+
if (response.getUploadStatus() == ResumableUploadStatus.UNKNOWN) {
168+
chunkExecutor.execute(() -> recoverChunk(null));
169+
} else if (response.getUploadStatus() == ResumableUploadStatus.FINAL) {
144170
result.set(response.getResponse());
145-
} else if (isFinal) {
171+
} else if (buffer.isFinal()) {
146172
result.setException(
147173
new IllegalStateException(
148174
"Upload stream ended and final chunk was transmitted, but server returned"
149175
+ " incomplete status for upload URL: "
150176
+ uploadUrl));
151177
} else {
178+
long nextOffset = buffer.getBufferBaseOffset() + buffer.getPayloadLength();
152179
chunkExecutor.execute(() -> transmitChunk(nextOffset));
153180
}
154181
}
@@ -158,6 +185,11 @@ public void onFailure(Throwable t) {
158185
if (t instanceof CancellationException || result.isDone()) {
159186
return;
160187
}
188+
if (ResumableUploadErrorClassifier.classify(t, ResumableUploadCommand.UPLOAD)
189+
== Category.RECOVERABLE) {
190+
chunkExecutor.execute(() -> recoverChunk(t));
191+
return;
192+
}
161193
result.setException(t);
162194
}
163195
},
@@ -166,4 +198,89 @@ public void onFailure(Throwable t) {
166198
result.setException(t);
167199
}
168200
}
201+
202+
private void recoverChunk(@Nullable Throwable failure) {
203+
if (result.isDone()) {
204+
return;
205+
}
206+
int maxAttempts =
207+
callContext.getRetrySettings() != null
208+
? callContext.getRetrySettings().getMaxAttempts()
209+
: 5;
210+
if (++recoveryAttempts >= maxAttempts) {
211+
result.setException(
212+
failure != null
213+
? failure
214+
: protocolViolation(
215+
"Chunk upload response missing X-Goog-Upload-Status header for upload URL: "
216+
+ uploadUrl));
217+
return;
218+
}
219+
220+
try {
221+
ApiFuture<QueryStatusResponse<ResponseT>> queryFuture =
222+
queryStatusCallable.futureCall(QueryStatusRequest.create(uploadUrl), callContext);
223+
this.currentChunkFuture = queryFuture;
224+
if (result.isCancelled()) {
225+
queryFuture.cancel(true);
226+
return;
227+
}
228+
229+
ApiFutures.addCallback(
230+
queryFuture,
231+
new ApiFutureCallback<QueryStatusResponse<ResponseT>>() {
232+
@Override
233+
public void onSuccess(QueryStatusResponse<ResponseT> queryResponse) {
234+
chunkExecutor.execute(() -> handleQueryResponse(queryResponse));
235+
}
236+
237+
@Override
238+
public void onFailure(Throwable t) {
239+
if (t instanceof CancellationException || result.isDone()) {
240+
return;
241+
}
242+
result.setException(t);
243+
}
244+
},
245+
MoreExecutors.directExecutor());
246+
} catch (Throwable t) {
247+
result.setException(t);
248+
}
249+
}
250+
251+
private void handleQueryResponse(QueryStatusResponse<ResponseT> queryResponse) {
252+
if (result.isDone()) {
253+
return;
254+
}
255+
if (queryResponse.getUploadStatus() == ResumableUploadStatus.UNKNOWN) {
256+
result.setException(
257+
protocolViolation(
258+
"Query status response missing X-Goog-Upload-Status header for upload URL: "
259+
+ uploadUrl));
260+
return;
261+
}
262+
if (queryResponse.getUploadStatus() == ResumableUploadStatus.FINAL) {
263+
result.set(queryResponse.getResponse());
264+
return;
265+
}
266+
Long committedOffset = queryResponse.getCommittedOffset();
267+
if (committedOffset == null) {
268+
result.setException(
269+
protocolViolation(
270+
"Incomplete query status response did not include a committed offset for upload URL: "
271+
+ uploadUrl));
272+
return;
273+
}
274+
try {
275+
buffer.realignTo(committedOffset);
276+
} catch (Throwable t) {
277+
result.setException(t);
278+
return;
279+
}
280+
sendCurrentBuffer();
281+
}
282+
283+
private static FailedPreconditionException protocolViolation(String message) {
284+
return new FailedPreconditionException(message, null, FAILED_PRECONDITION_STATUS_CODE, false);
285+
}
169286
}

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

Lines changed: 15 additions & 8 deletions
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 ClientContext clientContext;
@@ -75,22 +79,21 @@ final class ResumableUploadFutureImpl<ResponseT> implements ResumableUploadFutur
7579
@GuardedBy("lock")
7680
private @Nullable ApiFuture<?> inFlightFuture;
7781

78-
/**
79-
* Creates and initiates a new resumable upload future tracking session initiation and chunk
80-
* streaming.
81-
*
82-
* <p>The provided {@code payload} stream is managed by the returned future and will be closed
83-
* automatically upon completion, failure, or cancellation.
84-
*/
8582
static <ResponseT> ResumableUploadFutureImpl<ResponseT> create(
8683
ApiFuture<ResumableUploadSession> startFuture,
8784
UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>> uploadChunkCallable,
85+
UnaryCallable<QueryStatusRequest, QueryStatusResponse<ResponseT>> queryStatusCallable,
8886
InputStream payload,
8987
ResumableUploadCallSettings settings,
9088
ClientContext clientContext) {
9189
ResumableUploadFutureImpl<ResponseT> future =
9290
new ResumableUploadFutureImpl<>(
93-
startFuture, uploadChunkCallable, payload, settings, clientContext);
91+
startFuture,
92+
uploadChunkCallable,
93+
queryStatusCallable,
94+
payload,
95+
settings,
96+
clientContext);
9497
try {
9598
future.start();
9699
} catch (Throwable t) {
@@ -102,12 +105,15 @@ static <ResponseT> ResumableUploadFutureImpl<ResponseT> create(
102105
private ResumableUploadFutureImpl(
103106
ApiFuture<ResumableUploadSession> startFuture,
104107
UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>> uploadChunkCallable,
108+
UnaryCallable<QueryStatusRequest, QueryStatusResponse<ResponseT>> queryStatusCallable,
105109
InputStream payload,
106110
ResumableUploadCallSettings settings,
107111
ClientContext clientContext) {
108112
this.startFuture = checkNotNull(startFuture, "startFuture must not be null");
109113
this.uploadChunkCallable =
110114
checkNotNull(uploadChunkCallable, "uploadChunkCallable must not be null");
115+
this.queryStatusCallable =
116+
checkNotNull(queryStatusCallable, "queryStatusCallable must not be null");
111117
this.payload = checkNotNull(payload, "payload must not be null");
112118
this.settings = checkNotNull(settings, "settings must not be null");
113119
checkArgument(settings.getChunkSize() > 0, "chunkSize must be > 0");
@@ -128,6 +134,7 @@ public void onSuccess(ResumableUploadSession session) {
128134
ResumableUploadChunkCoordinator<ResponseT> coordinator =
129135
new ResumableUploadChunkCoordinator<>(
130136
uploadChunkCallable,
137+
queryStatusCallable,
131138
uploadSessionUrl,
132139
payload,
133140
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)