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 ceb06d982014 CAMEL-25094: camel-seda - keep the spooled stream cache
of a message sent from a Multicast or Split (#26995)
ceb06d982014 is described below
commit ceb06d982014989160b559697910e67a725d777b
Author: allthingssecurity <[email protected]>
AuthorDate: Mon Sep 28 19:18:08 2026 +0530
CAMEL-25094: camel-seda - keep the spooled stream cache of a message sent
from a Multicast or Split (#26995)
Multicast and Split set CamelStreamCacheUnitOfWork on their sub-exchanges,
so that the stream caches of the sub-routes are released with the unit of
work of the parent exchange. When a sub-exchange sends to seda: or
disruptor: without waiting (InOnly), the producer hands a copy of it to
the queue and gives the copy its own reference to the stream cache
(CAMEL-20866). The copy kept the property, so that reference was also
registered on the parent's unit of work. When the parent was done, the
spool file was deleted while the copy was still queued or being routed,
and the consumer route failed with "Cannot reset stream from file": the
message was lost. SedaConsumer logs this at WARN; the disruptor consumer
does not log it.
SedaProducer and DisruptorProducer now remove the property from the copy
before copying the stream cache, as the Wire Tap EIP does (CAMEL-12108).
The copy then releases the file when it is done. The path that waits for
the reply is not changed, as the producer waits for the copy there.
Adds SedaStreamCachingSpoolTest (camel-core) and
DisruptorStreamCachingSpoolTest: sequential and parallel Multicast, and
Split, sending to the queue, with the consumer route waiting until the
parent exchange is done before it reads the body; the body must be read
in full and no spool file may be left behind. A direct send and a
Recipient List, which does not set the property, are controls.
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../component/disruptor/DisruptorProducer.java | 5 +
.../disruptor/DisruptorStreamCachingSpoolTest.java | 144 +++++++++++++++++++++
.../apache/camel/component/seda/SedaProducer.java | 5 +
.../component/seda/SedaStreamCachingSpoolTest.java | 141 ++++++++++++++++++++
4 files changed, 295 insertions(+)
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 51169bbdf1b1..63241b1bed7c 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
@@ -24,6 +24,7 @@ import java.util.concurrent.atomic.AtomicBoolean;
import com.lmax.disruptor.InsufficientCapacityException;
import org.apache.camel.AsyncCallback;
import org.apache.camel.Exchange;
+import org.apache.camel.ExchangePropertyKey;
import org.apache.camel.ExchangeTimedOutException;
import org.apache.camel.StreamCache;
import org.apache.camel.WaitForTaskToComplete;
@@ -237,6 +238,10 @@ public class DisruptorProducer extends
DefaultAsyncProducer {
// set a new from endpoint to be the disruptor
target.getExchangeExtension().setFromEndpoint(endpoint);
if (copy) {
+ // the copy is routed independently of the original exchange, so
any stream cache it holds must be
+ // released when the copy is done, and not with the unit of work
of a parent (multicast/split) exchange
+ // (same as the Wire Tap EIP does, see CAMEL-12108)
+
target.removeProperty(ExchangePropertyKey.STREAM_CACHE_UNIT_OF_WORK);
// if the body is stream caching based we need to make a deep copy
if (target.getMessage().getBody() instanceof StreamCache sc) {
StreamCache newBody = sc.copy(target);
diff --git
a/components/camel-disruptor/src/test/java/org/apache/camel/component/disruptor/DisruptorStreamCachingSpoolTest.java
b/components/camel-disruptor/src/test/java/org/apache/camel/component/disruptor/DisruptorStreamCachingSpoolTest.java
new file mode 100644
index 000000000000..5479f491e49e
--- /dev/null
+++
b/components/camel-disruptor/src/test/java/org/apache/camel/component/disruptor/DisruptorStreamCachingSpoolTest.java
@@ -0,0 +1,144 @@
+/*
+ * 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.io.BufferedInputStream;
+import java.io.ByteArrayInputStream;
+import java.io.InputStream;
+import java.nio.file.Path;
+import java.util.List;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+
+import org.apache.camel.Exchange;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+
+/**
+ * A message sent InOnly to a disruptor from inside a Multicast or Split must
keep its spooled stream cache file until
+ * the disruptor route is done with it, and must not leave the file behind
afterwards.
+ */
+class DisruptorStreamCachingSpoolTest extends CamelTestSupport {
+
+ private static final byte[] DATA = createData(16 * 1024);
+
+ @TempDir
+ Path spoolDir;
+
+ // the disruptor route waits until the parent exchange is done, so the
file would already be deleted if the
+ // copy did not hold its own reference to it
+ private final CountDownLatch parentDone = new CountDownLatch(1);
+
+ @Test
+ void testMulticastToDisruptor() throws Exception {
+ sendAndAssert("direct:multicast", stream(), 1);
+ }
+
+ @Test
+ void testParallelMulticastToDisruptor() throws Exception {
+ sendAndAssert("direct:multicastParallel", stream(), 1);
+ }
+
+ @Test
+ void testRecipientListToDisruptor() throws Exception {
+ // the recipient list does not tie the stream caches of its copies to
the parent, so this has always worked
+ sendAndAssert("direct:recipientList", stream(), 1);
+ }
+
+ @Test
+ void testSplitToDisruptor() throws Exception {
+ sendAndAssert("direct:split", List.of(stream(), stream()), 2);
+ }
+
+ @Test
+ void testDirectToDisruptor() throws Exception {
+ // without an EIP creating sub exchanges this has always worked
+ sendAndAssert("direct:start", stream(), 1);
+ }
+
+ private void sendAndAssert(String uri, Object body, int expected) throws
Exception {
+ MockEndpoint result = getMockEndpoint("mock:result");
+ result.expectedMessageCount(expected);
+
+ template.sendBody(uri, body);
+ parentDone.countDown();
+
+ MockEndpoint.assertIsSatisfied(context);
+ for (Exchange exchange : result.getReceivedExchanges()) {
+ assertArrayEquals(DATA,
exchange.getMessage().getBody(byte[].class));
+ }
+
+ // the copy deletes its spool file when the disruptor route is done
+ await().atMost(5, TimeUnit.SECONDS).untilAsserted(() -> {
+ String[] files = spoolDir.toFile().list();
+ assertNotNull(files);
+ assertEquals(0, files.length, "Spool files left behind: " +
List.of(files));
+ });
+ }
+
+ private static InputStream stream() {
+ // a stream that is not converted to an in-memory cache, so it is
spooled to disk
+ return new BufferedInputStream(new ByteArrayInputStream(DATA));
+ }
+
+ private static byte[] createData(int size) {
+ byte[] data = new byte[size];
+ for (int i = 0; i < size; i++) {
+ data[i] = (byte) ('a' + i % 26);
+ }
+ return data;
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+
context.getStreamCachingStrategy().setSpoolDirectory(spoolDir.toFile());
+ context.getStreamCachingStrategy().setSpoolEnabled(true);
+ context.getStreamCachingStrategy().setSpoolThreshold(1024);
+
context.getStreamCachingStrategy().setRemoveSpoolDirectoryWhenStopping(false);
+ context.setStreamCaching(true);
+
+ from("direct:start").to("disruptor:b");
+
+ from("direct:multicast").multicast().to("disruptor:b",
"mock:other");
+
+
from("direct:multicastParallel").multicast().parallelProcessing().to("disruptor:b",
"mock:other");
+
+
from("direct:recipientList").recipientList(constant("disruptor:b,mock:other"));
+
+ // each part is spooled when it enters the part route, which
runs within the split
+ from("direct:split").split(body()).to("direct:part");
+ from("direct:part").to("disruptor:b");
+
+ from("disruptor:b")
+ .process(e -> parentDone.await(20, TimeUnit.SECONDS))
+ .convertBodyTo(byte[].class)
+ .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 5037c554b650..55aa8dbc95b7 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
@@ -25,6 +25,7 @@ import java.util.concurrent.atomic.AtomicBoolean;
import org.apache.camel.AsyncCallback;
import org.apache.camel.Exchange;
+import org.apache.camel.ExchangePropertyKey;
import org.apache.camel.ExchangeTimedOutException;
import org.apache.camel.StreamCache;
import org.apache.camel.WaitForTaskToComplete;
@@ -248,6 +249,10 @@ public class SedaProducer extends DefaultAsyncProducer {
// handover the completion so its the copy which performs that, as we
do not wait
if (copy) {
target = prepareCopy(exchange, true);
+ // the copy is routed independently of the original exchange, so
any stream cache it holds must be
+ // released when the copy is done, and not with the unit of work
of a parent (multicast/split) exchange
+ // (same as the Wire Tap EIP does, see CAMEL-12108)
+
target.removeProperty(ExchangePropertyKey.STREAM_CACHE_UNIT_OF_WORK);
// if the body is stream caching based we need to make a deep copy
if (target.getMessage().getBody() instanceof StreamCache sc) {
StreamCache newBody = sc.copy(target);
diff --git
a/core/camel-core/src/test/java/org/apache/camel/component/seda/SedaStreamCachingSpoolTest.java
b/core/camel-core/src/test/java/org/apache/camel/component/seda/SedaStreamCachingSpoolTest.java
new file mode 100644
index 000000000000..b4d4d43e22bd
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/component/seda/SedaStreamCachingSpoolTest.java
@@ -0,0 +1,141 @@
+/*
+ * 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.io.InputStream;
+import java.util.List;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Exchange;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+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.assertNotNull;
+
+/**
+ * A message sent InOnly to SEDA from inside a Multicast or Split must keep
its spooled stream cache file until the SEDA
+ * route is done with it, and must not leave the file behind afterwards.
+ */
+public class SedaStreamCachingSpoolTest extends ContextTestSupport {
+
+ private static final byte[] DATA = createData(16 * 1024);
+
+ // the seda route waits until the parent exchange is done, so the file
would already be deleted if the
+ // copy did not hold its own reference to it
+ private final CountDownLatch parentDone = new CountDownLatch(1);
+
+ @Test
+ public void testMulticastToSeda() throws Exception {
+ sendAndAssert("direct:multicast", stream(), 1);
+ }
+
+ @Test
+ public void testParallelMulticastToSeda() throws Exception {
+ sendAndAssert("direct:multicastParallel", stream(), 1);
+ }
+
+ @Test
+ public void testRecipientListToSeda() throws Exception {
+ // the recipient list does not tie the stream caches of its copies to
the parent, so this has always worked
+ sendAndAssert("direct:recipientList", stream(), 1);
+ }
+
+ @Test
+ public void testSplitToSeda() throws Exception {
+ sendAndAssert("direct:split", List.of(stream(), stream()), 2);
+ }
+
+ @Test
+ public void testDirectToSeda() throws Exception {
+ // without an EIP creating sub exchanges this has always worked
+ sendAndAssert("direct:start", stream(), 1);
+ }
+
+ private void sendAndAssert(String uri, Object body, int expected) throws
Exception {
+ MockEndpoint result = getMockEndpoint("mock:result");
+ result.expectedMessageCount(expected);
+
+ template.sendBody(uri, body);
+ parentDone.countDown();
+
+ assertMockEndpointsSatisfied();
+ for (Exchange exchange : result.getReceivedExchanges()) {
+ assertArrayEquals(DATA,
exchange.getMessage().getBody(byte[].class));
+ }
+
+ // the copy deletes its spool file when the seda route is done
+ File spoolDir = testDirectory().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));
+ });
+ }
+
+ private static InputStream stream() {
+ // a stream that is not converted to an in-memory cache, so it is
spooled to disk
+ return new BufferedInputStream(new ByteArrayInputStream(DATA));
+ }
+
+ private static byte[] createData(int size) {
+ byte[] data = new byte[size];
+ for (int i = 0; i < size; i++) {
+ data[i] = (byte) ('a' + i % 26);
+ }
+ return data;
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+
context.getStreamCachingStrategy().setSpoolDirectory(testDirectory().toFile());
+ context.getStreamCachingStrategy().setSpoolEnabled(true);
+ context.getStreamCachingStrategy().setSpoolThreshold(1024);
+
context.getStreamCachingStrategy().setRemoveSpoolDirectoryWhenStopping(false);
+ context.setStreamCaching(true);
+
+ from("direct:start").to("seda:b");
+
+ from("direct:multicast").multicast().to("seda:b",
"mock:other");
+
+
from("direct:multicastParallel").multicast().parallelProcessing().to("seda:b",
"mock:other");
+
+
from("direct:recipientList").recipientList(constant("seda:b,mock:other"));
+
+ // each part is spooled when it enters the part route, which
runs within the split
+ from("direct:split").split(body()).to("direct:part");
+ from("direct:part").to("seda:b");
+
+ from("seda:b")
+ .process(e -> parentDone.await(20, TimeUnit.SECONDS))
+ .convertBodyTo(byte[].class)
+ .to("mock:result");
+ }
+ };
+ }
+}