[ 
https://issues.apache.org/jira/browse/CAMEL-25297?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Claus Ibsen resolved CAMEL-25297.
---------------------------------
    Fix Version/s: 4.23.0
       Resolution: Fixed

Fixed on main via https://github.com/apache/camel/pull/27338

> camel-direct, camel-kamelet - while one exchange waits for the consumer of a 
> suspended or stopped route, the other exchanges of the same producer are sent 
> to that route
> ------------------------------------------------------------------------------------------------------------------------------------------------------------------------
>
>                 Key: CAMEL-25297
>                 URL: https://issues.apache.org/jira/browse/CAMEL-25297
>             Project: Camel
>          Issue Type: Bug
>          Components: camel-core, camel-kamelet
>            Reporter: shashank
>            Priority: Major
>             Fix For: 4.23.0
>
>
> {{DirectProducer}} caches the consumer of its endpoint in two fields that all 
> threads sending with the producer share (a route's {{to("direct:x")}} has one 
> producer):
> {code:java}
> if (consumer == null || stateCounter != component.getStateCounter()) {
>     stateCounter = component.getStateCounter();
>     consumer = component.getConsumer(key, block, timeout);
> }
> // send to consumer
> {code}
> {{DirectComponent}} increments {{stateCounter}} whenever a consumer is added 
> or removed, and a {{DirectConsumer}} removes itself when its route is 
> suspended ({{doSuspend}}) or stopped ({{doStop}}). With {{block=true}} (the 
> default) {{getConsumer}} then waits up to {{timeout}} (30 s by default) for 
> the consumer to come back.
> The first exchange after the route was suspended sets {{stateCounter}} to the 
> new value *before* it looks up the consumer, and then waits in 
> {{getConsumer}}. Every other exchange sent with the same producer in the 
> meantime sees {{stateCounter == component.getStateCounter()}} and the old, 
> non-null {{consumer}}, skips the lookup, and is processed by the consumer of 
> the suspended route. So suspending a direct route (route controller, JMX, a 
> route policy such as {{ThrottlingInflightRoutePolicy}}, which suspends the 
> consumer of a direct route so that its callers wait while too many exchanges 
> are inflight) does not stop concurrent senders: all but one keep going into 
> the suspended route, and only the first one waits as documented. The same 
> happens when the route is stopped: the exchanges go to the stopped route, 
> whose error handler fails them with a {{RejectedExecutionException}}, instead 
> of waiting for the consumer (and being processed once the route is started 
> again) or failing with {{DirectConsumerNotAvailableException}} after the 
> timeout. With {{block=false}} the window is short (the lookup does not wait), 
> but it is the same race; there two threads that both look up the consumer can 
> also write their results in the wrong order (a lookup that found the consumer 
> before it was removed writes it after a later lookup wrote {{null}}), and the 
> removed consumer then stays cached until the next add or remove.
> h3. Reproduction
> Test {{DirectProducerSuspendedConsumerTest}} (camel-core): 
> {{from("direct:start").to("direct:b?timeout=20000")}} and 
> {{from("direct:b").routeId("b").to("mock:b")}}. One exchange warms up the 
> producer, then route {{b}} is suspended. A first exchange is sent 
> asynchronously; a {{DirectComponent}} subclass counts down a latch in 
> {{getConsumer}}, so the test knows the first exchange is waiting for the 
> consumer. A second exchange is sent asynchronously. On main it is not 
> waiting: it reached the suspended route.
> {noformat}
> No exchange should reach the suspended route ==> expected: <1> but was: <2>
> {noformat}
> With the fix the second exchange also waits in {{getConsumer}}, and after 
> {{resumeRoute("b")}} both exchanges are delivered. The same test with 
> {{stopRoute("b")}} / {{startRoute("b")}}: on main the second exchange does 
> not wait, it fails at once with a {{RejectedExecutionException}} from the 
> stopped route's error handler:
> {noformat}
> The second exchange should wait for the consumer of route b ==> expected: 
> <false> but was: <true>
> {noformat}
> With the fix it waits and is delivered after {{startRoute("b")}}. No sleeps: 
> latches and Awaitility.
> The defect was found with a TLA+ model of the cache (two sending threads, 
> consumer removed/added, each Java statement one step). The property "an 
> exchange whose check starts after the consumer was removed is not sent to 
> that consumer" is violated: thread 1 sets the counter and starts the lookup, 
> the consumer is removed, thread 2's check passes with the old consumer. It 
> holds with one sending thread, and with the fix for 2 and 3 threads.
> h3. Proposed fix
> Keep the consumer and the counter in one immutable holder ({{record 
> CachedConsumer(DirectConsumer consumer, int stateCounter)}} in a {{volatile}} 
> field). The counter is read *before* the lookup and stored together with its 
> result, so another exchange either sees the old holder (counter differs, so 
> it looks up the consumer too, and waits) or the holder with the result of the 
> lookup. The cost on the fast path is the same as today (one volatile read of 
> the holder, one of the counter); a holder is allocated only when the consumer 
> changed.
> h3. camel-kamelet
> {{KameletProducer}} has the identical code ({{consumer}}/{{stateCounter}}, 
> {{KameletComponent.getConsumer}} with block/timeout, 
> {{KameletConsumer.doSuspend}}/{{doStop}} remove the consumer) and the same 
> defect: {{KameletProducerSuspendedConsumerTest}} (a route template 
> {{kamelet:source -> mock:kamelet}} used with 
> {{kamelet:echo/echo?timeout=20000}}, the kamelet route {{echo}} suspended or 
> stopped) fails on main the same way. The same holder fixes it.
> With the fix the new tests pass, and {{Direct*}}, {{*Suspend*}}, {{*Resume*}} 
> in camel-core (60 tests, 0 failures) and the {{Kamelet*Test}} tests of 
> camel-kamelet (84 tests, 3 skipped) pass.
> Affected: 4.14.x, 4.18.x and main (same code, checked with the GitHub 
> contents API).
> Duplicate check (2026-10-03): JIRA text "DirectProducer" with "stateCounter" 
> (none), "direct" with "suspend" since 2018 (15 issues, none about the 
> producer cache: CAMEL-14944 rest routes and controlbus, CAMEL-23494 
> preparingShutdown, ...), "DirectConsumerNotAvailableException" since 2020 
> (CAMEL-20867, CAMEL-15165, configuration issues). GitHub pull requests 
> "DirectProducer", "direct suspended": only #24698 (flaky 
> DirectProducerBlockingTest).
> _Filed with Claude Code on behalf of allthingssecurity._



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to