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

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

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

> camel-rocketmq - an InOut exchange is completed twice when its reply arrives 
> at the request timeout, or the reply is delivered twice
> ------------------------------------------------------------------------------------------------------------------------------------
>
>                 Key: CAMEL-25298
>                 URL: https://issues.apache.org/jira/browse/CAMEL-25298
>             Project: Camel
>          Issue Type: Bug
>          Components: camel-rocketmq
>            Reporter: shashank
>            Priority: Minor
>             Fix For: 4.23.0
>
>
> For InOut exchanges {{RocketMQProducer}} registers a reply handler in the 
> reply manager's timeout map, keyed by the generated message key. The reply 
> consumer hands each reply to 
> {{RocketMQReplyManagerSupport.handleReplyMessage}}:
> {code:java}
> ReplyHandler handler = timeoutMap.get(messageKey);
> if (handler != null) {
>     timeoutMap.remove(messageKey);
>     handler.onReply(messageKey, messageExt);   // -> processReply -> 
> callback.done(false)
> }
> {code}
> The lookup and the removal are two steps. In between:
> * the request timeout expires: the timeout map's purge task removes the entry 
> and calls {{onTimeout}}, which sets an {{ExchangeTimedOutException}} and 
> completes the exchange; then the reply consumer thread removes nothing, but 
> still calls {{onReply}}, which writes the reply into the exchange that was 
> already completed and completes it a second time;
> * a second copy of the reply is handled by another consumer thread (the reply 
> consumer is a {{MessageListenerConcurrently}} with several threads, and 
> RocketMQ delivers at least once): both threads get the handler and both 
> complete the exchange.
> * the reply manager stops: since CAMEL-24124 its timeout map drains the 
> pending handlers and calls {{onTimeout}} for each (rejecting the exchange); a 
> reply thread that got the handler just before completes the exchange a second 
> time.
> Completing an exchange twice continues the route twice after the producer 
> (the following steps, for example a database write, run twice), and mutates 
> the exchange while the first continuation uses it. camel-jms and camel-sjms 
> ({{QueueReplyManager}}, {{TemporaryQueueReplyManager}}) use 
> {{correlation.remove(id)}} once and only call the handler they removed, so 
> only one of the reply, a duplicate and the timeout can win.
> h3. Reproduction (unit test, no broker)
> {{RocketMQReplyManagerSupportTest}}: the reply manager with a timeout map 
> whose clock is controlled by the test, and without the reply topic consumer 
> (it needs a broker; the test calls {{handleReplyMessage}} itself). A hook in 
> the map's {{remove}} runs before the removal:
> * timeout: the hook lets the request time out (moves the clock, runs the 
> purge). On main the callback is called twice:
> {noformat}
> RocketMQReplyManagerSupportTest.testTimeoutWhileReplyIsHandled: The exchange 
> should be completed once ==> expected: <1> but was: <2>
> {noformat}
> * duplicate: the hook handles a second copy of the reply. On main:
> {noformat}
> RocketMQReplyManagerSupportTest.testDuplicateReply: The exchange should be 
> completed once ==> expected: <1> but was: <2>
> {noformat}
> * control: a single reply completes the exchange once with the reply body 
> (passes on main).
> The defect was found with a TLA+ model of the request-reply correlation 
> (producer send and send callback, responder, reply consumer threads, timeout 
> eviction; the map's get/remove/purge each atomic). The property "the exchange 
> is completed at most once" is violated in 7 steps (send, register, reply, 
> get, timeout, reply completes), and also without any timeout when the reply 
> is delivered twice. With the single removal it holds for up to 3 copies of 
> the reply, the timeout and a failing send (the other property of the model, 
> that a reply is never ignored while the exchange waits for it, still fails: 
> see "Not fixed here").
> h3. Proposed fix
> {{handleReplyMessage}}: {{ReplyHandler handler = 
> timeoutMap.remove(messageKey)}} and call {{onReply}} only when this call 
> removed the handler ({{DefaultTimeoutMap.remove}} and the purge take the same 
> lock, so exactly one of them gets the entry). {{cancelMessageKey}} also 
> removes in one step instead of {{get}} then {{remove}}.
> With the fix the 3 new tests pass. The module's other tests are ITs that need 
> Docker (not run). The module skips tests on aarch64 ({{skipTests.aarch64}}); 
> I ran them with {{-DskipTests.aarch64=false -DskipTests=false}}.
> h3. Not fixed here (same model)
> The model also shows that the producer registers the reply handler only in 
> the send callback ({{onSuccess}}), after the broker acknowledged the request. 
> A responder that answers before the client runs {{onSuccess}} (a busy 
> callback executor, a GC pause) gets its reply ignored ("Reply received for 
> unknown messageKey") and the exchange times out. camel-jms registers the 
> correlation before it sends. This needs the producer to be restructured (the 
> same method that CAMEL-25272 just changed), and I could not reproduce it 
> without a broker, so it is left for a follow-up.
> Affected: 4.14.x, 4.18.x and main (same code).
> Duplicate check (2026-10-03): JIRA text "rocketmq" since 2024 (6 issues: 
> CAMEL-25272, CAMEL-20368 batching consumers, docs and tracing issues), 
> "RocketMQReplyManagerSupport" (none). GitHub pull requests "rocketmq reply" 
> (#24742, CAMEL-24124 drain on shutdown), "RocketMQReplyManagerSupport" 
> (none); the open #22358 (getOut cleanup) touches {{processReceivedReply}} 
> only.
> _Filed with Claude Code on behalf of allthingssecurity._



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

Reply via email to