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

davsclaus pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git


The following commit(s) were added to refs/heads/main by this push:
     new 3cecd3a1997f CAMEL-25123: camel-seda - run the on completions of a 
message that is not added to the queue (#27032)
3cecd3a1997f is described below

commit 3cecd3a1997fdb6aa73b427c1817d6b25691bfa3
Author: allthingssecurity <[email protected]>
AuthorDate: Tue Sep 29 14:40:20 2026 +0530

    CAMEL-25123: camel-seda - run the on completions of a message that is not 
added to the queue (#27032)
    
    * CAMEL-25123: camel-seda - run the on completions of a message that is not 
added to the queue
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../component/disruptor/DisruptorProducer.java     |  10 +-
 .../DisruptorRingBufferFullOnCompletionTest.java   |  91 +++++++++++
 .../apache/camel/component/seda/SedaProducer.java  |  21 +++
 .../seda/SedaQueueFullOnCompletionTest.java        | 168 +++++++++++++++++++++
 .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc    |  13 ++
 5 files changed, 302 insertions(+), 1 deletion(-)

diff --git 
a/components/camel-disruptor/src/main/java/org/apache/camel/component/disruptor/DisruptorProducer.java
 
b/components/camel-disruptor/src/main/java/org/apache/camel/component/disruptor/DisruptorProducer.java
index 9ac80ff53e7d..9eaa481accb1 100644
--- 
a/components/camel-disruptor/src/main/java/org/apache/camel/component/disruptor/DisruptorProducer.java
+++ 
b/components/camel-disruptor/src/main/java/org/apache/camel/component/disruptor/DisruptorProducer.java
@@ -148,7 +148,15 @@ public class DisruptorProducer extends 
DefaultAsyncProducer {
                 // no wait, eg its a InOnly then just publish to the 
ringbuffer and return
                 // handover the completion so its the copy which performs 
that, as we do not wait
                 final Exchange copy = prepareCopy(exchange, true);
-                doPublish(copy);
+                try {
+                    doPublish(copy);
+                } catch (RuntimeException e) {
+                    // the copy is not published (such as when the ringbuffer 
is full), so the exchange takes back its on
+                    // completions (such as a consumer rolling back the 
message), which also releases the stream cache
+                    // of the copy
+                    copy.getExchangeExtension().handoverCompletions(exchange);
+                    throw e;
+                }
             }
         } catch (Exception e) {
             exchange.setException(e);
diff --git 
a/components/camel-disruptor/src/test/java/org/apache/camel/component/disruptor/DisruptorRingBufferFullOnCompletionTest.java
 
b/components/camel-disruptor/src/test/java/org/apache/camel/component/disruptor/DisruptorRingBufferFullOnCompletionTest.java
new file mode 100644
index 000000000000..d106deb8880d
--- /dev/null
+++ 
b/components/camel-disruptor/src/test/java/org/apache/camel/component/disruptor/DisruptorRingBufferFullOnCompletionTest.java
@@ -0,0 +1,91 @@
+/*
+ * 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.camel.component.disruptor;
+
+import java.util.List;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+
+import org.apache.camel.CamelExecutionException;
+import org.apache.camel.Exchange;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.apache.camel.support.SynchronizationAdapter;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.awaitility.Awaitility;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+
+/**
+ * When the ring buffer is full, the InOnly exchange is not published, and its 
on completions (such as the rollback of
+ * the consumer) must still run.
+ */
+public class DisruptorRingBufferFullOnCompletionTest extends CamelTestSupport {
+
+    private final List<String> events = new CopyOnWriteArrayList<>();
+    private final CountDownLatch release = new CountDownLatch(1);
+
+    @Test
+    void testRingBufferFull() throws Exception {
+        MockEndpoint mock = getMockEndpoint("mock:result");
+        mock.expectedBodiesReceived("A");
+
+        // the consumer holds the only slot of the ring buffer until it is 
released
+        template.sendBody("direct:start", "A");
+        try {
+            CamelExecutionException e = 
assertThrows(CamelExecutionException.class,
+                    () -> template.sendBody("direct:start", "B"));
+            assertInstanceOf(IllegalStateException.class, e.getCause());
+            assertEquals(List.of("failure:B"), events);
+        } finally {
+            release.countDown();
+        }
+        MockEndpoint.assertIsSatisfied(context);
+        Awaitility.await().atMost(5, TimeUnit.SECONDS)
+                .untilAsserted(() -> assertEquals(List.of("failure:B", 
"complete:A"), events));
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                from("direct:start").errorHandler(noErrorHandler())
+                        .process(e -> 
e.getExchangeExtension().addOnCompletion(new SynchronizationAdapter() {
+                            @Override
+                            public void onComplete(Exchange exchange) {
+                                events.add("complete:" + 
exchange.getMessage().getBody(String.class));
+                            }
+
+                            @Override
+                            public void onFailure(Exchange exchange) {
+                                events.add("failure:" + 
exchange.getMessage().getBody(String.class));
+                            }
+                        }))
+                        .to("disruptor:full?size=1&blockWhenFull=false");
+
+                from("disruptor:full?size=1")
+                        .process(e -> release.await(20, TimeUnit.SECONDS))
+                        .to("mock:result");
+            }
+        };
+    }
+}
diff --git 
a/components/camel-seda/src/main/java/org/apache/camel/component/seda/SedaProducer.java
 
b/components/camel-seda/src/main/java/org/apache/camel/component/seda/SedaProducer.java
index 55aa8dbc95b7..cd4e94808f69 100644
--- 
a/components/camel-seda/src/main/java/org/apache/camel/component/seda/SedaProducer.java
+++ 
b/components/camel-seda/src/main/java/org/apache/camel/component/seda/SedaProducer.java
@@ -263,12 +263,32 @@ public class SedaProducer extends DefaultAsyncProducer {
         }
 
         LOG.trace("Adding Exchange to queue: {}", target);
+        boolean added = false;
+        try {
+            added = offerToQueue(queue, target);
+        } finally {
+            if (copy && !added) {
+                // the copy is not queued (discarded, or failed to be added), 
so the exchange takes back its on
+                // completions (such as a consumer committing or rolling back 
the message), which also releases the
+                // stream cache of the copy
+                target.getExchangeExtension().handoverCompletions(exchange);
+            }
+        }
+    }
+
+    /**
+     * Adds the exchange to the queue
+     *
+     * @return {@code false} if the exchange is discarded as the queue is full
+     */
+    private boolean offerToQueue(BlockingQueue<Exchange> queue, Exchange 
target) {
         if (discardWhenFull) {
             try {
                 boolean added = queue.offer(target, 0, TimeUnit.MILLISECONDS);
                 if (!added) {
                     LOG.trace("Discarding Exchange as queue is full: {}", 
target);
                 }
+                return added;
             } catch (InterruptedException e) {
                 LOG.debug("Offer interrupted, are we stopping? {}", 
isStopping() || isStopped());
                 Thread.currentThread().interrupt();
@@ -298,6 +318,7 @@ public class SedaProducer extends DefaultAsyncProducer {
         } else {
             queue.add(target);
         }
+        return true;
     }
 
     private static RejectedExecutionException 
interruptedWhileAddingToQueue(InterruptedException cause) {
diff --git 
a/core/camel-core/src/test/java/org/apache/camel/component/seda/SedaQueueFullOnCompletionTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/component/seda/SedaQueueFullOnCompletionTest.java
new file mode 100644
index 000000000000..eca317867ddf
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/component/seda/SedaQueueFullOnCompletionTest.java
@@ -0,0 +1,168 @@
+/*
+ * 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.camel.component.seda;
+
+import java.io.BufferedInputStream;
+import java.io.ByteArrayInputStream;
+import java.io.File;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.util.List;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.TimeUnit;
+
+import org.apache.camel.CamelExecutionException;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Exchange;
+import org.apache.camel.Processor;
+import org.apache.camel.builder.NotifyBuilder;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.apache.camel.support.SynchronizationAdapter;
+import org.awaitility.Awaitility;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * When the SEDA producer does not add the copy of an InOnly exchange to the 
queue (discardWhenFull, offerTimeout, or a
+ * full queue), the on completions of the exchange (such as the commit or 
rollback of the consumer) must still run.
+ */
+public class SedaQueueFullOnCompletionTest extends ContextTestSupport {
+
+    private static final byte[] DATA = new byte[16 * 1024];
+
+    private final List<String> events = new CopyOnWriteArrayList<>();
+
+    @Test
+    public void testDiscardWhenFull() {
+        template.sendBody("direct:discard", "A");
+        // the queue is full, so this one is discarded, which completes the 
exchange
+        template.sendBody("direct:discard", "B");
+
+        assertEquals(List.of("complete:B"), events);
+    }
+
+    @Test
+    public void testOfferTimeout() {
+        template.sendBody("direct:offer", "A");
+        CamelExecutionException e = assertThrows(CamelExecutionException.class,
+                () -> template.sendBody("direct:offer", "B"));
+        assertInstanceOf(IllegalStateException.class, e.getCause());
+
+        assertEquals(List.of("failure:B"), events);
+    }
+
+    @Test
+    public void testQueueFull() {
+        template.sendBody("direct:add", "A");
+        CamelExecutionException e = assertThrows(CamelExecutionException.class,
+                () -> template.sendBody("direct:add", "B"));
+        assertInstanceOf(IllegalStateException.class, e.getCause());
+
+        assertEquals(List.of("failure:B"), events);
+    }
+
+    @Test
+    public void testDiscardWhenFullFileConsumer() throws Exception {
+        Path in = testDirectory("in", true);
+        NotifyBuilder notify = new 
NotifyBuilder(context).fromRoute("files").whenDone(2).create();
+        Files.writeString(in.resolve("a.txt"), "A");
+        Files.writeString(in.resolve("b.txt"), "B");
+        context.getRouteController().startRoute("files");
+        assertTrue(notify.matches(10, TimeUnit.SECONDS));
+
+        // a.txt is queued, and b.txt is discarded, which commits it (moves it 
to .camel)
+        Awaitility.await().atMost(5, TimeUnit.SECONDS)
+                .untilAsserted(() -> 
assertTrue(Files.exists(in.resolve(".camel/b.txt"))));
+        assertTrue(Files.exists(in.resolve("a.txt")));
+
+        MockEndpoint mock = getMockEndpoint("mock:files");
+        mock.expectedBodiesReceived("A");
+        context.getRouteController().startRoute("filesQueue");
+        assertMockEndpointsSatisfied();
+        Awaitility.await().atMost(5, TimeUnit.SECONDS)
+                .untilAsserted(() -> 
assertTrue(Files.exists(in.resolve(".camel/a.txt"))));
+    }
+
+    @Test
+    public void testDiscardWhenFullStreamCache() throws Exception {
+        template.sendBody("direct:spool", new BufferedInputStream(new 
ByteArrayInputStream(DATA)));
+        // discarded: the stream cache of its copy is released as well
+        template.sendBody("direct:spool", new BufferedInputStream(new 
ByteArrayInputStream(DATA)));
+
+        MockEndpoint mock = getMockEndpoint("mock:spool");
+        mock.expectedMessageCount(1);
+        context.getRouteController().startRoute("spoolQueue");
+        assertMockEndpointsSatisfied();
+        assertArrayEquals(DATA, 
mock.getReceivedExchanges().get(0).getMessage().getBody(byte[].class));
+
+        File spoolDir = testDirectory("spool").toFile();
+        Awaitility.await().atMost(5, TimeUnit.SECONDS).untilAsserted(() -> {
+            String[] files = spoolDir.list();
+            assertNotNull(files);
+            assertEquals(0, files.length, "Spool files left behind: " + 
List.of(files));
+        });
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                
context.getStreamCachingStrategy().setSpoolDirectory(testDirectory("spool").toFile());
+                context.getStreamCachingStrategy().setSpoolEnabled(true);
+                context.getStreamCachingStrategy().setSpoolThreshold(1024);
+                
context.getStreamCachingStrategy().setRemoveSpoolDirectoryWhenStopping(false);
+                context.setStreamCaching(true);
+
+                Processor record = e -> 
e.getExchangeExtension().addOnCompletion(new SynchronizationAdapter() {
+                    @Override
+                    public void onComplete(Exchange exchange) {
+                        events.add("complete:" + 
exchange.getMessage().getBody(String.class));
+                    }
+
+                    @Override
+                    public void onFailure(Exchange exchange) {
+                        events.add("failure:" + 
exchange.getMessage().getBody(String.class));
+                    }
+                });
+
+                // the queues have no consumers, so the first message fills 
them
+                
from("direct:discard").process(record).to("seda:discard?size=1&discardWhenFull=true");
+                from("direct:offer").errorHandler(noErrorHandler())
+                        
.process(record).to("seda:offer?size=1&blockWhenFull=true&offerTimeout=10");
+                from("direct:add").errorHandler(noErrorHandler())
+                        .process(record).to("seda:add?size=1");
+
+                
from(fileUri("in?initialDelay=0&delay=10&sortBy=file:name")).routeId("files").autoStartup(false)
+                        .to("seda:files?size=1&discardWhenFull=true");
+                
from("seda:files?size=1&discardWhenFull=true").routeId("filesQueue").autoStartup(false)
+                        .convertBodyTo(String.class).to("mock:files");
+
+                
from("direct:spool").to("seda:spool?size=1&discardWhenFull=true");
+                
from("seda:spool?size=1&discardWhenFull=true").routeId("spoolQueue").autoStartup(false)
+                        .convertBodyTo(byte[].class).to("mock:spool");
+            }
+        };
+    }
+}
diff --git 
a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc 
b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
index dbcec9ca9fe4..c69525a6512d 100644
--- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
+++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
@@ -1037,6 +1037,19 @@ An onCompletion with `parallelProcessing` that 
synchronously stops its own route
 shutdown timeout occurs, and the route is then stopped forcibly. Stop the 
route asynchronously instead, for example
 from a separate thread or with the Control Bus `async=true` option.
 
+=== camel-seda, camel-disruptor - a message that is not queued completes its 
own exchange
+
+When the SEDA producer does not add an InOnly message to the queue, because it 
is discarded (`discardWhenFull=true`)
+or adding it fails (the queue is full, the `offerTimeout` elapses, or the 
producer is interrupted), the on
+completions of the exchange now run when the exchange is done. The same 
applies to the Disruptor producer when the
+ring buffer is full (with `blockWhenFull=false`) or the Disruptor is not 
started. Previously they had been handed over to the copy of the exchange
+that was dropped, and never ran: for example the file consumer neither 
committed nor rolled back the file, which
+then stayed in its in-progress repository and was not picked up again until 
the route was restarted.
+
+Now a discarded message is committed like any other message whose route 
completed (for example the file consumer
+moves or deletes the file), and a message that could not be added is rolled 
back, so the consumer can pick it up
+again.
+
 === Component deprecation
 
 ==== camel-minio

Reply via email to