This is an automated email from the ASF dual-hosted git repository. Arsnael pushed a commit to branch master in repository https://gitbox.apache.org/repos/asf/james-project.git
commit 3977b551b17171a9f6c99a51d40e5d9ca16fff87 Author: Benoit TELLIER <[email protected]> AuthorDate: Fri Jul 24 09:08:53 2026 +0200 JAMES-4209 Leverage blob prefixes --- .../servers/pages/distributed/operate/backup.adoc | 12 ++++++++ server/apps/distributed-app/README.adoc | 11 +++++++ .../org/apache/james/RecoveryConfiguration.java | 36 ++++++++++++++-------- .../java/org/apache/james/S3RecoveryService.java | 7 +++-- .../apache/james/RecoveryConfigurationTest.java | 11 +++++++ .../org/apache/james/blob/aes/AESBlobStoreDAO.java | 5 +++ .../java/org/apache/james/blob/api/BlobStore.java | 12 ++++++++ .../org/apache/james/blob/api/BlobStoreDAO.java | 13 ++++++++ .../apache/james/blob/api/MetricableBlobStore.java | 5 +++ .../blob/api/ReadSaveBlobStoreDAOContract.java | 24 +++++++++++++++ .../blob/cassandra/cache/CachedBlobStore.java | 5 +++ .../blob/objectstorage/aws/S3BlobStoreDAO.java | 13 +++++++- .../deduplication/DeDuplicationBlobStore.scala | 2 ++ .../blob/deduplication/PassThroughBlobStore.scala | 2 ++ .../apache/james/blob/zstd/ZstdBlobStoreDAO.java | 5 +++ 15 files changed, 146 insertions(+), 17 deletions(-) diff --git a/docs/modules/servers/pages/distributed/operate/backup.adoc b/docs/modules/servers/pages/distributed/operate/backup.adoc index e11cc98ad4..961545a29e 100644 --- a/docs/modules/servers/pages/distributed/operate/backup.adoc +++ b/docs/modules/servers/pages/distributed/operate/backup.adoc @@ -163,6 +163,18 @@ restricts recovery to messages whose `Date` header is strictly after the given i ... org.apache.james.S3RecoveryMain --restore-after=2026-01-01T00:00:00Z ---- +An optional `--header-blob-prefix=<prefix>` argument (also settable via the +`RECOVERY_HEADER_BLOB_PREFIX` environment variable or the `recovery.header.blob.prefix` system +property) narrows the walk to the recovery sidecars whose header blob id starts with the given prefix, +pushing the filter down to S3's `ListObjectsV2`. Header blob ids are generation-aware +(`family_generation_...`), so a clever admin can pass e.g. `1_42_` to iterate solely the latest +generation instead of scanning the whole bucket: + +[source,bash] +---- +... org.apache.james.S3RecoveryMain --header-blob-prefix=1_42_ +---- + Notes: * Restored messages are re-stored (and get a fresh `recovery/` sidecar), so re-running the recovery diff --git a/server/apps/distributed-app/README.adoc b/server/apps/distributed-app/README.adoc index b81566e744..4acd0061f8 100644 --- a/server/apps/distributed-app/README.adoc +++ b/server/apps/distributed-app/README.adoc @@ -152,6 +152,17 @@ environment variable, or the `restore.messages.after` system property: $ java ... org.apache.james.S3RecoveryMain --restore-after=2026-01-01T00:00:00Z ---- +An optional `--header-blob-prefix=<prefix>` argument (also settable via the `RECOVERY_HEADER_BLOB_PREFIX` +environment variable or the `recovery.header.blob.prefix` system property) narrows the walk to the +recovery sidecars whose header blob id starts with the given prefix, and lets S3 filter server-side. +Because header blob ids are generation-aware (`family_generation_...`), a clever admin can pass e.g. +`1_42_` to iterate solely the latest generation instead of scanning the whole bucket: + +[source] +---- +$ java ... org.apache.james.S3RecoveryMain --header-blob-prefix=1_42_ +---- + Notes: * Restored messages are re-stored (and get a fresh `recovery/` sidecar), so re-running the recovery diff --git a/server/apps/distributed-app/src/main/java/org/apache/james/RecoveryConfiguration.java b/server/apps/distributed-app/src/main/java/org/apache/james/RecoveryConfiguration.java index 1234bf9231..860f314812 100644 --- a/server/apps/distributed-app/src/main/java/org/apache/james/RecoveryConfiguration.java +++ b/server/apps/distributed-app/src/main/java/org/apache/james/RecoveryConfiguration.java @@ -28,29 +28,39 @@ import java.util.Optional; * Configuration for the S3 blob store recovery run. * * <p>The optional {@code restoreAfter} instant restricts recovery to messages whose {@code Date} - * header is strictly after the given point in time. It can be provided (highest precedence first) as:</p> - * <ul> - * <li>a {@code --restore-after=<ISO-8601 instant>} program argument</li> - * <li>the {@code RESTORE_MESSAGES_AFTER} environment variable</li> - * <li>the {@code restore.messages.after} system property</li> - * </ul> + * header is strictly after the given point in time. It can be provided (highest precedence first) as a + * {@code --restore-after=<ISO-8601 instant>} program argument, the {@code RESTORE_MESSAGES_AFTER} + * environment variable, or the {@code restore.messages.after} system property.</p> + * + * <p>The optional {@code headerBlobPrefix} narrows the walk to the recovery sidecars whose header blob + * id starts with the given prefix. Because header blob ids are generation-aware + * ({@code family_generation_...}), a clever admin can pass e.g. {@code 1_42_} to iterate solely the + * latest generation instead of scanning the whole bucket. It defaults to the empty string (all + * recovery sidecars) and can be provided as a {@code --header-blob-prefix=<prefix>} program argument, + * the {@code RECOVERY_HEADER_BLOB_PREFIX} environment variable, or the {@code recovery.header.blob.prefix} + * system property.</p> */ -public record RecoveryConfiguration(Optional<Instant> restoreAfter) { +public record RecoveryConfiguration(Optional<Instant> restoreAfter, String headerBlobPrefix) { private static final String RESTORE_AFTER_ARG = "--restore-after="; private static final String RESTORE_AFTER_ENV = "RESTORE_MESSAGES_AFTER"; private static final String RESTORE_AFTER_PROPERTY = "restore.messages.after"; + private static final String HEADER_BLOB_PREFIX_ARG = "--header-blob-prefix="; + private static final String HEADER_BLOB_PREFIX_ENV = "RECOVERY_HEADER_BLOB_PREFIX"; + private static final String HEADER_BLOB_PREFIX_PROPERTY = "recovery.header.blob.prefix"; public static RecoveryConfiguration parse(String[] args) { - return new RecoveryConfiguration(restoreAfter(args).map(RecoveryConfiguration::parseInstant)); + return new RecoveryConfiguration( + option(args, RESTORE_AFTER_ARG, RESTORE_AFTER_ENV, RESTORE_AFTER_PROPERTY).map(RecoveryConfiguration::parseInstant), + option(args, HEADER_BLOB_PREFIX_ARG, HEADER_BLOB_PREFIX_ENV, HEADER_BLOB_PREFIX_PROPERTY).orElse("")); } - private static Optional<String> restoreAfter(String[] args) { + private static Optional<String> option(String[] args, String argPrefix, String envName, String propertyName) { return Arrays.stream(args) - .filter(arg -> arg.startsWith(RESTORE_AFTER_ARG)) - .map(arg -> arg.substring(RESTORE_AFTER_ARG.length())) + .filter(arg -> arg.startsWith(argPrefix)) + .map(arg -> arg.substring(argPrefix.length())) .findFirst() - .or(() -> Optional.ofNullable(System.getenv(RESTORE_AFTER_ENV))) - .or(() -> Optional.ofNullable(System.getProperty(RESTORE_AFTER_PROPERTY))) + .or(() -> Optional.ofNullable(System.getenv(envName))) + .or(() -> Optional.ofNullable(System.getProperty(propertyName))) .map(String::trim) .filter(value -> !value.isEmpty()); } diff --git a/server/apps/distributed-app/src/main/java/org/apache/james/S3RecoveryService.java b/server/apps/distributed-app/src/main/java/org/apache/james/S3RecoveryService.java index 703ddc10ce..d69b88daa4 100644 --- a/server/apps/distributed-app/src/main/java/org/apache/james/S3RecoveryService.java +++ b/server/apps/distributed-app/src/main/java/org/apache/james/S3RecoveryService.java @@ -117,10 +117,11 @@ public class S3RecoveryService { public Mono<Report> run() { BucketName bucket = blobStore.getDefaultBucketName(); - LOGGER.info("Starting S3 recovery on bucket {} (restore after: {})", bucket.asString(), configuration.restoreAfter()); - return Flux.from(blobStoreDAO.listBlobs(bucket)) + String prefix = RECOVERY_BLOB_PREFIX + configuration.headerBlobPrefix(); + LOGGER.info("Starting S3 recovery on bucket {} (prefix: {}, restore after: {})", + bucket.asString(), prefix, configuration.restoreAfter()); + return Flux.from(blobStoreDAO.listBlobs(bucket, prefix)) .map(BlobId::asString) - .filter(key -> key.startsWith(RECOVERY_BLOB_PREFIX)) .flatMap(recoveryKey -> restoreOne(bucket, recoveryKey), CONCURRENCY) .reduce(Report.empty(), Report::merge) .doOnNext(report -> LOGGER.info("S3 recovery finished: {}", report)); diff --git a/server/apps/distributed-app/src/test/java/org/apache/james/RecoveryConfigurationTest.java b/server/apps/distributed-app/src/test/java/org/apache/james/RecoveryConfigurationTest.java index 8fdd7f0a37..df3dae23d4 100644 --- a/server/apps/distributed-app/src/test/java/org/apache/james/RecoveryConfigurationTest.java +++ b/server/apps/distributed-app/src/test/java/org/apache/james/RecoveryConfigurationTest.java @@ -43,4 +43,15 @@ class RecoveryConfigurationTest { assertThatThrownBy(() -> RecoveryConfiguration.parse(new String[] {"--restore-after=not-a-date"})) .isInstanceOf(IllegalArgumentException.class); } + + @Test + void parseShouldDefaultHeaderBlobPrefixToEmpty() { + assertThat(RecoveryConfiguration.parse(new String[] {}).headerBlobPrefix()).isEmpty(); + } + + @Test + void parseShouldReadHeaderBlobPrefixArgument() { + assertThat(RecoveryConfiguration.parse(new String[] {"--header-blob-prefix=1_42_"}).headerBlobPrefix()) + .isEqualTo("1_42_"); + } } diff --git a/server/blob/blob-aes/src/main/java/org/apache/james/blob/aes/AESBlobStoreDAO.java b/server/blob/blob-aes/src/main/java/org/apache/james/blob/aes/AESBlobStoreDAO.java index 16dba25b4e..24672e1521 100644 --- a/server/blob/blob-aes/src/main/java/org/apache/james/blob/aes/AESBlobStoreDAO.java +++ b/server/blob/blob-aes/src/main/java/org/apache/james/blob/aes/AESBlobStoreDAO.java @@ -206,4 +206,9 @@ public class AESBlobStoreDAO implements BlobStoreDAO { public Publisher<BlobId> listBlobs(BucketName bucketName) { return underlying.listBlobs(bucketName); } + + @Override + public Publisher<BlobId> listBlobs(BucketName bucketName, String prefix) { + return underlying.listBlobs(bucketName, prefix); + } } diff --git a/server/blob/blob-api/src/main/java/org/apache/james/blob/api/BlobStore.java b/server/blob/blob-api/src/main/java/org/apache/james/blob/api/BlobStore.java index 0493e27743..e6089fc125 100644 --- a/server/blob/blob-api/src/main/java/org/apache/james/blob/api/BlobStore.java +++ b/server/blob/blob-api/src/main/java/org/apache/james/blob/api/BlobStore.java @@ -25,6 +25,7 @@ import org.reactivestreams.Publisher; import com.google.common.io.ByteSource; +import reactor.core.publisher.Flux; import reactor.util.function.Tuple2; /** @@ -97,4 +98,15 @@ public interface BlobStore { Publisher<Boolean> delete(BucketName bucketName, BlobId blobId); Publisher<BlobId> listBlobs(BucketName bucketName); + + /** + * Lists the blobs of a bucket whose id starts with the given prefix. + * + * <p>The default implementation filters the full listing. Implementations backed by a store able to + * push the prefix down (eg. S3 {@code ListObjectsV2}) should override this for efficiency.</p> + */ + default Publisher<BlobId> listBlobs(BucketName bucketName, String prefix) { + return Flux.from(listBlobs(bucketName)) + .filter(blobId -> blobId.asString().startsWith(prefix)); + } } diff --git a/server/blob/blob-api/src/main/java/org/apache/james/blob/api/BlobStoreDAO.java b/server/blob/blob-api/src/main/java/org/apache/james/blob/api/BlobStoreDAO.java index 41d2c257d2..223dafe3d8 100644 --- a/server/blob/blob-api/src/main/java/org/apache/james/blob/api/BlobStoreDAO.java +++ b/server/blob/blob-api/src/main/java/org/apache/james/blob/api/BlobStoreDAO.java @@ -38,6 +38,8 @@ import com.google.common.collect.ImmutableMap; import com.google.common.io.ByteSource; import com.google.common.io.FileBackedOutputStream; +import reactor.core.publisher.Flux; + /** * James virtual blob store abstraction. * @@ -291,4 +293,15 @@ public interface BlobStoreDAO { Publisher<BucketName> listBuckets(); Publisher<BlobId> listBlobs(BucketName bucketName); + + /** + * Lists the blobs of a bucket whose id starts with the given prefix (eg. {@link #RECOVERY_BLOB_PREFIX}). + * + * <p>The default implementation filters the full listing. Connectors able to push the prefix down to + * their backend (eg. S3 {@code ListObjectsV2}) should override this for efficiency.</p> + */ + default Publisher<BlobId> listBlobs(BucketName bucketName, String prefix) { + return Flux.from(listBlobs(bucketName)) + .filter(blobId -> blobId.asString().startsWith(prefix)); + } } diff --git a/server/blob/blob-api/src/main/java/org/apache/james/blob/api/MetricableBlobStore.java b/server/blob/blob-api/src/main/java/org/apache/james/blob/api/MetricableBlobStore.java index 0d60274937..b87bf1bfae 100644 --- a/server/blob/blob-api/src/main/java/org/apache/james/blob/api/MetricableBlobStore.java +++ b/server/blob/blob-api/src/main/java/org/apache/james/blob/api/MetricableBlobStore.java @@ -131,4 +131,9 @@ public class MetricableBlobStore implements BlobStore { public Publisher<BlobId> listBlobs(BucketName bucketName) { return blobStoreImpl.listBlobs(bucketName); } + + @Override + public Publisher<BlobId> listBlobs(BucketName bucketName, String prefix) { + return blobStoreImpl.listBlobs(bucketName, prefix); + } } diff --git a/server/blob/blob-api/src/test/java/org/apache/james/blob/api/ReadSaveBlobStoreDAOContract.java b/server/blob/blob-api/src/test/java/org/apache/james/blob/api/ReadSaveBlobStoreDAOContract.java index 640afb6a18..4ca331717d 100644 --- a/server/blob/blob-api/src/test/java/org/apache/james/blob/api/ReadSaveBlobStoreDAOContract.java +++ b/server/blob/blob-api/src/test/java/org/apache/james/blob/api/ReadSaveBlobStoreDAOContract.java @@ -295,6 +295,30 @@ public interface ReadSaveBlobStoreDAOContract { .containsOnly(TEST_BLOB_ID.asString(), OTHER_TEST_BLOB_ID.asString()); } + @Test + default void listWithPrefixShouldReturnOnlyMatchingBlobs() { + BlobStoreDAO store = testee(); + Mono.from(store.save(TEST_BUCKET_NAME, TEST_BLOB_ID, ELEVEN_KILOBYTES)).block(); + Mono.from(store.save(TEST_BUCKET_NAME, OTHER_TEST_BLOB_ID, ELEVEN_KILOBYTES)).block(); + + assertThat(Flux.from(testee().listBlobs(TEST_BUCKET_NAME, "test-")) + .map(BlobId::asString) + .collectList() + .block()) + .containsOnly(TEST_BLOB_ID.asString()); + } + + @Test + default void listWithPrefixShouldReturnEmptyWhenNoMatch() { + BlobStoreDAO store = testee(); + Mono.from(store.save(TEST_BUCKET_NAME, TEST_BLOB_ID, ELEVEN_KILOBYTES)).block(); + + assertThat(Flux.from(testee().listBlobs(TEST_BUCKET_NAME, "no-such-prefix-")) + .collectList() + .block()) + .isEmpty(); + } + static Stream<Arguments> blobs() { return Stream.of(new Object[]{"SHORT", SHORT_BYTEARRAY}, new Object[]{"LONG", ELEVEN_KILOBYTES}, new Object[]{"BIG", TWELVE_MEGABYTES}) .map(Arguments::of); diff --git a/server/blob/blob-cassandra/src/main/java/org/apache/james/blob/cassandra/cache/CachedBlobStore.java b/server/blob/blob-cassandra/src/main/java/org/apache/james/blob/cassandra/cache/CachedBlobStore.java index e3c95f7dce..b6c7bdb724 100644 --- a/server/blob/blob-cassandra/src/main/java/org/apache/james/blob/cassandra/cache/CachedBlobStore.java +++ b/server/blob/blob-cassandra/src/main/java/org/apache/james/blob/cassandra/cache/CachedBlobStore.java @@ -405,4 +405,9 @@ public class CachedBlobStore implements BlobStore { public Publisher<BlobId> listBlobs(BucketName bucketName) { return backend.listBlobs(bucketName); } + + @Override + public Publisher<BlobId> listBlobs(BucketName bucketName, String prefix) { + return backend.listBlobs(bucketName, prefix); + } } 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 1ed6a133ab..5176e47429 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 @@ -30,6 +30,7 @@ import java.util.List; import java.util.Map; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionException; +import java.util.function.Consumer; import jakarta.inject.Inject; import jakarta.inject.Singleton; @@ -68,6 +69,7 @@ import software.amazon.awssdk.services.s3.model.DeleteObjectsResponse; import software.amazon.awssdk.services.s3.model.GetObjectRequest; import software.amazon.awssdk.services.s3.model.GetObjectResponse; import software.amazon.awssdk.services.s3.model.ListBucketsResponse; +import software.amazon.awssdk.services.s3.model.ListObjectsV2Request; import software.amazon.awssdk.services.s3.model.ListObjectsV2Response; import software.amazon.awssdk.services.s3.model.NoSuchBucketException; import software.amazon.awssdk.services.s3.model.NoSuchKeyException; @@ -455,7 +457,16 @@ public class S3BlobStoreDAO implements BlobStoreDAO { @Override public Publisher<BlobId> listBlobs(BucketName bucketName) { - return Flux.from(client.listObjectsV2Paginator(builder -> builder.bucket(bucketName.asString()))) + return listBlobs(bucketName, builder -> builder.bucket(bucketName.asString())); + } + + @Override + public Publisher<BlobId> listBlobs(BucketName bucketName, String prefix) { + return listBlobs(bucketName, builder -> builder.bucket(bucketName.asString()).prefix(prefix)); + } + + private Publisher<BlobId> listBlobs(BucketName bucketName, Consumer<ListObjectsV2Request.Builder> request) { + return Flux.from(client.listObjectsV2Paginator(request)) .flatMapIterable(ListObjectsV2Response::contents) .map(S3Object::key) .map(blobIdFactory::parse) diff --git a/server/blob/blob-storage-strategy/src/main/scala/org/apache/james/server/blob/deduplication/DeDuplicationBlobStore.scala b/server/blob/blob-storage-strategy/src/main/scala/org/apache/james/server/blob/deduplication/DeDuplicationBlobStore.scala index 4b1a13e97d..25651d63c9 100644 --- a/server/blob/blob-storage-strategy/src/main/scala/org/apache/james/server/blob/deduplication/DeDuplicationBlobStore.scala +++ b/server/blob/blob-storage-strategy/src/main/scala/org/apache/james/server/blob/deduplication/DeDuplicationBlobStore.scala @@ -194,4 +194,6 @@ class DeDuplicationBlobStore @Inject()(blobStoreDAO: BlobStoreDAO, override def listBuckets(): Publisher[BucketName] = Flux.concat(blobStoreDAO.listBuckets(), Flux.just(defaultBucketName)).distinct() override def listBlobs(bucketName: BucketName): Publisher[BlobId] = blobStoreDAO.listBlobs(bucketName) + + override def listBlobs(bucketName: BucketName, prefix: String): Publisher[BlobId] = blobStoreDAO.listBlobs(bucketName, prefix) } diff --git a/server/blob/blob-storage-strategy/src/main/scala/org/apache/james/server/blob/deduplication/PassThroughBlobStore.scala b/server/blob/blob-storage-strategy/src/main/scala/org/apache/james/server/blob/deduplication/PassThroughBlobStore.scala index 6bff213406..9d5c677b4c 100644 --- a/server/blob/blob-storage-strategy/src/main/scala/org/apache/james/server/blob/deduplication/PassThroughBlobStore.scala +++ b/server/blob/blob-storage-strategy/src/main/scala/org/apache/james/server/blob/deduplication/PassThroughBlobStore.scala @@ -129,4 +129,6 @@ class PassThroughBlobStore @Inject()(blobStoreDAO: BlobStoreDAO, override def listBuckets(): Publisher[BucketName] = Flux.concat(blobStoreDAO.listBuckets(), Flux.just(defaultBucketName)).distinct() override def listBlobs(bucketName: BucketName): Publisher[BlobId] = blobStoreDAO.listBlobs(bucketName) + + override def listBlobs(bucketName: BucketName, prefix: String): Publisher[BlobId] = blobStoreDAO.listBlobs(bucketName, prefix) } diff --git a/server/blob/blob-zstd/src/main/java/org/apache/james/blob/zstd/ZstdBlobStoreDAO.java b/server/blob/blob-zstd/src/main/java/org/apache/james/blob/zstd/ZstdBlobStoreDAO.java index e9e837fae4..a51f79fc2d 100644 --- a/server/blob/blob-zstd/src/main/java/org/apache/james/blob/zstd/ZstdBlobStoreDAO.java +++ b/server/blob/blob-zstd/src/main/java/org/apache/james/blob/zstd/ZstdBlobStoreDAO.java @@ -190,6 +190,11 @@ public class ZstdBlobStoreDAO implements BlobStoreDAO { return underlying.listBlobs(bucketName); } + @Override + public Publisher<BlobId> listBlobs(BucketName bucketName, String prefix) { + return underlying.listBlobs(bucketName, prefix); + } + private InputStreamBlob decompress(InputStreamBlob blob) throws ObjectStoreIOException { try { metricRecorder.recordDecompression(); --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
