Skip to content

Commit 349c315

Browse files
authored
feat(gax): add RewindableStreamBuffer in prep for chunk upload recovery (#14423)
Introduces `RewindableStreamBuffer` to manage a single-chunk buffer for the user-provided `InputStream`. Recovery is not implemented in this PR - the buffer is wired up in the chunk coordinator here and the query-command-based recovery process will be implemented in the next.
1 parent e72365e commit 349c315

3 files changed

Lines changed: 379 additions & 29 deletions

File tree

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

Lines changed: 11 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -39,10 +39,8 @@
3939
import com.google.api.gax.resumable.ChunkUploadRequest;
4040
import com.google.api.gax.resumable.ChunkUploadResponse;
4141
import com.google.api.gax.resumable.ResumableUploadStatus;
42-
import com.google.common.io.ByteStreams;
4342
import com.google.common.util.concurrent.MoreExecutors;
4443
import java.io.InputStream;
45-
import java.util.Arrays;
4644
import java.util.concurrent.CancellationException;
4745
import org.jspecify.annotations.NullMarked;
4846
import org.jspecify.annotations.Nullable;
@@ -56,14 +54,10 @@
5654
@NullMarked
5755
final class ResumableUploadChunkCoordinator<ResponseT> {
5856

59-
private static final byte[] EMPTY_PAYLOAD = new byte[0];
60-
6157
private final UnaryCallable<ChunkUploadRequest, ChunkUploadResponse<ResponseT>>
6258
uploadChunkCallable;
6359
private final String uploadUrl;
64-
private final InputStream payload;
65-
private final byte[] buffer;
66-
private final int chunkSize;
60+
private final RewindableStreamBuffer buffer;
6761
private final ApiCallContext callContext;
6862
private final SettableApiFuture<ResponseT> uploadResultFuture = SettableApiFuture.create();
6963
private volatile @Nullable ApiFuture<?> inFlightFuture;
@@ -77,10 +71,9 @@ final class ResumableUploadChunkCoordinator<ResponseT> {
7771
this.uploadChunkCallable =
7872
checkNotNull(uploadChunkCallable, "uploadChunkCallable must not be null");
7973
this.uploadUrl = checkNotNull(uploadUrl, "uploadUrl must not be null");
80-
this.payload = checkNotNull(payload, "payload must not be null");
81-
this.chunkSize = chunkSize;
74+
checkNotNull(payload, "payload must not be null");
8275
this.callContext = checkNotNull(callContext, "callContext must not be null");
83-
this.buffer = new byte[chunkSize];
76+
this.buffer = new RewindableStreamBuffer(payload, chunkSize, uploadUrl);
8477
}
8578

8679
ApiFuture<ResponseT> getFuture() {
@@ -96,40 +89,30 @@ void start() {
9689
}
9790
},
9891
MoreExecutors.directExecutor());
99-
transmitChunk(0L);
92+
transmitChunk();
10093
}
10194

102-
private void transmitChunk(long currentOffset) {
95+
private void transmitChunk() {
10396
try {
10497
// Abort if the session was already completed or canceled.
10598
if (uploadResultFuture.isDone()) {
10699
return;
107100
}
108101

109102
// Read the next chunk slice from the payload stream.
110-
int bytesRead = ByteStreams.read(payload, buffer, 0, chunkSize);
103+
buffer.fill();
111104

112105
// 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-
}
122-
123106
ChunkUploadRequest chunkRequest =
124107
ChunkUploadRequest.newBuilder()
125108
.setUploadUrl(uploadUrl)
126-
.setPayload(chunkPayload)
127-
.setOffset(currentOffset)
128-
.setFinal(isFinal)
109+
.setPayload(buffer.getPayload())
110+
.setOffset(buffer.getBufferBaseOffset())
111+
.setFinal(buffer.isFinal())
129112
.build();
130113

131114
// Dispatch the chunk upload call and register the in-flight future for cancellation.
132-
long chunkLength = chunkPayload.length;
115+
boolean isFinal = chunkRequest.isFinal();
133116
ApiFuture<ChunkUploadResponse<ResponseT>> chunkFuture =
134117
uploadChunkCallable.futureCall(chunkRequest, callContext);
135118
if (!tryRegisterInFlightFuture(chunkFuture)) {
@@ -144,7 +127,6 @@ public void onSuccess(ChunkUploadResponse<ResponseT> response) {
144127
if (uploadResultFuture.isDone()) {
145128
return;
146129
}
147-
long nextOffset = currentOffset + chunkLength;
148130
if (response.getUploadStatus() == ResumableUploadStatus.FINAL) {
149131
uploadResultFuture.set(response.getResponse());
150132
} else if (isFinal) {
@@ -153,7 +135,7 @@ public void onSuccess(ChunkUploadResponse<ResponseT> response) {
153135
"Upload stream ended and final chunk was transmitted, but server returned"
154136
+ " incomplete status"));
155137
} else {
156-
transmitChunk(nextOffset);
138+
transmitChunk();
157139
}
158140
}
159141

Lines changed: 150 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,150 @@
1+
/*
2+
* Copyright 2026 Google LLC
3+
*
4+
* Redistribution and use in source and binary forms, with or without
5+
* modification, are permitted provided that the following conditions are
6+
* met:
7+
*
8+
* * Redistributions of source code must retain the above copyright
9+
* notice, this list of conditions and the following disclaimer.
10+
* * Redistributions in binary form must reproduce the above
11+
* copyright notice, this list of conditions and the following disclaimer
12+
* in the documentation and/or other materials provided with the
13+
* distribution.
14+
* * Neither the name of Google LLC nor the names of its
15+
* contributors may be used to endorse or promote products derived from
16+
* this software without specific prior written permission.
17+
*
18+
* THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS
19+
* "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT
20+
* LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR
21+
* A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT
22+
* OWNER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL,
23+
* SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT
24+
* LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE,
25+
* DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY
26+
* THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT
27+
* (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE
28+
* OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
29+
*/
30+
package com.google.api.gax.rpc;
31+
32+
import static com.google.common.base.Preconditions.checkArgument;
33+
import static com.google.common.base.Preconditions.checkNotNull;
34+
35+
import com.google.common.io.ByteStreams;
36+
import java.io.IOException;
37+
import java.io.InputStream;
38+
import java.util.Arrays;
39+
import org.jspecify.annotations.NullMarked;
40+
41+
/** Manages a single-chunk rewindable buffer over an {@link InputStream} for resumable uploads. */
42+
@NullMarked
43+
final class RewindableStreamBuffer {
44+
45+
private static final byte[] EMPTY_PAYLOAD = new byte[0];
46+
47+
private final InputStream inputStream;
48+
private final int chunkSize;
49+
private final String uploadUrl;
50+
private final byte[] buffer;
51+
52+
private long bufferBaseOffset;
53+
private int payloadLength;
54+
private boolean isFinal;
55+
private boolean streamExhausted;
56+
57+
RewindableStreamBuffer(InputStream inputStream, int chunkSize, String uploadUrl) {
58+
this.inputStream = checkNotNull(inputStream, "inputStream must not be null");
59+
checkArgument(chunkSize > 0, "chunkSize must be > 0");
60+
this.chunkSize = chunkSize;
61+
this.uploadUrl = checkNotNull(uploadUrl, "uploadUrl must not be null");
62+
this.buffer = new byte[chunkSize];
63+
this.bufferBaseOffset = 0L;
64+
this.payloadLength = 0;
65+
this.isFinal = false;
66+
this.streamExhausted = false;
67+
}
68+
69+
/**
70+
* Advances the buffer past the current payload, reading up to chunk size from the stream.
71+
*
72+
* @throws IOException if reading from the stream fails
73+
*/
74+
void fill() throws IOException {
75+
this.bufferBaseOffset += payloadLength;
76+
this.payloadLength = ByteStreams.read(inputStream, buffer, 0, chunkSize);
77+
this.isFinal = (payloadLength < chunkSize);
78+
if (this.isFinal) {
79+
this.streamExhausted = true;
80+
}
81+
}
82+
83+
/**
84+
* Realigns the buffer window to {@code committedOffset}.
85+
*
86+
* <p>Compacts forward within the existing buffer to discard already-committed bytes, and then
87+
* tops up the buffer to capacity from the underlying stream.
88+
*
89+
* @param committedOffset the server's committed byte offset
90+
* @throws IllegalStateException if {@code committedOffset} is below the buffer's base offset or
91+
* beyond the current buffer window
92+
* @throws IOException if reading from the stream fails
93+
*/
94+
void realignTo(long committedOffset) throws IOException {
95+
if (committedOffset < bufferBaseOffset) {
96+
throw new IllegalStateException(
97+
String.format(
98+
"Server committed offset %d is below buffer base offset %d for upload URL %s; cannot"
99+
+ " rewind stream before buffer base",
100+
committedOffset, bufferBaseOffset, uploadUrl));
101+
}
102+
103+
if (committedOffset > bufferBaseOffset + payloadLength) {
104+
throw new IllegalStateException(
105+
String.format(
106+
"Server committed offset %d is beyond current buffer window [%d, %d] for upload URL"
107+
+ " %s",
108+
committedOffset, bufferBaseOffset, bufferBaseOffset + payloadLength, uploadUrl));
109+
}
110+
111+
int committedWithinBuffer = (int) (committedOffset - bufferBaseOffset);
112+
int remainingBytes = payloadLength - committedWithinBuffer;
113+
114+
if (remainingBytes > 0 && committedWithinBuffer > 0) {
115+
System.arraycopy(buffer, committedWithinBuffer, buffer, 0, remainingBytes);
116+
}
117+
118+
this.bufferBaseOffset = committedOffset;
119+
this.payloadLength = remainingBytes;
120+
121+
if (!streamExhausted && payloadLength < chunkSize) {
122+
int space = chunkSize - payloadLength;
123+
int additionalRead = ByteStreams.read(inputStream, buffer, payloadLength, space);
124+
payloadLength += additionalRead;
125+
if (additionalRead < space) {
126+
streamExhausted = true;
127+
}
128+
}
129+
130+
this.isFinal = streamExhausted;
131+
}
132+
133+
byte[] getPayload() {
134+
if (payloadLength == buffer.length) {
135+
return buffer;
136+
}
137+
if (payloadLength == 0) {
138+
return EMPTY_PAYLOAD;
139+
}
140+
return Arrays.copyOf(buffer, payloadLength);
141+
}
142+
143+
long getBufferBaseOffset() {
144+
return bufferBaseOffset;
145+
}
146+
147+
boolean isFinal() {
148+
return isFinal;
149+
}
150+
}

0 commit comments

Comments
 (0)