feat(gax): add RewindableStreamBuffer in prep for chunk upload recovery - #14423
Conversation
ba8840b to
6d89325
Compare
6d89325 to
da1275c
Compare
589ef8e to
f130ccc
Compare
f130ccc to
0e94df8
Compare
0e94df8 to
00f7932
Compare
00f7932 to
c6c28b1
Compare
c6c28b1 to
4a4b890
Compare
4a4b890 to
14f6204
Compare
14f6204 to
0a389c4
Compare
0a389c4 to
7d192ff
Compare
7d192ff to
33b9862
Compare
2f7c89c to
715822d
Compare
0925229 to
11ab9b0
Compare
3ae0f6e to
97fb372
Compare
97fb372 to
3258c60
Compare
3258c60 to
074e557
Compare
| public abstract byte[] getPayload(); | ||
|
|
||
| /** The number of bytes within {@link #getPayload()} to upload. */ | ||
| public abstract int getPayloadLength(); |
There was a problem hiding this comment.
Do we have to expose getPayloadLength in the request? I think we can do the slicing in RewindableStreamBuffer and getPayload should always return the byte[] that we can upload directly.
| * offset or beyond the current buffer window | ||
| * @throws IOException if reading from the stream fails | ||
| */ | ||
| void realignTo(long committedOffset) throws IOException { |
There was a problem hiding this comment.
Is this method used in this PR?
There was a problem hiding this comment.
No, it's used in recovery (not implemented until #14424). This PR just introduces the buffer and wires it up.
074e557 to
c69d6b1
Compare
c69d6b1 to
7431c49
Compare
7431c49 to
1ba5e07
Compare
95a562e to
4c0fd44
Compare
| * @throws IOException if reading from the stream fails | ||
| */ | ||
| void fill(long targetOffset) throws IOException { | ||
| this.bufferBaseOffset = targetOffset; |
There was a problem hiding this comment.
Looks like this offset is always calculated from the info within the buffer and passed into here in follow up PR. I think we can get rid of this parameter and let RewindableStreamBuffer manage the targetOffset.
Introduce RewindableStreamBuffer managing a single-chunk buffer over an InputStream, supporting forward compaction and topping up upon recovery realignment without mark()/reset(). Enforces boundaries by throwing IllegalStateException when a server offset is below the base offset or beyond the current buffer window.
|
|





Introduces
RewindableStreamBufferto manage a single-chunk buffer for the user-providedInputStream.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.