This is an automated email from the ASF dual-hosted git repository.

btellier pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/james-project.git


The following commit(s) were added to refs/heads/master by this push:
     new b8d8128  JAMES-3150 S3BlobStoreDAO listBlob paging (#643)
b8d8128 is described below

commit b8d81280600d5bd310e1c532d747b5fa640e5b75
Author: Benoit TELLIER <[email protected]>
AuthorDate: Thu Sep 9 10:21:52 2021 +0700

    JAMES-3150 S3BlobStoreDAO listBlob paging (#643)
---
 .../blob/objectstorage/aws/S3BlobStoreDAO.java     | 13 +++++--------
 .../blob/objectstorage/aws/S3BlobStoreDAOTest.java | 22 ++++++++++++++++++++++
 2 files changed, 27 insertions(+), 8 deletions(-)

diff --git 
a/server/blob/blob-s3/src/main/java/org/apache/james/blob/objectstorage/aws/S3BlobStoreDAO.java
 
b/server/blob/blob-s3/src/main/java/org/apache/james/blob/objectstorage/aws/S3BlobStoreDAO.java
index cac27db..6baae74 100644
--- 
a/server/blob/blob-s3/src/main/java/org/apache/james/blob/objectstorage/aws/S3BlobStoreDAO.java
+++ 
b/server/blob/blob-s3/src/main/java/org/apache/james/blob/objectstorage/aws/S3BlobStoreDAO.java
@@ -69,7 +69,7 @@ import 
software.amazon.awssdk.services.s3.model.BucketAlreadyOwnedByYouException
 import software.amazon.awssdk.services.s3.model.DeleteObjectsResponse;
 import software.amazon.awssdk.services.s3.model.GetObjectResponse;
 import software.amazon.awssdk.services.s3.model.ListBucketsResponse;
-import software.amazon.awssdk.services.s3.model.ListObjectsResponse;
+import software.amazon.awssdk.services.s3.model.ListObjectsV2Response;
 import software.amazon.awssdk.services.s3.model.NoSuchBucketException;
 import software.amazon.awssdk.services.s3.model.NoSuchKeyException;
 import software.amazon.awssdk.services.s3.model.ObjectIdentifier;
@@ -282,14 +282,12 @@ public class S3BlobStoreDAO implements BlobStoreDAO, 
Startable, Closeable {
     }
 
     private Mono<BucketName> emptyBucket(BucketName bucketName) {
-        return Mono.fromFuture(() -> client.listObjects(builder -> 
builder.bucket(bucketName.asString())))
+        return Flux.from(client.listObjectsV2Paginator(builder -> 
builder.bucket(bucketName.asString())))
             .flatMap(response -> Flux.fromIterable(response.contents())
                 .window(EMPTY_BUCKET_BATCH_SIZE)
                 .flatMap(this::buildListForBatch, DEFAULT_CONCURRENCY)
                 .flatMap(identifiers -> deleteObjects(bucketName, 
identifiers), DEFAULT_CONCURRENCY)
                 .then(Mono.just(response)))
-            .flux()
-            .takeUntil(list -> !list.isTruncated())
             .then(Mono.just(bucketName));
     }
 
@@ -324,12 +322,11 @@ public class S3BlobStoreDAO implements BlobStoreDAO, 
Startable, Closeable {
 
     @Override
     public Publisher<BlobId> listBlobs(BucketName bucketName) {
-        return Mono.fromFuture(() -> client.listObjects(builder -> 
builder.bucket(bucketName.asString())))
-            .flux()
-            .takeUntil(list -> !list.isTruncated())
-            .flatMapIterable(ListObjectsResponse::contents)
+        return Flux.from(client.listObjectsV2Paginator(builder -> 
builder.bucket(bucketName.asString())))
+            .flatMapIterable(ListObjectsV2Response::contents)
             .map(S3Object::key)
             .map(blobIdFactory::from)
+            .onErrorResume(e -> e.getCause() instanceof NoSuchBucketException, 
e -> Flux.empty())
             .onErrorResume(NoSuchBucketException.class, e -> Flux.empty());
     }
 }
diff --git 
a/server/blob/blob-s3/src/test/java/org/apache/james/blob/objectstorage/aws/S3BlobStoreDAOTest.java
 
b/server/blob/blob-s3/src/test/java/org/apache/james/blob/objectstorage/aws/S3BlobStoreDAOTest.java
index 01922a5..2786acd 100644
--- 
a/server/blob/blob-s3/src/test/java/org/apache/james/blob/objectstorage/aws/S3BlobStoreDAOTest.java
+++ 
b/server/blob/blob-s3/src/test/java/org/apache/james/blob/objectstorage/aws/S3BlobStoreDAOTest.java
@@ -18,14 +18,23 @@
  ****************************************************************/
 package org.apache.james.blob.objectstorage.aws;
 
+import static org.apache.james.blob.api.BlobStoreDAOFixture.ELEVEN_KILOBYTES;
+import static org.apache.james.blob.api.BlobStoreDAOFixture.TEST_BUCKET_NAME;
+import static org.assertj.core.api.Assertions.assertThat;
+
 import org.apache.james.blob.api.BlobStoreDAO;
 import org.apache.james.blob.api.BlobStoreDAOContract;
 import org.apache.james.blob.api.TestBlobId;
 import org.junit.jupiter.api.AfterAll;
 import org.junit.jupiter.api.AfterEach;
 import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.extension.ExtendWith;
 
+import com.google.common.io.ByteSource;
+
+import reactor.core.publisher.Flux;
+
 @ExtendWith(DockerAwsS3Extension.class)
 public class S3BlobStoreDAOTest implements BlobStoreDAOContract {
     private static S3BlobStoreDAO testee;
@@ -60,4 +69,17 @@ public class S3BlobStoreDAOTest implements 
BlobStoreDAOContract {
     public BlobStoreDAO testee() {
         return testee;
     }
+
+    @Test
+    void listingManyBlobsShouldSucceedWhenExceedingPageSize() {
+        BlobStoreDAO store = testee();
+
+        final int count = 1500;
+        Flux.range(0, count)
+            .flatMap(i -> store.save(TEST_BUCKET_NAME, new 
TestBlobId("test-blob-id-" + i), ByteSource.wrap(ELEVEN_KILOBYTES)))
+            .blockLast();
+
+        
assertThat(Flux.from(testee().listBlobs(TEST_BUCKET_NAME)).count().block())
+            .isEqualTo(count);
+    }
 }

---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to