allthingssecurity opened a new pull request, #26995: URL: https://github.com/apache/camel/pull/26995
# Description [CAMEL-25094](https://issues.apache.org/jira/browse/CAMEL-25094) With stream caching spooled to disk, a message sent InOnly to `seda:` or `disruptor:` from inside a Multicast or Split loses its spool file when the parent exchange completes. The consumer route then fails with `Cannot reset stream from file` and the message is lost: ``` from("direct:start").multicast().to("seda:b", "mock:other"); from("seda:b").convertBodyTo(byte[].class).to("mock:result"); WARN SedaConsumer - Error processing exchange. ... StreamCacheException ... Cannot reset stream from file ... ``` Multicast and Split set `CamelStreamCacheUnitOfWork` on 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.addToQueue` and `DisruptorProducer.prepareCopy` give the queued copy its own reference to the stream cache with `sc.copy(target)` (CAMEL-20866), but the copy keeps the property. So `FileInputStreamCache.TempFileManager.addExchange` registers 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. `SedaConsumer` logs 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.addToQueue` and `DisruptorProducer.prepareCopy` remove `ExchangePropertyKey.STREAM_CACHE_UNIT_OF_WORK` from 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: one `sc.copy` adds one reference and one countdown. - Only the branch that copies the exchange is changed. The path that waits for the reply (InOut, `waitForTaskToComplete`) is left alone, because the producer waits for the copy there. camel-stub extends `SedaProducer` and gets the fix too. - It also stops stream caches that are created inside the seda or disruptor route from being registered on the parent's unit of work (the second symptom that CAMEL-12108 describes for the Wire Tap). I found this by reading the code; there is no test for it. The Recipient List is not affected: `RecipientListProcessor` creates 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 `discardWhenFull` drops, or that fails the `offerTimeout`, never completes, so its reference to the spool file is never released. That already happens for a plain `to("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: - New `SedaStreamCachingSpoolTest` (camel-core, where the seda tests live) and `DisruptorStreamCachingSpoolTest` (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 after `sendBody` returns, 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. - Without the change, the three Multicast/Split tests fail in both classes (`mock://result Received message count. Expected: <1> but was: <0>`, and `SedaConsumer` logs the `StreamCacheException` above). The controls pass on main. - With the change: camel-core `*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 - [x] I checked that the commit is targeting the correct branch (Camel 4 uses the `main` branch) # Tracking - [x] If this is a large change, bug fix, or code improvement, I checked there is a [JIRA issue](https://issues.apache.org/jira/browse/CAMEL) filed for the change (usually before you start working on it). # Apache Camel coding standards and style - [x] I checked that each commit in the pull request has a meaningful subject line and body. - [ ] I have run `mvn clean install -DskipTests` locally 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 - [x] If this PR includes AI-generated code, commits have proper co-authorship attribution (e.g., `Co-authored-by` trailers) 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-By` trailer. _Claude Code on behalf of allthingssecurity_ 🤖 Generated with [Claude Code](https://claude.com/claude-code) -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
