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 0ee35cd6c6 [ENHANCEMENT] Store shall preserve ObjectNotFoundException 
(#3222)
0ee35cd6c6 is described below

commit 0ee35cd6c691613b0995c4e0e7968ea1364d0757
Author: Benoit TELLIER <[email protected]>
AuthorDate: Thu Oct 1 11:24:28 2026 +0200

    [ENHANCEMENT] Store shall preserve ObjectNotFoundException (#3222)
---
 .../main/java/org/apache/james/blob/api/Store.java | 11 ++++++++++
 .../james/blob/mail/MimeMessageStoreTest.java      | 24 ++++++++++++++++++++++
 2 files changed, 35 insertions(+)

diff --git 
a/server/blob/blob-common/src/main/java/org/apache/james/blob/api/Store.java 
b/server/blob/blob-common/src/main/java/org/apache/james/blob/api/Store.java
index 229d0925cc..ad97813c76 100644
--- a/server/blob/blob-common/src/main/java/org/apache/james/blob/api/Store.java
+++ b/server/blob/blob-common/src/main/java/org/apache/james/blob/api/Store.java
@@ -110,6 +110,8 @@ public interface Store<T, I> {
                 .flatMap(entry -> readByteSource(bucketName, entry.getValue(), 
entry.getKey().getStoragePolicy())
                     .map(result -> Pair.of(entry.getKey(), result)))
                 .collectMap(Map.Entry::getKey, Pair::getValue)
+                // A downstream onErrorContinue can silently drop a failed 
blob read
+                .flatMap(streams -> ensureAllPartsRead(blobIds, streams))
                 // Critical to correctly propagate errors.
                 // Replacing by `map` would cause the error not to be catch 
downstream. No idea why, failed to reproduce with a test.
                 // Impact: unacknowledged messages for RabbitMQ mailQueue that 
eventually piles up to interruption of service.
@@ -118,6 +120,15 @@ public interface Store<T, I> {
                     .subscribeOn(ReactorUtils.BLOCKING_CALL_WRAPPER));
         }
 
+        private Mono<Map<BlobType, CloseableByteSource>> ensureAllPartsRead(I 
blobIds, Map<BlobType, CloseableByteSource> streams) {
+            if (streams.keySet().containsAll(blobIds.asMap().keySet())) {
+                return Mono.just(streams);
+            }
+            return Mono.fromRunnable(() -> 
streams.forEach(Throwing.biConsumer((blobType, byteSource) -> 
byteSource.close())))
+                .subscribeOn(ReactorUtils.BLOCKING_CALL_WRAPPER)
+                .then(Mono.error(() -> new ObjectNotFoundException("Missing 
blob parts for " + blobIds.asMap())));
+        }
+
         private Mono<CloseableByteSource> readByteSource(BucketName 
bucketName, BlobId blobId, StoragePolicy storagePolicy) {
             return Mono.usingWhen(blobStore.readReactive(bucketName, blobId, 
storagePolicy),
                 Throwing.function(in -> {
diff --git 
a/server/blob/mail-store/src/test/java/org/apache/james/blob/mail/MimeMessageStoreTest.java
 
b/server/blob/mail-store/src/test/java/org/apache/james/blob/mail/MimeMessageStoreTest.java
index 2ecad6ba33..5fb07e07fa 100644
--- 
a/server/blob/mail-store/src/test/java/org/apache/james/blob/mail/MimeMessageStoreTest.java
+++ 
b/server/blob/mail-store/src/test/java/org/apache/james/blob/mail/MimeMessageStoreTest.java
@@ -43,6 +43,7 @@ import org.assertj.core.api.SoftAssertions;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
 
+import reactor.core.publisher.Flux;
 import reactor.core.publisher.Mono;
 
 class MimeMessageStoreTest {
@@ -106,6 +107,29 @@ class MimeMessageStoreTest {
             .isInstanceOf(ObjectNotFoundException.class);
     }
 
+    @Test
+    void 
readShouldThrowObjectNotFoundWhenOnePartIsMissingAndDownstreamUsesOnErrorContinue()
 throws Exception {
+        MimeMessage message = MimeMessageBuilder.mimeMessageBuilder()
+            .addFrom("[email protected]")
+            .addToRecipient("[email protected]")
+            .setSubject("Important Mail")
+            .setText("Important mail content")
+            .build();
+
+        MimeMessagePartsId parts = testee.save(message).block();
+        Mono.from(blobStore.delete(blobStore.getDefaultBucketName(), 
parts.getHeaderBlobId())).block();
+
+        // Mimics JamesMailSpooler: an onErrorContinue downstream must not 
truncate the blob parts read
+        Throwable error = Flux.just(parts)
+            .flatMap(partsId -> testee.read(partsId)
+                .then(Mono.<Throwable>empty())
+                .onErrorResume(Mono::just))
+            .onErrorContinue((e, o) -> { })
+            .blockFirst();
+
+        assertThat(error).isInstanceOf(ObjectNotFoundException.class);
+    }
+
     @Test
     void shouldSupportStoringMimeMessageWrapperWithLFInHeaders() {
         MimeMessageSource mimeMessageSource = new MimeMessageSource() {


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

Reply via email to