Skip to content

Commit fd4401e

Browse files
authored
refactor(gax): remove circular ref between resumable upload future and chunk coordinator (#14421)
This cleans up and helps clarify the layering and responsibility structure ahead of introducing non-happy path features. In general the Future coordinates the overall upload lifecycle while delegating details of specific operations (e.g. chunk uploads, status listeners, global timeout) to the layer below. Actors on that layer don't maintain explicit references to the Future.
1 parent 28baf15 commit fd4401e

3 files changed

Lines changed: 139 additions & 84 deletions

File tree

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

Lines changed: 65 additions & 45 deletions
Original file line numberDiff line numberDiff line change
@@ -35,16 +35,17 @@
3535
import com.google.api.core.ApiFutureCallback;
3636
import com.google.api.core.ApiFutures;
3737
import com.google.api.core.InternalApi;
38+
import com.google.api.core.SettableApiFuture;
3839
import com.google.api.gax.resumable.ChunkUploadRequest;
3940
import com.google.api.gax.resumable.ChunkUploadResponse;
4041
import com.google.api.gax.resumable.ResumableUploadStatus;
4142
import com.google.common.io.ByteStreams;
4243
import com.google.common.util.concurrent.MoreExecutors;
43-
import java.io.IOException;
4444
import java.io.InputStream;
4545
import java.util.Arrays;
4646
import java.util.concurrent.CancellationException;
4747
import org.jspecify.annotations.NullMarked;
48+
import org.jspecify.annotations.Nullable;
4849

4950
/**
5051
* Coordinates chunk transmission steps of a resumable upload session.
@@ -64,84 +65,90 @@ final class ResumableUploadChunkCoordinator<ResponseT> {
6465
private final byte[] buffer;
6566
private final int chunkSize;
6667
private final ApiCallContext callContext;
67-
private final ResumableUploadFutureImpl<ResponseT> sessionFuture;
68+
private final SettableApiFuture<ResponseT> uploadResultFuture = SettableApiFuture.create();
69+
private volatile @Nullable ApiFuture<?> inFlightFuture;
6870

6971
ResumableUploadChunkCoordinator(
7072
UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>> uploadChunkCallable,
7173
String uploadUrl,
7274
InputStream payload,
7375
int chunkSize,
74-
ApiCallContext callContext,
75-
ResumableUploadFutureImpl<ResponseT> sessionFuture) {
76+
ApiCallContext callContext) {
7677
this.uploadChunkCallable =
7778
checkNotNull(uploadChunkCallable, "uploadChunkCallable must not be null");
7879
this.uploadUrl = checkNotNull(uploadUrl, "uploadUrl must not be null");
7980
this.payload = checkNotNull(payload, "payload must not be null");
8081
this.chunkSize = chunkSize;
8182
this.callContext = checkNotNull(callContext, "callContext must not be null");
82-
this.sessionFuture = checkNotNull(sessionFuture, "sessionFuture must not be null");
8383
this.buffer = new byte[chunkSize];
8484
}
8585

86+
ApiFuture<ResponseT> getFuture() {
87+
return uploadResultFuture;
88+
}
89+
8690
void start() {
91+
uploadResultFuture.addListener(
92+
() -> {
93+
ApiFuture<?> inFlight = inFlightFuture;
94+
if (uploadResultFuture.isCancelled() && inFlight != null) {
95+
inFlight.cancel(true);
96+
}
97+
},
98+
MoreExecutors.directExecutor());
8799
transmitChunk(0L);
88100
}
89101

90102
private void transmitChunk(long currentOffset) {
91-
// Abort if the session was already completed or canceled.
92-
if (sessionFuture.isDone()) {
93-
return;
94-
}
95-
96-
// Read the next chunk slice from the payload stream.
97-
int bytesRead;
98103
try {
99-
bytesRead = ByteStreams.read(payload, buffer, 0, chunkSize);
100-
} catch (IOException e) {
101-
sessionFuture.fail(e);
102-
return;
103-
}
104+
// Abort if the session was already completed or canceled.
105+
if (uploadResultFuture.isDone()) {
106+
return;
107+
}
104108

105-
// Determine if this is the final chunk and build the chunk request.
106-
boolean isFinal = bytesRead < chunkSize;
107-
byte[] chunkPayload;
108-
if (bytesRead == chunkSize) {
109-
chunkPayload = buffer;
110-
} else if (bytesRead == 0) {
111-
chunkPayload = EMPTY_PAYLOAD;
112-
} else {
113-
chunkPayload = Arrays.copyOf(buffer, bytesRead);
114-
}
109+
// Read the next chunk slice from the payload stream.
110+
int bytesRead = ByteStreams.read(payload, buffer, 0, chunkSize);
115111

116-
ChunkUploadRequest chunkRequest =
117-
ChunkUploadRequest.newBuilder()
118-
.setUploadUrl(uploadUrl)
119-
.setPayload(chunkPayload)
120-
.setOffset(currentOffset)
121-
.setFinal(isFinal)
122-
.build();
112+
// Determine if this is the final chunk and build the chunk request.
113+
boolean isFinal = bytesRead < chunkSize;
114+
byte[] chunkPayload;
115+
if (bytesRead == chunkSize) {
116+
chunkPayload = buffer;
117+
} else if (bytesRead == 0) {
118+
chunkPayload = EMPTY_PAYLOAD;
119+
} else {
120+
chunkPayload = Arrays.copyOf(buffer, bytesRead);
121+
}
123122

124-
// Dispatch the chunk upload call and register the in-flight future for cancellation.
125-
long chunkLength = chunkPayload.length;
126-
try {
123+
ChunkUploadRequest chunkRequest =
124+
ChunkUploadRequest.newBuilder()
125+
.setUploadUrl(uploadUrl)
126+
.setPayload(chunkPayload)
127+
.setOffset(currentOffset)
128+
.setFinal(isFinal)
129+
.build();
130+
131+
// Dispatch the chunk upload call and register the in-flight future for cancellation.
132+
long chunkLength = chunkPayload.length;
127133
ApiFuture<ChunkUploadResponse<ResponseT>> chunkFuture =
128134
uploadChunkCallable.futureCall(chunkRequest, callContext);
129-
sessionFuture.setInFlightFuture(chunkFuture);
135+
if (!tryRegisterInFlightFuture(chunkFuture)) {
136+
return;
137+
}
130138

131-
// Asynchronously handle the response: complete, fail, or chain the next chunk.
132139
ApiFutures.addCallback(
133140
chunkFuture,
134141
new ApiFutureCallback<ChunkUploadResponse<ResponseT>>() {
135142
@Override
136143
public void onSuccess(ChunkUploadResponse<ResponseT> response) {
137-
if (sessionFuture.isDone()) {
144+
if (uploadResultFuture.isDone()) {
138145
return;
139146
}
140147
long nextOffset = currentOffset + chunkLength;
141148
if (response.getUploadStatus() == ResumableUploadStatus.FINAL) {
142-
sessionFuture.succeed(response.getResponse());
149+
uploadResultFuture.set(response.getResponse());
143150
} else if (isFinal) {
144-
sessionFuture.fail(
151+
uploadResultFuture.setException(
145152
new IllegalStateException(
146153
"Upload stream ended and final chunk was transmitted, but server returned"
147154
+ " incomplete status"));
@@ -152,15 +159,28 @@ public void onSuccess(ChunkUploadResponse<ResponseT> response) {
152159

153160
@Override
154161
public void onFailure(Throwable t) {
155-
if (t instanceof CancellationException || sessionFuture.isDone()) {
162+
if (t instanceof CancellationException || uploadResultFuture.isDone()) {
156163
return;
157164
}
158-
sessionFuture.fail(t);
165+
uploadResultFuture.setException(t);
159166
}
160167
},
161168
MoreExecutors.directExecutor());
162169
} catch (Throwable t) {
163-
sessionFuture.fail(t);
170+
uploadResultFuture.setException(t);
171+
}
172+
}
173+
174+
/**
175+
* Registers the in-flight future for possible cancellation, returning false if the upload was
176+
* already cancelled.
177+
*/
178+
private boolean tryRegisterInFlightFuture(ApiFuture<?> future) {
179+
this.inFlightFuture = future;
180+
if (uploadResultFuture.isCancelled()) {
181+
future.cancel(true);
182+
return false;
164183
}
184+
return true;
165185
}
166186
}

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

Lines changed: 30 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -72,6 +72,7 @@ final class ResumableUploadFutureImpl<ResponseT> implements ResumableUploadFutur
7272

7373
private volatile @Nullable String uploadSessionUrl;
7474

75+
// Tracks the current operation's Future (start, chunk upload) to propagate cancellation.
7576
@GuardedBy("lock")
7677
private @Nullable ApiFuture<?> inFlightFuture;
7778

@@ -121,23 +122,36 @@ private void start() {
121122
new ApiFutureCallback<ResumableUploadSession>() {
122123
@Override
123124
public void onSuccess(ResumableUploadSession session) {
124-
if (resultFuture.isDone()) {
125-
return;
126-
}
127-
uploadSessionUrl = session.getUploadUrl();
125+
String sessionUrl = session.getUploadUrl();
128126
ResumableUploadChunkCoordinator<ResponseT> coordinator =
129127
new ResumableUploadChunkCoordinator<>(
130-
uploadChunkCallable,
131-
session.getUploadUrl(),
132-
payload,
133-
settings.getChunkSize(),
134-
callContext,
135-
ResumableUploadFutureImpl.this);
136-
try {
137-
coordinator.start();
138-
} catch (Throwable t) {
139-
fail(t);
128+
uploadChunkCallable, sessionUrl, payload, settings.getChunkSize(), callContext);
129+
ApiFuture<ResponseT> uploadFuture = coordinator.getFuture();
130+
synchronized (lock) {
131+
if (resultFuture.isDone()) {
132+
return;
133+
}
134+
uploadSessionUrl = sessionUrl;
135+
inFlightFuture = uploadFuture;
140136
}
137+
ApiFutures.addCallback(
138+
uploadFuture,
139+
new ApiFutureCallback<ResponseT>() {
140+
@Override
141+
public void onSuccess(ResponseT response) {
142+
succeed(response);
143+
}
144+
145+
@Override
146+
public void onFailure(Throwable t) {
147+
if (t instanceof CancellationException) {
148+
return;
149+
}
150+
fail(t);
151+
}
152+
},
153+
MoreExecutors.directExecutor());
154+
coordinator.start();
141155
}
142156

143157
@Override
@@ -151,33 +165,15 @@ public void onFailure(Throwable t) {
151165
MoreExecutors.directExecutor());
152166
}
153167

154-
/**
155-
* Registers the active in-flight future for cancellation. If this session future has already been
156-
* canceled, the supplied future is canceled immediately.
157-
*/
158-
void setInFlightFuture(ApiFuture<?> inFlightFuture) {
159-
boolean shouldCancel = false;
160-
synchronized (lock) {
161-
if (resultFuture.isDone()) {
162-
shouldCancel = resultFuture.isCancelled();
163-
} else {
164-
this.inFlightFuture = inFlightFuture;
165-
}
166-
}
167-
if (shouldCancel) {
168-
inFlightFuture.cancel(true);
169-
}
170-
}
171-
172-
void succeed(@Nullable ResponseT result) {
168+
private void succeed(@Nullable ResponseT result) {
173169
synchronized (lock) {
174170
inFlightFuture = null;
175171
}
176172
closePayload();
177173
resultFuture.set(result);
178174
}
179175

180-
void fail(Throwable t) {
176+
private void fail(Throwable t) {
181177
synchronized (lock) {
182178
inFlightFuture = null;
183179
}

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

Lines changed: 44 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,7 @@
4141
import static org.mockito.Mockito.withSettings;
4242

4343
import com.google.api.core.ApiFutures;
44+
import com.google.api.core.ForwardingApiFuture;
4445
import com.google.api.core.SettableApiFuture;
4546
import com.google.api.gax.resumable.ChunkUploadRequest;
4647
import com.google.api.gax.resumable.ChunkUploadResponse;
@@ -57,6 +58,7 @@
5758
import java.util.concurrent.CountDownLatch;
5859
import java.util.concurrent.ExecutionException;
5960
import java.util.concurrent.TimeUnit;
61+
import java.util.concurrent.atomic.AtomicReference;
6062
import org.junit.jupiter.api.BeforeEach;
6163
import org.junit.jupiter.api.Test;
6264
import org.junit.jupiter.api.extension.ExtendWith;
@@ -230,17 +232,54 @@ void testUploadCallable_cancelInFlight_haltsUpload() throws Exception {
230232
}
231233

232234
@Test
233-
void testUploadCallable_setInFlightFutureAfterCancel_immediatelyCancelsFuture() {
235+
void testUploadCallable_cancelBeforeStartCompletes_abortsChunkUpload() {
234236
SettableApiFuture<ResumableUploadSession> startFuture = SettableApiFuture.create();
235-
when(mockStartCallable.futureCall(any(), any())).thenReturn(startFuture);
237+
// Ignore cancellation on startFuture to simulate start completing concurrently with cancel()
238+
when(mockStartCallable.futureCall(any(), any()))
239+
.thenReturn(
240+
new ForwardingApiFuture<ResumableUploadSession>(startFuture) {
241+
@Override
242+
public boolean cancel(boolean mayInterruptIfRunning) {
243+
return false;
244+
}
245+
});
246+
236247
ResumableUploadFuture<String> future =
237248
callable.futureCall("resource-path", streamOf("data"), null);
238249
assertThat(future.cancel(true)).isTrue();
239250
assertThat(future.isCancelled()).isTrue();
240251

241-
SettableApiFuture<String> lateFuture = SettableApiFuture.create();
242-
((ResumableUploadFutureImpl<String>) future).setInFlightFuture(lateFuture);
243-
assertThat(lateFuture.isCancelled()).isTrue();
252+
startFuture.set(
253+
ResumableUploadSession.newBuilder().setUploadUrl("https://upload.url/late").build());
254+
verifyNoInteractions(mockChunkCallable);
255+
}
256+
257+
@Test
258+
void testUploadCallable_cancelDuringChunkDispatch_immediatelyCancelsChunkFuture() {
259+
stubStartSession("https://upload.url/cancel-dispatch");
260+
SettableApiFuture<ChunkUploadResponse<String>> lateChunkFuture = SettableApiFuture.create();
261+
AtomicReference<ResumableUploadFuture<String>> futureRef = new AtomicReference<>();
262+
when(mockChunkCallable.futureCall(any(ChunkUploadRequest.class), any()))
263+
.thenAnswer(
264+
inv -> {
265+
futureRef.get().cancel(true);
266+
return lateChunkFuture;
267+
});
268+
269+
// Defer startFuture completion until futureRef is populated
270+
SettableApiFuture<ResumableUploadSession> startFuture = SettableApiFuture.create();
271+
when(mockStartCallable.futureCall(any(), any())).thenReturn(startFuture);
272+
273+
ResumableUploadFuture<String> future =
274+
callable.futureCall("resource-path", streamOf("data"), null);
275+
futureRef.set(future);
276+
startFuture.set(
277+
ResumableUploadSession.newBuilder()
278+
.setUploadUrl("https://upload.url/cancel-dispatch")
279+
.build());
280+
281+
assertThat(future.isCancelled()).isTrue();
282+
assertThat(lateChunkFuture.isCancelled()).isTrue();
244283
}
245284

246285
@Test

0 commit comments

Comments
 (0)