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]