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