Skip to content

Commit 95df758

Browse files
authored
fix: update gRPC writeAndClose to only set finish_write on the last message (#2163)
As of 2.26.0 it would set finish_write on every message emitted by writeAndClose Update ITGapicUnbufferedWritableByteChannelTest to also include checksum values in its requests/responses.
1 parent e9746f8 commit 95df758

3 files changed

Lines changed: 48 additions & 50 deletions

File tree

google-cloud-storage/src/main/java/com/google/cloud/storage/GapicUnbufferedWritableByteChannel.java

Lines changed: 21 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,7 @@
3434
import java.nio.channels.ClosedChannelException;
3535
import java.util.ArrayList;
3636
import java.util.List;
37+
import org.checkerframework.checker.nullness.qual.NonNull;
3738

3839
final class GapicUnbufferedWritableByteChannel<
3940
RequestFactoryT extends WriteObjectRequestBuilderFactory>
@@ -82,16 +83,7 @@ public boolean isOpen() {
8283
@Override
8384
public void close() throws IOException {
8485
if (!finished) {
85-
long offset = writeCtx.getTotalSentBytes().get();
86-
Crc32cLengthKnown crc32cValue = writeCtx.getCumulativeCrc32c().get();
87-
88-
WriteObjectRequest.Builder b =
89-
writeCtx.newRequestBuilder().setFinishWrite(true).setWriteOffset(offset);
90-
if (crc32cValue != null) {
91-
b.setObjectChecksums(
92-
ObjectChecksums.newBuilder().setCrc32C(crc32cValue.getValue()).build());
93-
}
94-
WriteObjectRequest message = b.build();
86+
WriteObjectRequest message = finishMessage();
9587
try {
9688
flusher.close(message);
9789
finished = true;
@@ -139,7 +131,7 @@ private long internalWrite(ByteBuffer[] srcs, int srcsOffset, int srcsLength, bo
139131
.newRequestBuilder()
140132
.setWriteOffset(offset)
141133
.setChecksummedData(checksummedData.build());
142-
if (!datum.isOnlyFullBlocks() || finalize) {
134+
if (!datum.isOnlyFullBlocks()) {
143135
builder.setFinishWrite(true);
144136
if (cumulative != null) {
145137
builder.setObjectChecksums(
@@ -152,6 +144,10 @@ private long internalWrite(ByteBuffer[] srcs, int srcsOffset, int srcsLength, bo
152144
messages.add(build);
153145
bytesConsumed += contentSize;
154146
}
147+
if (finalize && !finished) {
148+
messages.add(finishMessage());
149+
finished = true;
150+
}
155151

156152
try {
157153
flusher.flush(messages);
@@ -162,4 +158,18 @@ private long internalWrite(ByteBuffer[] srcs, int srcsOffset, int srcsLength, bo
162158

163159
return bytesConsumed;
164160
}
161+
162+
@NonNull
163+
private WriteObjectRequest finishMessage() {
164+
long offset = writeCtx.getTotalSentBytes().get();
165+
Crc32cLengthKnown crc32cValue = writeCtx.getCumulativeCrc32c().get();
166+
167+
WriteObjectRequest.Builder b =
168+
writeCtx.newRequestBuilder().setFinishWrite(true).setWriteOffset(offset);
169+
if (crc32cValue != null) {
170+
b.setObjectChecksums(ObjectChecksums.newBuilder().setCrc32C(crc32cValue.getValue()).build());
171+
}
172+
WriteObjectRequest message = b.build();
173+
return message;
174+
}
165175
}

google-cloud-storage/src/main/java/com/google/cloud/storage/WriteFlushStrategy.java

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -113,7 +113,11 @@ private static WriteObjectRequest possiblyPairDownRequest(
113113
}
114114

115115
if (message.getWriteOffset() > 0) {
116-
b.clearWriteObjectSpec().clearObjectChecksums();
116+
b.clearWriteObjectSpec();
117+
}
118+
119+
if (message.getWriteOffset() > 0 && !message.getFinishWrite()) {
120+
b.clearObjectChecksums();
117121
}
118122
return b.build();
119123
}

google-cloud-storage/src/test/java/com/google/cloud/storage/ITGapicUnbufferedWritableByteChannelTest.java

Lines changed: 22 additions & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@
3232
import com.google.common.collect.ImmutableSet;
3333
import com.google.protobuf.ByteString;
3434
import com.google.storage.v2.Object;
35+
import com.google.storage.v2.ObjectChecksums;
3536
import com.google.storage.v2.StartResumableWriteRequest;
3637
import com.google.storage.v2.StartResumableWriteResponse;
3738
import com.google.storage.v2.StorageClient;
@@ -63,8 +64,9 @@ public final class ITGapicUnbufferedWritableByteChannelTest {
6364
private static final Logger LOGGER =
6465
Logger.getLogger(ITGapicUnbufferedWritableByteChannelTest.class.getName());
6566

67+
private static final Hasher HASHER = Hasher.enabled();
6668
private static final ChunkSegmenter segmenter =
67-
new ChunkSegmenter(Hasher.noop(), ByteStringStrategy.copy(), 10, 5);
69+
new ChunkSegmenter(HASHER, ByteStringStrategy.copy(), 10, 5);
6870

6971
private static final String uploadId = "upload-id";
7072

@@ -80,31 +82,35 @@ public final class ITGapicUnbufferedWritableByteChannelTest {
8082
private static final WriteObjectRequest req1 =
8183
WriteObjectRequest.newBuilder()
8284
.setUploadId(uploadId)
83-
.setChecksummedData(TestUtils.getChecksummedData(ByteString.copyFrom(bytes, 0, 10)))
85+
.setChecksummedData(getChecksummedData(ByteString.copyFrom(bytes, 0, 10), HASHER))
8486
.build();
8587
private static final WriteObjectRequest req2 =
8688
WriteObjectRequest.newBuilder()
8789
.setUploadId(uploadId)
8890
.setWriteOffset(10)
89-
.setChecksummedData(TestUtils.getChecksummedData(ByteString.copyFrom(bytes, 10, 10)))
91+
.setChecksummedData(getChecksummedData(ByteString.copyFrom(bytes, 10, 10), HASHER))
9092
.build();
9193
private static final WriteObjectRequest req3 =
9294
WriteObjectRequest.newBuilder()
9395
.setUploadId(uploadId)
9496
.setWriteOffset(20)
95-
.setChecksummedData(TestUtils.getChecksummedData(ByteString.copyFrom(bytes, 20, 10)))
97+
.setChecksummedData(getChecksummedData(ByteString.copyFrom(bytes, 20, 10), HASHER))
9698
.build();
9799
private static final WriteObjectRequest req4 =
98100
WriteObjectRequest.newBuilder()
99101
.setUploadId(uploadId)
100102
.setWriteOffset(30)
101-
.setChecksummedData(TestUtils.getChecksummedData(ByteString.copyFrom(bytes, 30, 10)))
103+
.setChecksummedData(getChecksummedData(ByteString.copyFrom(bytes, 30, 10), HASHER))
102104
.build();
103105
private static final WriteObjectRequest req5 =
104106
WriteObjectRequest.newBuilder()
105107
.setUploadId(uploadId)
106108
.setWriteOffset(40)
107109
.setFinishWrite(true)
110+
.setObjectChecksums(
111+
ObjectChecksums.newBuilder()
112+
.setCrc32C(HASHER.hash(ByteBuffer.wrap(bytes)).getValue())
113+
.build())
108114
.build();
109115

110116
private static final WriteObjectResponse resp1 =
@@ -123,35 +129,24 @@ public final class ITGapicUnbufferedWritableByteChannelTest {
123129

124130
@Test
125131
public void directUpload() throws IOException, InterruptedException, ExecutionException {
126-
Object obj = Object.newBuilder().setBucket("buck").setName("obj").build();
127-
WriteObjectSpec spec = WriteObjectSpec.newBuilder().setResource(obj).build();
128132

129133
byte[] bytes = DataGenerator.base64Characters().genBytes(40);
130134
WriteObjectRequest req1 =
131-
WriteObjectRequest.newBuilder()
135+
ITGapicUnbufferedWritableByteChannelTest.req1
136+
.toBuilder()
137+
.clearUploadId()
132138
.setWriteObjectSpec(spec)
133-
.setChecksummedData(TestUtils.getChecksummedData(ByteString.copyFrom(bytes, 0, 10)))
134139
.build();
135140
WriteObjectRequest req2 =
136-
WriteObjectRequest.newBuilder()
137-
.setWriteOffset(10)
138-
.setChecksummedData(TestUtils.getChecksummedData(ByteString.copyFrom(bytes, 10, 10)))
139-
.build();
141+
ITGapicUnbufferedWritableByteChannelTest.req2.toBuilder().clearUploadId().build();
140142
WriteObjectRequest req3 =
141-
WriteObjectRequest.newBuilder()
142-
.setWriteOffset(20)
143-
.setChecksummedData(TestUtils.getChecksummedData(ByteString.copyFrom(bytes, 20, 10)))
144-
.build();
143+
ITGapicUnbufferedWritableByteChannelTest.req3.toBuilder().clearUploadId().build();
145144
WriteObjectRequest req4 =
146-
WriteObjectRequest.newBuilder()
147-
.setWriteOffset(30)
148-
.setChecksummedData(TestUtils.getChecksummedData(ByteString.copyFrom(bytes, 30, 10)))
149-
.build();
145+
ITGapicUnbufferedWritableByteChannelTest.req4.toBuilder().clearUploadId().build();
150146
WriteObjectRequest req5 =
151-
WriteObjectRequest.newBuilder().setWriteOffset(40).setFinishWrite(true).build();
147+
ITGapicUnbufferedWritableByteChannelTest.req5.toBuilder().clearUploadId().build();
152148

153-
WriteObjectResponse resp =
154-
WriteObjectResponse.newBuilder().setResource(obj.toBuilder().setSize(40)).build();
149+
WriteObjectResponse resp = resp5;
155150

156151
WriteObjectRequest base = WriteObjectRequest.newBuilder().setWriteObjectSpec(spec).build();
157152
WriteObjectRequestBuilderFactory reqFactory = WriteObjectRequestBuilderFactory.simple(base);
@@ -314,10 +309,7 @@ public boolean shouldRetry(Throwable t, Object ignore) {
314309

315310
@Test
316311
public void resumableUpload_finalizeWhenWriteAndCloseCalledEvenWhenQuantumAligned()
317-
throws IOException, InterruptedException, ExecutionException {
318-
int quantum = 10;
319-
ChunkSegmenter segmenter =
320-
new ChunkSegmenter(Hasher.noop(), ByteStringStrategy.copy(), 50, quantum);
312+
throws IOException {
321313
SettableApiFuture<WriteObjectResponse> result = SettableApiFuture.create();
322314

323315
AtomicReference<List<WriteObjectRequest>> actualFlush = new AtomicReference<>();
@@ -342,18 +334,10 @@ public void close(@Nullable WriteObjectRequest req) {
342334
}
343335
});
344336

345-
byte[] bytes = DataGenerator.base64Characters().genBytes(quantum);
346-
347337
long written = c.writeAndClose(ByteBuffer.wrap(bytes));
348-
WriteObjectRequest expectedRequest =
349-
WriteObjectRequest.newBuilder()
350-
.setUploadId(uploadId)
351-
.setChecksummedData(getChecksummedData(ByteString.copyFrom(bytes), Hasher.noop()))
352-
.setFinishWrite(true)
353-
.build();
354338

355-
assertThat(written).isEqualTo(10);
356-
assertThat(actualFlush.get()).isEqualTo(ImmutableList.of(expectedRequest));
339+
assertThat(written).isEqualTo(40);
340+
assertThat(actualFlush.get()).isEqualTo(ImmutableList.of(req1, req2, req3, req4, req5));
357341
// calling close is okay, as long as the provided request is null
358342
assertThat(actualClose.get()).isAnyOf(closeRequestSentinel, null);
359343
}

0 commit comments

Comments
 (0)