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

chibenwa 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 3e54de48f4 [ENHANCEMENT] Apache James Blob File Sharding (#3197)
3e54de48f4 is described below

commit 3e54de48f4181ca0a80f5564bd7cdd3b60d67264
Author: ilya terskov <[email protected]>
AuthorDate: Thu Oct 1 16:27:11 2026 +0700

    [ENHANCEMENT] Apache James Blob File Sharding (#3197)
    
    Reuses semantic of MinIO blobstore and enables it by default onto the file 
blob store
---
 .../servers/partials/configure/blobstore.adoc      |  17 +-
 .../sample-configuration/jvm.properties            |   5 +-
 .../apache/james/blob/file/FileBlobStoreDAO.java   |  68 ++++----
 .../james/blob/file/FileBlobStoreDAOTest.java      |  55 ++++++
 .../blob/file/FileBlobStoreGCAlgorithmTest.java    |  10 ++
 .../blob/file/FileBlobStorePassThroughTest.java    |  14 ++
 .../blob/file/FileWithFolderHierarchyTest.java     | 184 +++++++++++++++++++++
 .../blobstore/BlobDeduplicationGCModule.java       |  16 +-
 .../blobstore/BlobDeduplicationGCModuleTest.java   |  99 +++++++++++
 9 files changed, 426 insertions(+), 42 deletions(-)

diff --git a/docs/modules/servers/partials/configure/blobstore.adoc 
b/docs/modules/servers/partials/configure/blobstore.adoc
index 1ef24ae231..4ed0caa956 100644
--- a/docs/modules/servers/partials/configure/blobstore.adoc
+++ b/docs/modules/servers/partials/configure/blobstore.adoc
@@ -204,19 +204,26 @@ This bucket name is used as is: 
`objectstorage.bucketPrefix` is not applied to i
 
 |===
 
-==== Improve listing support for MinIO
+[[_improve_listing_support_for_minio]]
+==== Improve listing and filesystem performance with blob folder hierarchy
 
-Due to blobs being stored in folder, adding `/` in blobs name emulates folder 
and avoids blobs to be all stored in a
-same folder, thus improving listing.
+Due to blobs being stored in a single flat namespace by default, adding `/` in 
blob names creates subfolders
+and avoids millions of blobs from being stored in a single folder. This 
significantly improves listing on MinIO
+and prevents directory entry bloat and file system performance degradation in 
`FileBlobStore`.
 
 Instead of `1_628_36825033-d835-4490-9f5a-eef120b1e85c` the following blob id 
will be used: `1/628/3/6/8/2/5033-d835-4490-9f5a-eef120b1e85c`
 
-To enable blob hierarchy compatible with MinIO add in `jvm.properties`:
+To enable blob folder hierarchy (on by default in Postgres App distribution), 
add in `jvm.properties`:
 
 ----
-james.s3.minio.compatibility.mode=true
+james.blobstore.folder.hierarchy=true
 ----
 
+NOTE: The legacy property `james.s3.minio.compatibility.mode=true` is also 
supported as a backward-compatible alias.
+
+NOTE: The `FileBlobStore` implementation is officially not supported on 
Windows operating systems due to Windows mandatory file locking semantics upon 
concurrent file replacement/deletion. UNIX/Linux POSIX-compliant environments 
are recommended for production deployments.
+
+
 ==== Unordered listing for Ceph RADOS Gateway
 
 Ceph RADOS Gateway supports an `allow-unordered` extension on bucket listings: 
instead of merging the entries of
diff --git a/server/apps/postgres-app/sample-configuration/jvm.properties 
b/server/apps/postgres-app/sample-configuration/jvm.properties
index b2fe9d6abc..b9c65fb444 100644
--- a/server/apps/postgres-app/sample-configuration/jvm.properties
+++ b/server/apps/postgres-app/sample-configuration/jvm.properties
@@ -52,4 +52,7 @@ james.jmx.credential.generation=true
 jmx.remote.x.mlet.allow.getMBeansFromURL=false
 
 # Integer. Optional, defaults to 5000. In case of large data, this argument 
specifies the maximum number of rows to return in a single batch set when 
executing query.
-#query.batch.size=5000
\ No newline at end of file
+#query.batch.size=5000
+
+# Enable folder hierarchy across subdirectories for blob storage (S3/MinIO and 
FileBlobStore) to avoid directory bloat
+james.blobstore.folder.hierarchy=true
diff --git 
a/server/blob/blob-file/src/main/java/org/apache/james/blob/file/FileBlobStoreDAO.java
 
b/server/blob/blob-file/src/main/java/org/apache/james/blob/file/FileBlobStoreDAO.java
index c1db4d532a..b80c1f8d3b 100644
--- 
a/server/blob/blob-file/src/main/java/org/apache/james/blob/file/FileBlobStoreDAO.java
+++ 
b/server/blob/blob-file/src/main/java/org/apache/james/blob/file/FileBlobStoreDAO.java
@@ -66,11 +66,13 @@ import reactor.core.scheduler.Schedulers;
 import reactor.util.retry.Retry;
 
 public class FileBlobStoreDAO implements BlobStoreDAO {
+
     private static final Logger LOGGER = 
LoggerFactory.getLogger(FileBlobStoreDAO.class);
     private static final String JAMES_BLOB_METADATA_ATTRIBUTE_PREFIX = 
"james-blob-metadata-";
+    private static final String STAGING_PREFIX = ".james-staging-";
 
     private final File root;
-    private final  BlobId.Factory blobIdFactory;
+    private final BlobId.Factory blobIdFactory;
 
     @Inject
     public FileBlobStoreDAO(FileSystem fileSystem, BlobId.Factory 
blobIdFactory) throws FileNotFoundException {
@@ -80,8 +82,7 @@ public class FileBlobStoreDAO implements BlobStoreDAO {
 
     @Override
     public InputStreamBlob read(BucketName bucketName, BlobId blobId) throws 
ObjectStoreIOException, ObjectNotFoundException {
-        File bucketRoot = getBucketRoot(bucketName);
-        File blob = new File(bucketRoot, blobId.asString());
+        File blob = getBlobFile(bucketName, blobId);
         try {
             return InputStreamBlob.of(new FileInputStream(blob), 
readMetadata(blob.toPath()));
         } catch (FileNotFoundException e) {
@@ -89,6 +90,13 @@ public class FileBlobStoreDAO implements BlobStoreDAO {
         }
     }
 
+    private File getBlobFile(BucketName bucketName, BlobId blobId) {
+        File bucketRoot = getBucketRoot(bucketName);
+        File blob = new File(bucketRoot, blobId.asString());
+        Preconditions.checkArgument(!isStagingFile(blob.toPath()), "Blob name 
uses reserved staging prefix: %s", blobId.asString());
+        return blob;
+    }
+
     private File getBucketRoot(BucketName bucketName) {
         File bucketRoot = new File(root, bucketName.asString());
         if (!bucketRoot.exists()) {
@@ -110,8 +118,7 @@ public class FileBlobStoreDAO implements BlobStoreDAO {
     @Override
     public Publisher<BytesBlob> readBytes(BucketName bucketName, BlobId 
blobId) {
         return Mono.fromCallable(() -> {
-                File bucketRoot = getBucketRoot(bucketName);
-                File blob = new File(bucketRoot, blobId.asString());
+                File blob = getBlobFile(bucketName, blobId);
                 return BytesBlob.of(FileUtils.readFileToByteArray(blob), 
readMetadata(blob.toPath()));
             }).onErrorResume(NoSuchFileException.class, e -> Mono.error(new 
ObjectNotFoundException(String.format("Cannot locate %s within %s", 
blobId.asString(), bucketName.asString()), e)))
             .subscribeOn(Schedulers.boundedElastic());
@@ -129,22 +136,14 @@ public class FileBlobStoreDAO implements BlobStoreDAO {
     public Mono<Void> save(BucketName bucketName, BlobId blobId, byte[] data, 
BlobMetadata metadata) {
         Preconditions.checkNotNull(data);
 
-        return Mono.fromRunnable(() -> {
-                File bucketRoot = getBucketRoot(bucketName);
-                File blob = new File(bucketRoot, blobId.asString());
-                save(data, blob, metadata);
-            })
+        return Mono.fromRunnable(() -> save(data, getBlobFile(bucketName, 
blobId), metadata))
             .subscribeOn(Schedulers.boundedElastic())
             .then();
     }
 
     public Mono<Void> save(BucketName bucketName, BlobId blobId, InputStream 
inputStream, BlobMetadata metadata) {
         Preconditions.checkNotNull(inputStream);
-        return Mono.fromRunnable(() -> {
-                File bucketRoot = getBucketRoot(bucketName);
-                File blob = new File(bucketRoot, blobId.asString());
-                save(inputStream, blob, metadata);
-            })
+        return Mono.fromRunnable(() -> save(inputStream, 
getBlobFile(bucketName, blobId), metadata))
             .subscribeOn(Schedulers.boundedElastic())
             .then()
             .retryWhen(Retry.backoff(10, Duration.ofMillis(100))
@@ -182,15 +181,18 @@ public class FileBlobStoreDAO implements BlobStoreDAO {
         }
     }
 
+
     public Mono<Void> save(BucketName bucketName, BlobId blobId, ByteSource 
content, BlobMetadata metadata) {
-        return Mono.fromCallable(() -> {
+        return Mono.using(
+            () -> {
                 try {
-                    return content.read();
+                    return content.openStream();
                 } catch (IOException e) {
                     throw new ObjectStoreIOException("IOException occurred", 
e);
                 }
-            })
-            .flatMap(bytes -> save(bucketName, blobId, bytes, metadata));
+            },
+            is -> save(bucketName, blobId, is, metadata),
+            Throwing.consumer(InputStream::close).sneakyThrow());
     }
 
     @Override
@@ -198,9 +200,12 @@ public class FileBlobStoreDAO implements BlobStoreDAO {
         Preconditions.checkNotNull(bucketName);
 
         return Mono.fromRunnable(Throwing.runnable(() -> {
-                File bucketRoot = getBucketRoot(bucketName);
-                File blob = new File(bucketRoot, blobId.asString());
-                FileUtils.deleteQuietly(blob);
+                File blob = getBlobFile(bucketName, blobId);
+                try {
+                    Files.deleteIfExists(blob.toPath());
+                } catch (IOException e) {
+                    throw new ObjectStoreIOException("Error deleting blob", e);
+                }
             }))
             .subscribeOn(Schedulers.boundedElastic())
             .then();
@@ -235,16 +240,18 @@ public class FileBlobStoreDAO implements BlobStoreDAO {
     @Override
     public Publisher<BlobId> listBlobs(BucketName bucketName) {
         return Mono.fromCallable(() -> {
-                File bucketRoot = getBucketRoot(bucketName);
+                File bucketRoot = new File(root, bucketName.asString());
                 Path rootPath = bucketRoot.toPath();
-                // Blob ids may contain '/' (eg. recovery sidecar keys) and 
are then stored in nested
+                // Blob ids may contain '/' (eg. recovery sidecar keys or 
hierarchy-aware blob IDs) and are then stored in nested
                 // directories, so we walk the tree and rebuild the id from 
the bucket-root-relative path.
                 return Files.walk(rootPath)
                     .filter(Files::isRegularFile)
+                    .filter(path -> !isStagingFile(path))
                     .map(path -> 
blobIdFactory.parse(toBlobId(rootPath.relativize(path))));
             })
             .flatMapMany(Flux::fromStream)
-            .subscribeOn(Schedulers.boundedElastic());
+            .subscribeOn(Schedulers.boundedElastic())
+            .onErrorResume(NoSuchFileException.class, e -> Flux.empty());
     }
 
     private String toBlobId(Path relativePath) {
@@ -311,12 +318,15 @@ public class FileBlobStoreDAO implements BlobStoreDAO {
         return StandardCharsets.UTF_8.decode(byteBuffer).toString();
     }
 
+    private boolean isStagingFile(Path path) {
+        return path.getFileName().toString().startsWith(STAGING_PREFIX);
+    }
+
     private File createTempFile(File blob) {
+        Path parentPath = blob.getParentFile().toPath();
         try {
-            // Blob ids may contain '/' (eg. recovery sidecar keys), 
introducing nested directories
-            // that must exist before the temp file is created alongside the 
target blob.
-            Files.createDirectories(blob.getParentFile().toPath());
-            return Files.createTempFile(blob.getParentFile().toPath(), 
blob.getName(), ".tmp").toFile();
+            Files.createDirectories(parentPath);
+            return Files.createTempFile(parentPath, STAGING_PREFIX, 
"").toFile();
         } catch (IOException e) {
             throw new ObjectStoreIOException("IOException occurred", e);
         }
diff --git 
a/server/blob/blob-file/src/test/java/org/apache/james/blob/file/FileBlobStoreDAOTest.java
 
b/server/blob/blob-file/src/test/java/org/apache/james/blob/file/FileBlobStoreDAOTest.java
index 9fd98c0747..8de2169cd3 100644
--- 
a/server/blob/blob-file/src/test/java/org/apache/james/blob/file/FileBlobStoreDAOTest.java
+++ 
b/server/blob/blob-file/src/test/java/org/apache/james/blob/file/FileBlobStoreDAOTest.java
@@ -26,6 +26,11 @@ import org.apache.james.blob.api.PlainBlobId;
 import org.apache.james.server.core.filesystem.FileSystemImpl;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Disabled;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.condition.DisabledOnOs;
+import org.junit.jupiter.api.condition.OS;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.MethodSource;
 
 class FileBlobStoreDAOTest implements BlobStoreDAOContract, 
MetadataAwareBlobStoreDAOContract {
 
@@ -42,8 +47,58 @@ class FileBlobStoreDAOTest implements BlobStoreDAOContract, 
MetadataAwareBlobSto
     }
 
     @Override
+    @Test
     @Disabled("Not supported")
     public void mixingSaveReadAndDeleteShouldReturnConsistentState() {
 
     }
+
+    @Override
+    @DisabledOnOs(OS.WINDOWS)
+    @ParameterizedTest(name = "[{index}] {0}")
+    @MethodSource("blobs")
+    public void concurrentSaveBytesShouldReturnConsistentValues(String 
description, BlobStoreDAO.BytesBlob bytes) throws 
java.util.concurrent.ExecutionException, InterruptedException {
+        
BlobStoreDAOContract.super.concurrentSaveBytesShouldReturnConsistentValues(description,
 bytes);
+    }
+
+    @Override
+    @DisabledOnOs(OS.WINDOWS)
+    @ParameterizedTest(name = "[{index}] {0}")
+    @MethodSource("blobs")
+    public void concurrentSaveInputStreamShouldReturnConsistentValues(String 
description, BlobStoreDAO.BytesBlob bytes) throws 
java.util.concurrent.ExecutionException, InterruptedException {
+        
BlobStoreDAOContract.super.concurrentSaveInputStreamShouldReturnConsistentValues(description,
 bytes);
+    }
+
+    @Override
+    @DisabledOnOs(OS.WINDOWS)
+    @ParameterizedTest(name = "[{index}] {0}")
+    @MethodSource("blobs")
+    public void concurrentSaveByteSourceShouldReturnConsistentValues(String 
description, BlobStoreDAO.BytesBlob bytes) throws 
java.util.concurrent.ExecutionException, InterruptedException {
+        
BlobStoreDAOContract.super.concurrentSaveByteSourceShouldReturnConsistentValues(description,
 bytes);
+    }
+
+    @Override
+    @Test
+    @DisabledOnOs(OS.WINDOWS)
+    public void 
readBytesShouldNotReadPartiallyWhenDeletingConcurrentlyBigBlob() throws 
Exception {
+        
BlobStoreDAOContract.super.readBytesShouldNotReadPartiallyWhenDeletingConcurrentlyBigBlob();
+    }
+
+    @Override
+    @Test
+    @DisabledOnOs(OS.WINDOWS)
+    public void readShouldNotReadPartiallyWhenDeletingConcurrentlyBigBlob() 
throws Exception {
+        
BlobStoreDAOContract.super.readShouldNotReadPartiallyWhenDeletingConcurrentlyBigBlob();
+    }
+
+    @Test
+    void saveShouldRejectBlobIdWithReservedStagingPrefix() {
+        org.apache.james.blob.api.BucketName bucketName = 
org.apache.james.blob.api.BucketName.of("test-bucket");
+        org.apache.james.blob.api.BlobId nestedStagingId = new 
PlainBlobId.Factory().of("folder/.james-staging-evil");
+
+        org.assertj.core.api.Assertions.assertThatThrownBy(() ->
+                reactor.core.publisher.Mono.from(blobStore.save(bucketName, 
nestedStagingId, BlobStoreDAO.BytesBlob.of(new byte[]{1, 2, 3}))).block())
+            .isInstanceOf(IllegalArgumentException.class)
+            .hasMessageContaining("Blob name uses reserved staging prefix");
+    }
 }
diff --git 
a/server/blob/blob-file/src/test/java/org/apache/james/blob/file/FileBlobStoreGCAlgorithmTest.java
 
b/server/blob/blob-file/src/test/java/org/apache/james/blob/file/FileBlobStoreGCAlgorithmTest.java
index fe72ec9444..a03b9d8b75 100644
--- 
a/server/blob/blob-file/src/test/java/org/apache/james/blob/file/FileBlobStoreGCAlgorithmTest.java
+++ 
b/server/blob/blob-file/src/test/java/org/apache/james/blob/file/FileBlobStoreGCAlgorithmTest.java
@@ -24,6 +24,9 @@ import org.apache.james.blob.api.PlainBlobId;
 import 
org.apache.james.server.blob.deduplication.BloomFilterGCAlgorithmContract;
 import org.apache.james.server.core.filesystem.FileSystemImpl;
 import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.condition.DisabledOnOs;
+import org.junit.jupiter.api.condition.OS;
 
 public class FileBlobStoreGCAlgorithmTest implements 
BloomFilterGCAlgorithmContract {
 
@@ -38,4 +41,11 @@ public class FileBlobStoreGCAlgorithmTest implements 
BloomFilterGCAlgorithmContr
     public BlobStoreDAO blobStoreDAO() {
         return blobStoreDAO;
     }
+
+    @Override
+    @Test
+    @DisabledOnOs(OS.WINDOWS)
+    public void allOrphanBlobIdsShouldRemovedAfterMultipleRunningTimesGC() {
+        
BloomFilterGCAlgorithmContract.super.allOrphanBlobIdsShouldRemovedAfterMultipleRunningTimesGC();
+    }
 }
diff --git 
a/server/blob/blob-file/src/test/java/org/apache/james/blob/file/FileBlobStorePassThroughTest.java
 
b/server/blob/blob-file/src/test/java/org/apache/james/blob/file/FileBlobStorePassThroughTest.java
index 0538af2990..a252b20a98 100644
--- 
a/server/blob/blob-file/src/test/java/org/apache/james/blob/file/FileBlobStorePassThroughTest.java
+++ 
b/server/blob/blob-file/src/test/java/org/apache/james/blob/file/FileBlobStorePassThroughTest.java
@@ -51,4 +51,18 @@ public class FileBlobStorePassThroughTest implements 
DeleteBlobStoreContract, Me
     public BlobId.Factory blobIdFactory() {
         return BLOB_ID_FACTORY;
     }
+
+    @Override
+    @org.junit.jupiter.api.Test
+    
@org.junit.jupiter.api.condition.DisabledOnOs(org.junit.jupiter.api.condition.OS.WINDOWS)
+    public void readShouldNotReadPartiallyWhenDeletingConcurrentlyBigBlob() 
throws Exception {
+        
DeleteBlobStoreContract.super.readShouldNotReadPartiallyWhenDeletingConcurrentlyBigBlob();
+    }
+
+    @Override
+    @org.junit.jupiter.api.Test
+    
@org.junit.jupiter.api.condition.DisabledOnOs(org.junit.jupiter.api.condition.OS.WINDOWS)
+    public void 
readBytesShouldNotReadPartiallyWhenDeletingConcurrentlyBigBlob() throws 
Exception {
+        
DeleteBlobStoreContract.super.readBytesShouldNotReadPartiallyWhenDeletingConcurrentlyBigBlob();
+    }
 }
diff --git 
a/server/blob/blob-file/src/test/java/org/apache/james/blob/file/FileWithFolderHierarchyTest.java
 
b/server/blob/blob-file/src/test/java/org/apache/james/blob/file/FileWithFolderHierarchyTest.java
new file mode 100644
index 0000000000..061cf2b6e1
--- /dev/null
+++ 
b/server/blob/blob-file/src/test/java/org/apache/james/blob/file/FileWithFolderHierarchyTest.java
@@ -0,0 +1,184 @@
+/****************************************************************
+ * Licensed to the Apache Software Foundation (ASF) under one   *
+ * or more contributor license agreements.  See the NOTICE file *
+ * distributed with this work for additional information        *
+ * regarding copyright ownership.  The ASF licenses this file   *
+ * to you under the Apache License, Version 2.0 (the            *
+ * "License"); you may not use this file except in compliance   *
+ * with the License.  You may obtain a copy of the License at   *
+ *                                                              *
+ *   http://www.apache.org/licenses/LICENSE-2.0                 *
+ *                                                              *
+ * Unless required by applicable law or agreed to in writing,   *
+ * software distributed under the License is distributed on an  *
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY       *
+ * KIND, either express or implied.  See the License for the    *
+ * specific language governing permissions and limitations      *
+ * under the License.                                           *
+ ****************************************************************/
+
+package org.apache.james.blob.file;
+
+import static org.apache.james.blob.api.BlobStore.StoragePolicy.LOW_COST;
+import static org.assertj.core.api.Assertions.assertThat;
+
+import java.io.File;
+import java.nio.charset.StandardCharsets;
+import java.time.Instant;
+import java.util.List;
+import java.util.UUID;
+
+import org.apache.commons.io.FileUtils;
+import org.apache.james.blob.api.BlobId;
+import org.apache.james.blob.api.BlobStore;
+import org.apache.james.blob.api.BlobStoreContract;
+import org.apache.james.blob.api.BucketName;
+import org.apache.james.blob.api.PlainBlobId;
+import org.apache.james.server.blob.deduplication.GenerationAwareBlobId;
+import org.apache.james.server.blob.deduplication.MinIOGenerationAwareBlobId;
+import org.apache.james.server.core.filesystem.FileSystemImpl;
+import org.apache.james.utils.UpdatableTickingClock;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Nested;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.MethodSource;
+
+import reactor.core.publisher.Flux;
+import reactor.core.publisher.Mono;
+
+class FileWithFolderHierarchyTest implements BlobStoreContract {
+    private static final UpdatableTickingClock CLOCK = new 
UpdatableTickingClock(Instant.parse("2021-08-19T10:15:30.00Z"));
+
+    private FileSystemImpl fileSystem;
+    private FileBlobStoreDAO fileBlobStoreDAO;
+    private BlobStore testee;
+    private BlobId.Factory blobIdFactory;
+
+    @BeforeEach
+    void beforeEach() throws Exception {
+        fileSystem = FileSystemImpl.forTesting();
+        blobIdFactory = new MinIOGenerationAwareBlobId.Factory(CLOCK, 
GenerationAwareBlobId.Configuration.DEFAULT, new PlainBlobId.Factory());
+        fileBlobStoreDAO = new FileBlobStoreDAO(fileSystem, blobIdFactory);
+        testee = createBlobStore(blobIdFactory);
+    }
+
+    @AfterEach
+    void tearDown() throws Exception {
+        FileUtils.deleteQuietly(fileSystem.getFile("file://var/blob"));
+    }
+
+    @Override
+    public BlobStore testee() {
+        return testee;
+    }
+
+    @Override
+    public BlobId.Factory blobIdFactory() {
+        return blobIdFactory;
+    }
+
+    public BlobStore createBlobStore(BlobId.Factory blobIdFactory) {
+        return new FileBlobStoreFactory(fileSystem).builder()
+            .blobIdFactory(blobIdFactory)
+            .defaultBucketName()
+            .deduplication();
+    }
+
+    @ParameterizedTest
+    @MethodSource("storagePolicies")
+    void saveShouldReturnBlobIdOfString(BlobStore.StoragePolicy storagePolicy) 
{
+        BlobStore store = testee();
+        BucketName defaultBucketName = store.getDefaultBucketName();
+
+        BlobId blobId = Mono.from(store.save(defaultBucketName, "toto", 
storagePolicy)).block();
+        String blobIdString = blobId.asString();
+
+        assertThat(blobIdString).isEqualTo("1/628/M/f/emXjFVhqwZi9eYtmKc5A");
+        assertThat(blobId).isEqualTo(blobIdFactory().parse(blobIdString));
+    }
+
+    @Test
+    void deleteShouldDeleteBlobFile() throws Exception {
+        BlobStore store = testee();
+        BucketName defaultBucketName = store.getDefaultBucketName();
+
+        BlobId blobId = Mono.from(store.save(defaultBucketName, "toto", 
LOW_COST)).block();
+        File blobFile = new File(fileSystem.getFile("file://var/blob/" + 
defaultBucketName.asString()), blobId.asString());
+        assertThat(blobFile).exists();
+
+        Mono.from(fileBlobStoreDAO.delete(defaultBucketName, blobId)).block();
+
+        assertThat(blobFile).doesNotExist();
+    }
+
+    @Nested
+    class Compatible {
+
+        private BlobStore withGenerationAwareBlobId;
+        private BlobStore withMinIOGenerationAwareBlobId;
+        private BucketName defaultBucketName;
+
+        @BeforeEach
+        void setup() {
+            BlobId.Factory plainBlobIdFactory = new PlainBlobId.Factory();
+            withGenerationAwareBlobId = createBlobStore(new 
GenerationAwareBlobId.Factory(CLOCK, plainBlobIdFactory, 
GenerationAwareBlobId.Configuration.DEFAULT));
+            withMinIOGenerationAwareBlobId = createBlobStore(new 
MinIOGenerationAwareBlobId.Factory(CLOCK, 
GenerationAwareBlobId.Configuration.DEFAULT, plainBlobIdFactory));
+            defaultBucketName = 
withGenerationAwareBlobId.getDefaultBucketName();
+        }
+
+        @Test
+        void 
readWithMinIOGenerationAwareShouldSuccessWhenBlobWasStoredByGenerationAware() {
+            String originalData = "toto" + UUID.randomUUID();
+            BlobId blobId = 
Mono.from(withGenerationAwareBlobId.save(defaultBucketName, originalData, 
LOW_COST)).block();
+
+            assertThat(blobId).isInstanceOf(GenerationAwareBlobId.class);
+
+            byte[] readAsByte = 
Mono.from(withMinIOGenerationAwareBlobId.readBytes(defaultBucketName, 
blobId)).block();
+
+            assertThat(new String(readAsByte, 
StandardCharsets.UTF_8)).isEqualTo(originalData);
+        }
+
+        @Test
+        void 
listBlobsShouldReturnCorrectBlobIdWhenBlobWasStoredByGenerationAware() {
+            String originalData = "toto" + UUID.randomUUID();
+            BlobId blobId = 
Mono.from(withGenerationAwareBlobId.save(defaultBucketName, originalData, 
LOW_COST)).block();
+            assertThat(blobId).isInstanceOf(GenerationAwareBlobId.class);
+
+            List<BlobId> blobIdList = 
Flux.from(withMinIOGenerationAwareBlobId.listBlobs(defaultBucketName)).collectList().block();
+            assertThat(blobIdList).hasSize(1);
+            
assertThat(blobIdList.getFirst()).isInstanceOf(GenerationAwareBlobId.class);
+
+            byte[] readAsByte = 
Mono.from(withMinIOGenerationAwareBlobId.readBytes(defaultBucketName, 
blobIdList.getFirst())).block();
+
+            assertThat(new String(readAsByte, 
StandardCharsets.UTF_8)).isEqualTo(originalData);
+        }
+
+        @Test
+        void 
readWithGenerationAwareShouldSuccessWhenBlobWasStoredByMinIOGenerationAware() {
+            String originalData = "toto" + UUID.randomUUID();
+            BlobId blobId = 
Mono.from(withMinIOGenerationAwareBlobId.save(defaultBucketName, originalData, 
LOW_COST)).block();
+            assertThat(blobId).isInstanceOf(MinIOGenerationAwareBlobId.class);
+
+            byte[] readAsByte = 
Mono.from(withGenerationAwareBlobId.readBytes(defaultBucketName, 
blobId)).block();
+
+            assertThat(new String(readAsByte, 
StandardCharsets.UTF_8)).isEqualTo(originalData);
+        }
+
+        @Test
+        void 
listBlobsShouldReturnCorrectBlobIdWhenBlobWasStoredByMinIOGenerationAware() {
+            String originalData = "toto" + UUID.randomUUID();
+            BlobId blobId = 
Mono.from(withMinIOGenerationAwareBlobId.save(defaultBucketName, originalData, 
LOW_COST)).block();
+            assertThat(blobId).isInstanceOf(MinIOGenerationAwareBlobId.class);
+
+            List<BlobId> blobIdList = 
Flux.from(withGenerationAwareBlobId.listBlobs(defaultBucketName)).collectList().block();
+            assertThat(blobIdList).hasSize(1);
+            
assertThat(blobIdList.getFirst()).isInstanceOf(GenerationAwareBlobId.class);
+
+            byte[] readAsByte = 
Mono.from(withGenerationAwareBlobId.readBytes(defaultBucketName, 
blobIdList.getFirst())).block();
+
+            assertThat(new String(readAsByte, 
StandardCharsets.UTF_8)).isEqualTo(originalData);
+        }
+    }
+}
diff --git 
a/server/container/guice/blob/deduplication-gc/src/main/java/org/apache/james/modules/blobstore/BlobDeduplicationGCModule.java
 
b/server/container/guice/blob/deduplication-gc/src/main/java/org/apache/james/modules/blobstore/BlobDeduplicationGCModule.java
index 073c1aa91f..1d4ec6143b 100644
--- 
a/server/container/guice/blob/deduplication-gc/src/main/java/org/apache/james/modules/blobstore/BlobDeduplicationGCModule.java
+++ 
b/server/container/guice/blob/deduplication-gc/src/main/java/org/apache/james/modules/blobstore/BlobDeduplicationGCModule.java
@@ -70,14 +70,16 @@ public class BlobDeduplicationGCModule extends 
AbstractModule {
     @Singleton
     @Provides
     public BlobId.Factory generationAwareBlobIdFactory(Clock clock, 
PlainBlobId.Factory delegate, GenerationAwareBlobId.Configuration 
configuration) {
-            String property = 
System.getProperty("james.s3.minio.compatibility.mode");
-            boolean compatibilityModeActivated = 
Optional.ofNullable(property).map(Boolean::parseBoolean).orElse(false);
+        boolean compatibilityModeActivated = 
Optional.ofNullable(System.getProperty("james.blobstore.folder.hierarchy"))
+            .or(() -> 
Optional.ofNullable(System.getProperty("james.s3.minio.compatibility.mode")))
+            .map(Boolean::parseBoolean)
+            .orElse(false);
 
-            if (compatibilityModeActivated) {
-                return new MinIOGenerationAwareBlobId.Factory(clock, 
configuration, delegate);
-            } else {
-                return new GenerationAwareBlobId.Factory(clock, delegate, 
configuration);
-            }
+        if (compatibilityModeActivated) {
+            return new MinIOGenerationAwareBlobId.Factory(clock, 
configuration, delegate);
+        } else {
+            return new GenerationAwareBlobId.Factory(clock, delegate, 
configuration);
+        }
     }
 
     @Singleton
diff --git 
a/server/container/guice/blob/deduplication-gc/src/test/java/org/apache/james/modules/blobstore/BlobDeduplicationGCModuleTest.java
 
b/server/container/guice/blob/deduplication-gc/src/test/java/org/apache/james/modules/blobstore/BlobDeduplicationGCModuleTest.java
new file mode 100644
index 0000000000..38877653d7
--- /dev/null
+++ 
b/server/container/guice/blob/deduplication-gc/src/test/java/org/apache/james/modules/blobstore/BlobDeduplicationGCModuleTest.java
@@ -0,0 +1,99 @@
+/****************************************************************
+ * Licensed to the Apache Software Foundation (ASF) under one   *
+ * or more contributor license agreements.  See the NOTICE file *
+ * distributed with this work for additional information        *
+ * regarding copyright ownership.  The ASF licenses this file   *
+ * to you under the Apache License, Version 2.0 (the            *
+ * "License"); you may not use this file except in compliance   *
+ * with the License.  You may obtain a copy of the License at   *
+ *                                                              *
+ *   http://www.apache.org/licenses/LICENSE-2.0                 *
+ *                                                              *
+ * Unless required by applicable law or agreed to in writing,   *
+ * software distributed under the License is distributed on an  *
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY       *
+ * KIND, either express or implied.  See the License for the    *
+ * specific language governing permissions and limitations      *
+ * under the License.                                           *
+ ****************************************************************/
+
+package org.apache.james.modules.blobstore;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+import java.time.Clock;
+
+import org.apache.james.blob.api.BlobId;
+import org.apache.james.blob.api.PlainBlobId;
+import org.apache.james.server.blob.deduplication.GenerationAwareBlobId;
+import org.apache.james.server.blob.deduplication.MinIOGenerationAwareBlobId;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+class BlobDeduplicationGCModuleTest {
+
+    private static final String FOLDER_HIERARCHY_PROPERTY = 
"james.blobstore.folder.hierarchy";
+    private static final String S3_MINIO_COMPATIBILITY_PROPERTY = 
"james.s3.minio.compatibility.mode";
+
+    private BlobDeduplicationGCModule module;
+
+    @BeforeEach
+    @AfterEach
+    void clearProperties() {
+        System.clearProperty(FOLDER_HIERARCHY_PROPERTY);
+        System.clearProperty(S3_MINIO_COMPATIBILITY_PROPERTY);
+    }
+
+    @BeforeEach
+    void setUp() {
+        module = new BlobDeduplicationGCModule();
+    }
+
+    @Test
+    void shouldReturnDefaultFactoryWhenNoPropertiesSet() {
+        BlobId.Factory factory = module.generationAwareBlobIdFactory(
+            Clock.systemUTC(),
+            new PlainBlobId.Factory(),
+            GenerationAwareBlobId.Configuration.DEFAULT);
+
+        
assertThat(factory).isExactlyInstanceOf(GenerationAwareBlobId.Factory.class);
+    }
+
+    @Test
+    void shouldReturnMinIOFactoryWhenFolderHierarchyPropertyIsTrue() {
+        System.setProperty(FOLDER_HIERARCHY_PROPERTY, "true");
+
+        BlobId.Factory factory = module.generationAwareBlobIdFactory(
+            Clock.systemUTC(),
+            new PlainBlobId.Factory(),
+            GenerationAwareBlobId.Configuration.DEFAULT);
+
+        
assertThat(factory).isExactlyInstanceOf(MinIOGenerationAwareBlobId.Factory.class);
+    }
+
+    @Test
+    void shouldReturnMinIOFactoryWhenS3MinioCompatibilityModeIsTrue() {
+        System.setProperty(S3_MINIO_COMPATIBILITY_PROPERTY, "true");
+
+        BlobId.Factory factory = module.generationAwareBlobIdFactory(
+            Clock.systemUTC(),
+            new PlainBlobId.Factory(),
+            GenerationAwareBlobId.Configuration.DEFAULT);
+
+        
assertThat(factory).isExactlyInstanceOf(MinIOGenerationAwareBlobId.Factory.class);
+    }
+
+    @Test
+    void shouldPrioritizeFolderHierarchyPropertyOverCompatibilityMode() {
+        System.setProperty(FOLDER_HIERARCHY_PROPERTY, "false");
+        System.setProperty(S3_MINIO_COMPATIBILITY_PROPERTY, "true");
+
+        BlobId.Factory factory = module.generationAwareBlobIdFactory(
+            Clock.systemUTC(),
+            new PlainBlobId.Factory(),
+            GenerationAwareBlobId.Configuration.DEFAULT);
+
+        
assertThat(factory).isExactlyInstanceOf(GenerationAwareBlobId.Factory.class);
+    }
+}


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

Reply via email to