allthingssecurity opened a new pull request, #26996: URL: https://github.com/apache/camel/pull/26996
# Description [CAMEL-25095](https://issues.apache.org/jira/browse/CAMEL-25095) `JmsProducer.processInOut` registers the reply handler in the correlation map inside `MessageCreator.createMessage`, before `MessageProducer.send()` runs. From then on, two other threads can complete the exchange: the request timeout (the timeout checker evicts the handler and `onTimeout` runs `processReply`, which calls `done(false)`), and the reply listener, when the request already reached the broker and was answered. CAMEL-24073 (#24730, 4.22.0) added `replyManager.cancelCorrelationId(...)` to the catch block around the send, so a failed send no longer leaves the handler behind for the timeout to fire later. The catch block still rethrows in every case, though, and `process()` then sets the send exception and calls `callback.done(true)`. When the timeout or the reply removed the handler while `send()` was still running, the exchange had already been completed, and the send failure completes it a second time: ``` from("direct:start") .doTry().to("jms:queue:req?requestTimeout=500") .doCatch(Exception.class)... send() to req blocks 1500 ms, then throws JMSException: ~550 ms timeout thread: catch block runs, ExchangeCompleted, on-completions run as onComplete ~1520 ms caller thread: ExchangeFailed, the same on-completions run as onFailure, the caller gets UncategorizedJmsException although the route handled the timeout inflight repository size after 3 such requests: -1, -2, -3 ``` A send that outlives `requestTimeout` and then fails is what a broker that stops answering produces: the Artemis client gives up after `callTimeout` (30 s by default), and the default `requestTimeout` of camel-jms is 20 s. In the second variant, the request is sent and answered, and `send()` throws afterwards (a lost send acknowledgement): the reply and the send failure both complete the exchange. This change: - `ReplyManager.cancelCorrelationId` returns `true` when it removed a pending handler, and `false` otherwise. The removal from the correlation map is already atomic in the three paths (timeout eviction, reply, cancel), so whoever removes the handler owns the completion of the exchange. - The catch block in `processInOut`: when a handler was registered and the cancel returns `false`, the timeout or the reply completes the exchange (`onTimeout` may still be queued on the timeout thread pool). The send failure is logged at WARN and `processInOut` returns `false` without touching the exchange. Otherwise it rethrows as before, so the CAMEL-24073 case (the send fails before the timeout) is unchanged, and so is a failure before the handler was registered. - When the reply manager stops, its correlation map evicts all pending handlers and completes their exchanges with a `RejectedExecutionException`, so an exchange whose cancel returned `false` is always completed by the other thread. (The only exception exists already: an `onTimeout` task rejected by a thread pool that is shutting down.) - Upgrade guide 4.23: the exchange keeps the outcome of the timeout or the reply, and `cancelCorrelationId` (added to the public `ReplyManager` interface in 4.22.0) now returns `boolean`. A custom `ReplyManager` has to change the return type. Tests, in `JmsInOutSendFailureCallbackTest` (the CAMEL-24073 test). The proxied `MessageProducer` now picks its behaviour by queue, and it uses a latch instead of sleeps: the send waits until an on-completion on the exchange has run. - `testCallbackInvokedOnceWhenTimeoutFiresDuringFailingSend`: `requestTimeout=100`, `requestTimeoutCheckerInterval=50`. The send waits until the timeout has completed the exchange, then throws. Asserts `ExchangeTimedOutException` on the result, exactly one `onFailure` and no `onComplete`, and an inflight count of 0. - `testCallbackInvokedOnceWhenHandledTimeoutFiresDuringFailingSend`: the same inside `doTry`/`doCatch(ExchangeTimedOutException.class)`. Asserts no exception, the body set by the catch block, exactly one `onComplete` and no `onFailure`, and an inflight count of 0. - `testCallbackInvokedOnceWhenReplyArrivesDuringFailingSend`: the send really sends the request, waits until the reply (from a route consuming the request queue) has completed the exchange, then throws. Asserts the reply body, no exception, exactly one `onComplete`, and an inflight count of 0. - `testCallbackInvokedOnceOnSendFailure` (CAMEL-24073) is unchanged. Without the change in `JmsProducer` and `ReplyManagerSupport`, the three new tests fail (for example `expected: <org.apache.camel.ExchangeTimedOutException> but was: <org.springframework.jms.UncategorizedJmsException>`, and for the handled timeout and the reply: `expected: <null> but was: <...UncategorizedJmsException...>`), and the CAMEL-24073 test passes. With it, all four pass, and the send failure is logged at WARN once per test. The camel-jms tests `*InOut*Test,*ReplyTo*Test,*RequestReply*Test,*Timeout*Test` pass: 111 tests (103 + 6 serial + 2 exclusive), 0 failures, 1 skipped. The 11 integration tests with the same names (`*IT`) pass as well. (`JmsDeadLetterChannelInOutIT` times out when it runs in the same parallel surefire execution as the unit tests, on `main` too; it passes on its own and with the other ITs.) Found with a TLA+ model of the producer thread, the timeout checker, the reply listener and the broker, which checks that the callback is called at most once. It is violated by the two orders above, it reproduces CAMEL-24073 on the code before that fix, and the fixed model holds it and completes every exchange. I then reproduced both orders against the real classes with an embedded Artemis broker. Not in this PR: camel-sjms (`SjmsProducer.processInOut`) looks like it still has the pattern from before CAMEL-24073 (the reply is registered inside the message creator, and a send failure completes the exchange with the handler still registered). I only read that code and did not test it. # 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 camel-jms and the modules it depends on, including the formatter and import-sort plugins, and ran its request/reply tests. 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]
