CAMEL-25082: camel-disruptor - mark the published copy, not the caller's exchange, as ignored on a request/reply timeout - #26980
Conversation
…r's exchange, as ignored on a request/reply timeout Cause: on a request/reply timeout DisruptorProducer set the disruptor.ignoreExchange property on the caller's exchange instead of the copy it had published into the ring buffer. The consumer ignored any exchange that had the property (containsKey, whatever the value). Effect: a timed out copy that the consumer had not started yet was still processed, and every later copy of the caller's exchange inherited the property and was dropped by the consumer: a redelivery after the timeout timed out again, and an InOnly fallback to another disruptor endpoint was lost silently. Fix: the producer puts its completed flag (the AtomicBoolean that the reply, the timeout and the interrupt already claim) on the published copy before publishing it, so the consumer ignores the copy only when the producer no longer waits for it. The flag is not copied into the routed exchange, not copied back into the caller's exchange with the reply, and not inherited by a new copy. The consumer and the reconfiguration buffer check the value of the flag. camel-seda is not affected, as SedaProducer removes the timed out copy from its queue. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Signed-off-by: Claus Ibsen <claus.ibsen@gmail.com>
|
🌟 Thank you for your contribution to the Apache Camel project! 🌟 🐫 Apache Camel Committers, please review the following items:
|
gnodet-bot
left a comment
There was a problem hiding this comment.
Solid fix. The core bug was that the timeout path marked the caller's exchange instead of the published copy — so the consumer never saw the flag on the timed-out copy, and every subsequent copy of the caller (redelivery, fallback) inherited the mark and was silently dropped.
The fix correctly reuses the completed AtomicBoolean (already shared between timeout/reply/interrupt paths since CAMEL-25038) as the ignore flag on the published copy. The happens-before relationship through AtomicBoolean.get()/compareAndSet() is correct for the producer→consumer handoff. All leak paths are cleaned:
prepareCopy()strips it from new copies (no inheritance)prepareExchange()strips it from the routed copy (no leak into the route)onDonestrips it aftercopyResults(no leak back to caller)
isIgnoreExchange() defensively handles both the new AtomicBoolean and the legacy Boolean.TRUE value — good.
Tests cover all four scenarios from the description (redelivery, InOnly fallback, timed-out-before-consumer, no-leak-on-success). Upgrade guide paragraph is clear.
No issues found.
This review was generated by an AI agent, Hermès on behalf of @gnodet.
|
🧪 CI tested the following changed modules:
🔬 Scalpel shadow comparison — Scalpel: 9 of 695 tested, 26 compile-only — current: 9 all testedMaveniverse Scalpel detected 9 affected modules (current approach: 9). Skip-tests mode would test 9 modules (2 direct + 8 downstream), skip tests for 26 (generated code, meta-modules) Modules Scalpel would test (9)
Modules with tests skipped (26)
All tested modules (36 modules, 5m 12s total)Total reactor time: 5m 12s
Top 20 slowest modules:
|
Description
CAMEL-25082, found in the review of CAMEL-25038 (#26913).
On a request/reply timeout
DisruptorProducersetDisruptorEndpoint.DISRUPTOR_IGNORE_EXCHANGEon the caller's exchange, not on the copy it had published into the ring buffer.DisruptorConsumerignores any exchange that has the property (containsKey, whatever its value). So:prepareCopycopies the properties) and was dropped by the consumer, with only a TRACE log:doCatch(ExchangeTimedOutException.class).to("disruptor:fallback"), was dropped: an InOnly message was lost silently, an InOut one timed out again.Change
completedflag (theAtomicBooleanthat the reply, the timeout and the interrupt already claim since CAMEL-25038) on the published copy before it publishes it. The consumer ignores the copy only when the producer no longer waits for it, which now also covers an interrupted wait. The property is set before publishing, so the property map is not changed while the consumer reads it; only theAtomicBooleanchanges afterwards.copyResultscopies the properties), and is not inherited by a new copy of the caller's exchange.DisruptorConsumerand the reconfiguration buffer inDisruptorReferencecheck the value of the flag throughDisruptorEndpoint.isIgnoreExchange.Compared with camel-seda
camel-seda is not affected: on a timeout (and on an interrupt)
SedaProducerremoves the timed out copy from its queue (endpoint.getQueue().remove(copy)) and does not mark the caller's exchange.Tests
New
DisruptorTimeoutIgnoreExchangeTest:Without the main-code change the first three fail (
CamelExecutionExceptionfrom the timed out redeliveries,mock://fallback ... Expected: <1> but was: <0>,mock://busy ... Expected: <1> but was: <2>). With it, the whole camel-disruptor suite passes (110 tests).Claude Code on behalf of davsclaus
🤖 Generated with Claude Code