Skip to content

Commit 3689a54

Browse files
committed
feat(gax): wire progress listener into coordinator and attempt lifecycle
1 parent 7b55d31 commit 3689a54

7 files changed

Lines changed: 357 additions & 8 deletions

File tree

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

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -71,6 +71,7 @@ class ChunkAttemptCallable<ResponseT> implements Callable<ChunkUploadResponse<Re
7171
private final RewindableStreamBuffer buffer;
7272
private final String uploadUrl;
7373
private final ApiCallContext originalCallContext;
74+
private final UploadProgressTracker progressTracker;
7475
private final long deadlineNanos;
7576
private final ApiClock clock;
7677

@@ -89,7 +90,8 @@ class ChunkAttemptCallable<ResponseT> implements Callable<ChunkUploadResponse<Re
8990
String uploadUrl,
9091
ChunkUploadRequest request,
9192
ApiCallContext callContext,
92-
UploadCommand command) {
93+
UploadCommand command,
94+
UploadProgressTracker progressTracker) {
9395
this(
9496
uploadChunkCallable,
9597
queryStatusCallable,
@@ -98,6 +100,7 @@ class ChunkAttemptCallable<ResponseT> implements Callable<ChunkUploadResponse<Re
98100
request,
99101
callContext,
100102
command,
103+
progressTracker,
101104
Long.MAX_VALUE,
102105
NanoClock.getDefaultClock());
103106
}
@@ -110,6 +113,7 @@ class ChunkAttemptCallable<ResponseT> implements Callable<ChunkUploadResponse<Re
110113
ChunkUploadRequest request,
111114
ApiCallContext callContext,
112115
UploadCommand command,
116+
UploadProgressTracker progressTracker,
113117
long deadlineNanos,
114118
ApiClock clock) {
115119
this.uploadChunkCallable =
@@ -121,6 +125,7 @@ class ChunkAttemptCallable<ResponseT> implements Callable<ChunkUploadResponse<Re
121125
this.currentRequest = checkNotNull(request, "request must not be null");
122126
this.originalCallContext = checkNotNull(callContext, "callContext must not be null");
123127
this.currentCommand = checkNotNull(command, "command must not be null");
128+
this.progressTracker = checkNotNull(progressTracker, "progressTracker must not be null");
124129
this.deadlineNanos = deadlineNanos;
125130
this.clock = checkNotNull(clock, "clock must not be null");
126131
}
@@ -158,6 +163,8 @@ private void prepareAttempt(
158163
SettableApiFuture<ChunkUploadResponse<ResponseT>> attemptFuture,
159164
ApiCallContext attemptContext,
160165
RetryingFuture<ChunkUploadResponse<ResponseT>> currentRetryingFuture) {
166+
progressTracker.onRecovering(lastFailure);
167+
161168
// Per GAX-R7: query uses sensible unary defaults trimmed to the remaining global deadline.
162169
long remainingNanos =
163170
deadlineNanos == Long.MAX_VALUE
@@ -254,6 +261,7 @@ private void handleQuerySuccess(
254261
// Normal path: realign buffer to committedOffset, compact and top up.
255262
try {
256263
buffer.realignTo(committedOffset);
264+
progressTracker.onOffsetReceived(committedOffset);
257265
} catch (Throwable e) {
258266
failAttempt(attemptFuture, e);
259267
return;

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

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,7 @@
5252
import java.io.InputStream;
5353
import java.time.Duration;
5454
import java.util.concurrent.CancellationException;
55+
import java.util.concurrent.Executor;
5556
import java.util.concurrent.ScheduledExecutorService;
5657
import java.util.concurrent.ScheduledFuture;
5758
import java.util.concurrent.TimeUnit;
@@ -94,6 +95,7 @@ final class ResumableUploadChunkCoordinator<ResponseT> {
9495
private final int chunkSize;
9596
private final ApiCallContext callContext;
9697
private final ClientContext clientContext;
98+
private final UploadProgressTracker progressTracker = new UploadProgressTracker();
9799

98100
private volatile @Nullable String uploadSessionUrl;
99101
private volatile @Nullable RewindableStreamBuffer buffer;
@@ -166,6 +168,7 @@ public void onSuccess(ResumableUploadSession session) {
166168
}
167169
}
168170
uploadSessionUrl = session.getUploadUrl();
171+
progressTracker.onStarted(uploadSessionUrl);
169172
buffer = new RewindableStreamBuffer(payload, chunkSize, uploadSessionUrl);
170173
scheduleNextChunk(0L);
171174
}
@@ -181,6 +184,14 @@ public void onFailure(Throwable t) {
181184
MoreExecutors.directExecutor());
182185
}
183186

187+
void addProgressListener(ResumableUploadProgressListener listener, Executor executor) {
188+
progressTracker.addListener(listener, executor);
189+
}
190+
191+
ResumableUploadStatus getStatus() {
192+
return progressTracker.getStatus();
193+
}
194+
184195
private void onTimeout() {
185196
synchronized (lock) {
186197
if (done) {
@@ -233,6 +244,7 @@ void cancel(boolean mayInterruptIfRunning) {
233244
if (inFlight != null) {
234245
inFlight.cancel(mayInterruptIfRunning);
235246
}
247+
progressTracker.onFailed(new CancellationException("Upload was cancelled"), uploadSessionUrl);
236248
closePayload();
237249
}
238250

@@ -257,11 +269,15 @@ private void finish(@Nullable ResponseT response, @Nullable Throwable error) {
257269
}
258270
IOException closeError = closePayload();
259271
if (error == null) {
272+
long totalBytes =
273+
buffer != null ? buffer.getBufferBaseOffset() + buffer.getPayloadLength() : 0L;
274+
progressTracker.onFinalized(totalBytes);
260275
result.set(response);
261276
} else {
262277
if (closeError != null) {
263278
error.addSuppressed(closeError);
264279
}
280+
progressTracker.onFailed(error, uploadSessionUrl);
265281
result.setException(error);
266282
}
267283
}
@@ -364,6 +380,7 @@ private void transmitSingleChunk(long currentOffset) {
364380
chunkRequest,
365381
chunkCallContext,
366382
command,
383+
progressTracker,
367384
deadlineNanos,
368385
clientContext.getClock());
369386

@@ -383,6 +400,7 @@ public void onSuccess(ChunkUploadResponse<ResponseT> response) {
383400
}
384401
}
385402
long nextOffset = streamBuffer.getBufferBaseOffset() + streamBuffer.getPayloadLength();
403+
progressTracker.onChunkUploaded(nextOffset);
386404
if (response.isComplete()) {
387405
finish(response.getResponse(), null);
388406
} else if (streamBuffer.isFinal()) {

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

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -31,6 +31,7 @@
3131

3232
import com.google.api.core.ApiFuture;
3333
import com.google.api.core.BetaApi;
34+
import java.util.concurrent.Executor;
3435
import org.jspecify.annotations.NullMarked;
3536
import org.jspecify.annotations.Nullable;
3637

@@ -48,4 +49,18 @@ public interface ResumableUploadFuture<ResponseT> extends ApiFuture<ResponseT> {
4849

4950
/** Returns the upload session URL, or {@code null} if session initiation is in progress. */
5051
@Nullable String getUploadSessionUrl();
52+
53+
/**
54+
* Registers a listener to receive progress and state transition notifications for this upload.
55+
*
56+
* <p>A snapshot of the current upload status is dispatched to the listener immediately upon
57+
* subscription on the provided executor. Subsequent status updates are delivered in order.
58+
*
59+
* @param listener callback listener to receive progress notifications
60+
* @param executor executor on which the listener callbacks are dispatched
61+
*/
62+
void addProgressListener(ResumableUploadProgressListener listener, Executor executor);
63+
64+
/** Returns the current status snapshot of the upload session. */
65+
ResumableUploadStatus getStatus();
5166
}

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

Lines changed: 12 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -93,6 +93,16 @@ static <ResponseT> ResumableUploadFutureImpl<ResponseT> create(
9393
return coordinator.getUploadSessionUrl();
9494
}
9595

96+
@Override
97+
public void addProgressListener(ResumableUploadProgressListener listener, Executor executor) {
98+
coordinator.addProgressListener(listener, executor);
99+
}
100+
101+
@Override
102+
public ResumableUploadStatus getStatus() {
103+
return coordinator.getStatus();
104+
}
105+
96106
@Override
97107
public void addListener(Runnable listener, Executor executor) {
98108
result.addListener(listener, executor);
@@ -115,12 +125,12 @@ public boolean isDone() {
115125
}
116126

117127
@Override
118-
public ResponseT get() throws InterruptedException, ExecutionException {
128+
public @Nullable ResponseT get() throws InterruptedException, ExecutionException {
119129
return result.get();
120130
}
121131

122132
@Override
123-
public ResponseT get(long timeout, TimeUnit unit)
133+
public @Nullable ResponseT get(long timeout, TimeUnit unit)
124134
throws InterruptedException, ExecutionException, TimeoutException {
125135
return result.get(timeout, unit);
126136
}

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

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -130,7 +130,8 @@ void call_successfulChunk_setsAttemptFuture() throws Exception {
130130
"https://upload.url/test",
131131
request,
132132
callContext,
133-
UploadCommand.UPLOAD);
133+
UploadCommand.UPLOAD,
134+
new UploadProgressTracker());
134135

135136
callable.setRetryingFuture(mockExternalFuture);
136137
ChunkUploadResponse<String> callResult = callable.call();
@@ -178,7 +179,8 @@ void call_returnsWithoutBlocking_andPropagatesCancellation() throws Exception {
178179
"https://upload.url/test",
179180
request,
180181
callContext,
181-
UploadCommand.UPLOAD);
182+
UploadCommand.UPLOAD,
183+
new UploadProgressTracker());
182184

183185
List<Runnable> listeners = new ArrayList<>();
184186
doAnswer(
@@ -251,7 +253,8 @@ void call_perAttemptDeadline_appliesRpcTimeoutToCallContext() throws Exception {
251253
"https://upload.url/test",
252254
request,
253255
callContext,
254-
UploadCommand.UPLOAD);
256+
UploadCommand.UPLOAD,
257+
new UploadProgressTracker());
255258

256259
callable.setRetryingFuture(mockExternalFuture);
257260
callable.call();
@@ -292,7 +295,8 @@ void call_nonBlockingExecution_callingThreadMakesImmediateProgress() throws Exce
292295
"https://upload.url/test",
293296
request,
294297
callContext,
295-
UploadCommand.UPLOAD);
298+
UploadCommand.UPLOAD,
299+
new UploadProgressTracker());
296300

297301
callable.setRetryingFuture(mockExternalFuture);
298302

0 commit comments

Comments
 (0)