Skip to content

Commit ee8a737

Browse files
committed
test: add bidiWriteObjectRedirectedError_redirectCounterResetOnResponse to ITAppendableUploadFakeTest
Verify that successful responses on the BidiWriteObject stream (such as a successful chunk persist or reconnect state lookup response) correctly reset the client's consecutive redirect counter to 0. This ensures that the client is not blocked by the max consecutive redirect limit (3) when redirects are spread out. [Generated-by: AI]
1 parent 5dc69d5 commit ee8a737

1 file changed

Lines changed: 205 additions & 0 deletions

File tree

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

Lines changed: 205 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -291,6 +291,211 @@ public void bidiWriteObjectRedirectedError_maxAttempts() throws Exception {
291291
}
292292
}
293293

294+
@Test
295+
public void bidiWriteObjectRedirectedError_redirectCounterResetOnResponse() throws Exception {
296+
String routingToken1 = "routingToken1";
297+
String routingToken2 = "routingToken2";
298+
String routingToken3 = "routingToken3";
299+
String routingToken4 = "routingToken4";
300+
String routingToken5 = "routingToken5";
301+
302+
BidiWriteHandle writeHandle1 =
303+
BidiWriteHandle.newBuilder()
304+
.setHandle(ByteString.copyFromUtf8("handle-1"))
305+
.build();
306+
BidiWriteHandle writeHandle2 =
307+
BidiWriteHandle.newBuilder()
308+
.setHandle(ByteString.copyFromUtf8("handle-2"))
309+
.build();
310+
BidiWriteHandle writeHandle3 =
311+
BidiWriteHandle.newBuilder()
312+
.setHandle(ByteString.copyFromUtf8("handle-3"))
313+
.build();
314+
315+
// 1. First write of "ABC".
316+
BidiWriteObjectRequest req_open_abc = BidiUploadTestUtils.withFlushAndStateLookup(open_abc);
317+
BidiWriteObjectResponse res_abc_h1 = res_abc.toBuilder().setWriteHandle(writeHandle1).build();
318+
319+
// 2. Second write of "DEF".
320+
BidiWriteObjectRequest req_def = BidiUploadTestUtils.withFlushAndStateLookup(def);
321+
BidiWriteObjectResponse res_def_success =
322+
BidiWriteObjectResponse.newBuilder()
323+
.setWriteHandle(writeHandle1)
324+
.setPersistedSize(6)
325+
.build();
326+
327+
// 3. Reconnect to routingToken1
328+
BidiWriteObjectRequest reconnect_token1 =
329+
BidiWriteObjectRequest.newBuilder()
330+
.setAppendObjectSpec(
331+
AppendObjectSpec.newBuilder()
332+
.setBucket(METADATA.getBucket())
333+
.setObject(METADATA.getName())
334+
.setGeneration(METADATA.getGeneration())
335+
.setWriteHandle(writeHandle1)
336+
.setRoutingToken(routingToken1)
337+
.build())
338+
.setStateLookup(true)
339+
.build();
340+
BidiWriteObjectResponse res_lookup_token1 =
341+
BidiWriteObjectResponse.newBuilder()
342+
.setWriteHandle(writeHandle1)
343+
.setPersistedSize(3)
344+
.build();
345+
346+
// 5. Third write of "GHI".
347+
BidiWriteObjectRequest req_ghi = BidiUploadTestUtils.withFlushAndStateLookup(ghi);
348+
BidiWriteObjectResponse res_ghi_success =
349+
BidiWriteObjectResponse.newBuilder()
350+
.setWriteHandle(writeHandle2)
351+
.setPersistedSize(9)
352+
.build();
353+
354+
// 6. Reconnect to routingToken2
355+
BidiWriteObjectRequest reconnect_token2 =
356+
BidiWriteObjectRequest.newBuilder()
357+
.setAppendObjectSpec(
358+
AppendObjectSpec.newBuilder()
359+
.setBucket(METADATA.getBucket())
360+
.setObject(METADATA.getName())
361+
.setGeneration(METADATA.getGeneration())
362+
.setWriteHandle(writeHandle1)
363+
.setRoutingToken(routingToken2)
364+
.build())
365+
.setStateLookup(true)
366+
.build();
367+
BidiWriteObjectResponse res_lookup_token2 =
368+
BidiWriteObjectResponse.newBuilder()
369+
.setWriteHandle(writeHandle2)
370+
.setPersistedSize(6)
371+
.build();
372+
373+
// 8. Fourth write of "J" (finalize/finish write)
374+
BidiWriteObjectRequest req_j_finish = j_finish;
375+
BidiWriteObjectResponse res_j_final = resource_10.toBuilder().setWriteHandle(writeHandle3).build();
376+
377+
// 9. Reconnect to routingToken3
378+
BidiWriteObjectRequest reconnect_token3 =
379+
BidiWriteObjectRequest.newBuilder()
380+
.setAppendObjectSpec(
381+
AppendObjectSpec.newBuilder()
382+
.setBucket(METADATA.getBucket())
383+
.setObject(METADATA.getName())
384+
.setGeneration(METADATA.getGeneration())
385+
.setWriteHandle(writeHandle2)
386+
.setRoutingToken(routingToken3)
387+
.build())
388+
.setStateLookup(true)
389+
.build();
390+
391+
// 10. Reconnect to routingToken4
392+
BidiWriteObjectRequest reconnect_token4 =
393+
BidiWriteObjectRequest.newBuilder()
394+
.setAppendObjectSpec(
395+
AppendObjectSpec.newBuilder()
396+
.setBucket(METADATA.getBucket())
397+
.setObject(METADATA.getName())
398+
.setGeneration(METADATA.getGeneration())
399+
.setWriteHandle(writeHandle2)
400+
.setRoutingToken(routingToken4)
401+
.build())
402+
.setStateLookup(true)
403+
.build();
404+
405+
// 11. Reconnect to routingToken5
406+
BidiWriteObjectRequest reconnect_token5 =
407+
BidiWriteObjectRequest.newBuilder()
408+
.setAppendObjectSpec(
409+
AppendObjectSpec.newBuilder()
410+
.setBucket(METADATA.getBucket())
411+
.setObject(METADATA.getName())
412+
.setGeneration(METADATA.getGeneration())
413+
.setWriteHandle(writeHandle2)
414+
.setRoutingToken(routingToken5)
415+
.build())
416+
.setStateLookup(true)
417+
.build();
418+
BidiWriteObjectResponse res_lookup_token5 =
419+
BidiWriteObjectResponse.newBuilder()
420+
.setWriteHandle(writeHandle3)
421+
.setPersistedSize(9)
422+
.build();
423+
424+
AtomicInteger defCounter = new AtomicInteger();
425+
Consumer<StreamObserver<BidiWriteObjectResponse>> defHandler =
426+
respond -> {
427+
if (defCounter.getAndIncrement() == 0) {
428+
respond.onError(packRedirectIntoAbortedException(makeRedirect(routingToken1)));
429+
} else {
430+
respond.onNext(res_def_success);
431+
}
432+
};
433+
434+
AtomicInteger ghiCounter = new AtomicInteger();
435+
Consumer<StreamObserver<BidiWriteObjectResponse>> ghiHandler =
436+
respond -> {
437+
if (ghiCounter.getAndIncrement() == 0) {
438+
respond.onError(packRedirectIntoAbortedException(makeRedirect(routingToken2)));
439+
} else {
440+
respond.onNext(res_ghi_success);
441+
}
442+
};
443+
444+
AtomicInteger jCounter = new AtomicInteger();
445+
Consumer<StreamObserver<BidiWriteObjectResponse>> jHandler =
446+
respond -> {
447+
if (jCounter.getAndIncrement() == 0) {
448+
respond.onError(packRedirectIntoAbortedException(makeRedirect(routingToken3)));
449+
} else {
450+
respond.onNext(res_j_final);
451+
respond.onCompleted();
452+
}
453+
};
454+
455+
FakeStorage fake =
456+
FakeStorage.of(
457+
ImmutableMap.<BidiWriteObjectRequest, Consumer<StreamObserver<BidiWriteObjectResponse>>>builder()
458+
.put(req_open_abc, respond -> respond.onNext(res_abc_h1))
459+
.put(req_def, defHandler)
460+
.put(reconnect_token1, respond -> respond.onNext(res_lookup_token1))
461+
.put(req_ghi, ghiHandler)
462+
.put(reconnect_token2, respond -> respond.onNext(res_lookup_token2))
463+
.put(req_j_finish, jHandler)
464+
.put(reconnect_token3, respond -> respond.onError(packRedirectIntoAbortedException(makeRedirect(routingToken4))))
465+
.put(reconnect_token4, respond -> respond.onError(packRedirectIntoAbortedException(makeRedirect(routingToken5))))
466+
.put(reconnect_token5, respond -> respond.onNext(res_lookup_token5))
467+
.build());
468+
469+
try (FakeServer fakeServer = FakeServer.of(fake);
470+
Storage storage =
471+
fakeServer.getGrpcStorageOptions().toBuilder()
472+
.setRetrySettings(
473+
fakeServer.getGrpcStorageOptions().getRetrySettings().toBuilder()
474+
.setRetryDelayMultiplier(1.0)
475+
.setInitialRetryDelayDuration(Duration.ofMillis(10))
476+
.build())
477+
.build()
478+
.getService()) {
479+
480+
BlobId id = BlobId.of("b", "o");
481+
BlobAppendableUploadConfig config =
482+
BlobAppendableUploadConfig.of()
483+
.withFlushPolicy(FlushPolicy.maxFlushSize(3))
484+
.withCloseAction(CloseAction.FINALIZE_WHEN_CLOSING);
485+
BlobAppendableUpload b =
486+
storage.blobAppendableUpload(BlobInfo.newBuilder(id).build(), config);
487+
try (AppendableUploadWriteableByteChannel channel = b.open()) {
488+
ByteBuffer wrap = ByteBuffer.wrap(content.getBytes());
489+
Buffers.emptyTo(wrap, channel);
490+
}
491+
492+
// Verification
493+
ApiFuture<BlobInfo> resultFuture = b.getResult();
494+
BlobInfo finalMetadata = resultFuture.get(3, TimeUnit.SECONDS);
495+
assertThat(finalMetadata.getSize()).isEqualTo(10L);
496+
}
497+
}
498+
294499
/**
295500
* We use a small segmenter (3 byte segments) and flush "ABCDEFGHIJ". We make sure that this
296501
* resolves to segments of "ABC"/"DEF"/"GHI"/"J".

0 commit comments

Comments
 (0)