CAMEL-25094: camel-seda - keep the spooled stream cache of a message sent from a Multicast or Split - #26995
Conversation
…sent from a Multicast or Split 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 <noreply@anthropic.com>
oscerd
left a comment
There was a problem hiding this comment.
Correct, and it reuses the right precedent. A message sent InOnly to seda:/disruptor: from a Multicast/Split is routed on its own unit of work asynchronously, but its spooled stream cache was still pinned to the parent's UoW via STREAM_CACHE_UNIT_OF_WORK, so the parent completing deleted the spool file before the consumer read it (Cannot reset stream from file). Removing that property from the copy detaches the cache from the parent's lifecycle, so it's released with the copy's own UoW when the consumer route finishes instead — exactly what the Wire Tap EIP does (CAMEL-12108), and the comment cites it. The cache is still cleaned up (by the consumer's UoW), so there's no leak, just a correct hand-off.
Applying the same one-line removal in both SedaProducer and DisruptorProducer keeps the two InOnly async producers consistent, and the SedaStreamCachingSpoolTest/DisruptorStreamCachingSpoolTest pin that the spool file survives and is readable downstream. LGTM (pending CI — fork PR).
This review was generated with AI assistance and reviewed/issued by the human operator. Claude Code on behalf of oscerd
gnodet-bot
left a comment
There was a problem hiding this comment.
Solid fix. The change is minimal, correctly scoped, and exactly mirrors the established WireTapProcessor pattern (CAMEL-12108).
What this fixes: When a message is sent InOnly to seda: or disruptor: from inside a Multicast or Split, the queued copy inherits STREAM_CACHE_UNIT_OF_WORK from the parent. This causes sc.copy(target) to register the copy's stream cache reference on the parent's unit of work. When the parent completes (immediately after the InOnly send), it deletes the spool file — before the consumer reads it.
Why it's correct: Removing STREAM_CACHE_UNIT_OF_WORK before sc.copy(target) is essential — copy() calls TempFileManager.addExchange() which reads that property to decide which UoW owns the reference. The placement in SedaProducer.addToQueue() (vs prepareCopy()) is intentional: prepareCopy() is protected and used by the InOut path too, while the property removal only applies to the fire-and-forget copy. In DisruptorProducer, prepareCopy() is private with a if (copy) guard, so the placement there is equally safe.
Tests are well-designed: the CountDownLatch pattern deterministically reproduces the race (parent completes before consumer reads), and the control tests (direct, recipientList) confirm no regressions. The Awaitility assertion on spool file cleanup verifies no resource leaks.
This review was generated by an AI agent, Hermès on behalf of @gnodet.
davsclaus
left a comment
There was a problem hiding this comment.
Thanks. Removing CamelStreamCacheUnitOfWork from the independent queued copy before sc.copy(target) is the same fix WireTapProcessor has had since CAMEL-12108, and it's correctly limited to the copy path, so InOut and wait-for-reply sends are unchanged. The spool tests with direct and recipient list as controls are good. LGTM.
The discardWhenFull/offerTimeout case you mention, where a dropped copy never releases its spool file, is worth a follow-up JIRA.
Claude Code on behalf of davsclaus
|
🌟 Thank you for your contribution to the Apache Camel project! 🌟 🐫 Apache Camel Committers, please review the following items:
|
|
🧪 CI tested the following changed modules:
🔬 Scalpel shadow comparison — Scalpel: 480 of 695 tested, 28 compile-only — current: 484 all testedMaveniverse Scalpel detected 480 affected modules (current approach: 484). Skip-tests mode would test 480 modules (3 direct + 477 downstream), skip tests for 28 (generated code, meta-modules) Modules only in current approach (4)
Modules Scalpel would test (480)
Modules with tests skipped (28)
Build reactor — dependencies compiled but only changed modules were tested (3 modules, 6.5s total)Total reactor time: 6.5s
Top 20 slowest modules:
|
oscerd
left a comment
There was a problem hiding this comment.
CI is green on both build (17) and build (25), matching what I traced earlier: the async InOnly copy detaches the spooled stream cache from the parent unit of work with target.removeProperty(ExchangePropertyKey.STREAM_CACHE_UNIT_OF_WORK), so a stream cached to disk by a Multicast or Split survives until the seda/disruptor consumer route processes it, instead of being closed with the parent exchange (the "Cannot reset stream from file" case). Both SedaProducer and DisruptorProducer carry the identical change, mirroring the Wire Tap EIP fix (CAMEL-12108).
Approving.
This review was generated with AI assistance and reviewed/issued by the human operator. Claude Code on behalf of oscerd
Description
CAMEL-25094
With stream caching spooled to disk, a message sent InOnly to
seda:ordisruptor:from inside a Multicast or Split loses its spool file when the parent exchange completes. The consumer route then fails withCannot reset stream from fileand the message is lost:Multicast and Split set
CamelStreamCacheUnitOfWorkon their sub-exchanges, so that the stream caches of the sub-routes are released with the parent's unit of work (MulticastProcessor,Splitter). For an InOnly send,SedaProducer.addToQueueandDisruptorProducer.prepareCopygive the queued copy its own reference to the stream cache withsc.copy(target)(CAMEL-20866), but the copy keeps the property. SoFileInputStreamCache.TempFileManager.addExchangeregisters the copy's reference on the parent's unit of work too, and the parent deletes the file as soon as it is done, which is right after the InOnly send returned.SedaConsumerlogs the failure at WARN. The disruptor consumer does not log it at all (it processes the exchange with a no-op callback), so there the message disappears silently.The Wire Tap EIP had the same problem and removes the property from its copy (CAMEL-12108,
WireTapProcessor). In CAMEL-20866, SEDA (InOnly), Disruptor (InOnly) and Wire Tap were named as the cases where the copy is executed independently, but only the Wire Tap drops the property.This change:
SedaProducer.addToQueueandDisruptorProducer.prepareCopyremoveExchangePropertyKey.STREAM_CACHE_UNIT_OF_WORKfrom the copy before copying the stream cache, as the Wire Tap does. The copy's reference is then registered on the copy, and the file is deleted when the copy is done. There is no double release: onesc.copyadds one reference and one countdown.waitForTaskToComplete) is left alone, because the producer waits for the copy there. camel-stub extendsSedaProducerand gets the fix too.The Recipient List is not affected:
RecipientListProcessorcreates its own sub-exchanges and does not set the property. The new tests include it as a control. A Recipient List inside a Split still inherits the property from the split's sub-exchange, and that case is covered by this fix.One caveat, which exists already and is not changed here: a copy that
discardWhenFulldrops, or that fails theofferTimeout, never completes, so its reference to the spool file is never released. That already happens for a plainto("seda:x?discardWhenFull=true")and leaves the file until the spool directory is removed when the context stops. With this change the Multicast and Split cases behave the same way, instead of losing the message. The dropped copy also takes the original's handed-over on completions with it, which is a separate issue and is not changed here.Tests:
SedaStreamCachingSpoolTest(camel-core, where the seda tests live) andDisruptorStreamCachingSpoolTest(camel-disruptor). Spooling is enabled with a 1 KB threshold into a test directory, and the body is a 16 KB stream. Routes: a sequential Multicast, a parallel Multicast and a Split of two streams, each sending InOnly to the queue. The consumer route first waits on a latch that the test releases aftersendBodyreturns, so the parent exchange is done before the body is read, and the test is deterministic. Each test asserts that the full body is read in the consumer route, and then that the spool directory is empty (Awaitility). A direct send and a Recipient List are controls.mock://result Received message count. Expected: <1> but was: <0>, andSedaConsumerlogs theStreamCacheExceptionabove). The controls pass on main.*StreamCach*,*Multicast*,*Split*,*WireTap*plus all the seda tests (org/apache/camel/component/seda/**): 662 tests, 0 failures, 0 errors, 6 skipped. camel-disruptor, all tests: 111 tests, 0 failures, 0 errors.This merges cleanly with the open #26793 (
SedaProducer, other hunks) and #26980 (DisruptorProducer, other hunks).Found with a TLA+ model of the reference counter (parent, sub-exchange, queued copy, consumer), which finds the trace where the parent is done before the consumer reads, and holds for the fixed model, including that the file is always deleted in the end. I then reproduced the failure with the real classes, and checked that the 4.6.0
SedaProducer, from before the deep copy, fails the same way, so this is not a regression of CAMEL-20866.Target
mainbranch)Tracking
Apache Camel coding standards and style
mvn clean install -DskipTestslocally from root folder and I have committed all auto-generated changes.(I built and tested the affected modules, including the formatter and import-sort plugins. I did not run the full root build.)
AI-assisted contributions
Co-authored-bytrailers) and the PR description identifies the AI tool used.This PR was prepared with Claude Code (Claude Opus 5.5). The commit carries a
Co-Authored-Bytrailer.Claude Code on behalf of allthingssecurity
🤖 Generated with Claude Code