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

Claus Ibsen updated CAMEL-25128:
--------------------------------
    Fix Version/s: 4.23.0

> camel-netty - with a TimeoutCorrelationManagerSupport correlation manager, a 
> failed write and the request timeout both complete the exchange, so it is 
> completed twice
> ----------------------------------------------------------------------------------------------------------------------------------------------------------------------
>
>                 Key: CAMEL-25128
>                 URL: https://issues.apache.org/jira/browse/CAMEL-25128
>             Project: Camel
>          Issue Type: Bug
>          Components: camel-netty
>            Reporter: shashank
>            Priority: Minor
>             Fix For: 4.23.0
>
>
> For concurrent request/reply the producer can use a {{correlationManager}} 
> that extends {{TimeoutCorrelationManagerSupport}}. 
> {{NettyProducer.processWithConnectedChannel}} registers the 
> {{NettyCamelState}} with {{putState}} (a {{TimeoutMap}} entry with the 
> timeout) and then writes the request. Two paths can then complete the 
> exchange, and they do not know about each other:
> * the write listener: when the write fails, it sets the cause on the exchange 
> and calls {{state.onExceptionCaughtOnce(false)}}, which calls 
> {{callbackDoneOnce}} (added by CAMEL-16718, 3.11.0). The correlation entry 
> stays in the map;
> * the timeout: the {{TimeoutMap}} evicts the entry and 
> {{TimeoutCorrelationManagerSupport.onEviction}} sets an 
> {{ExchangeTimedOutException}} and calls the raw {{callback.done(false)}} of 
> the state. It does not use {{callbackDoneOnce}}, so the "callback called" 
> flag of the state is not set.
> So:
> * *write fails first* (any write failure: an encoder that rejects the 
> message, such as a frame too long for a {{LengthFieldPrepender}}, a closed or 
> reset connection): the exchange completes with the write exception and 
> routing continues; {{timeout}} ms later (30 s by default) the eviction sets 
> {{ExchangeTimedOutException}} on the same exchange and calls the callback 
> again. Nothing removes the entry on a failed write.
> * *timeout first*: a large request that stays in the outbound buffer (the 
> server does not read) times out, the eviction completes the exchange; when 
> the connection is then reset, the pending write fails and completes the 
> exchange again, because the eviction did not set the flag that 
> {{callbackDoneOnce}} checks.
> The reply path is safe: {{getState}} removes the entry from the map, which is 
> atomic with the eviction.
> Effect on a route 
> {{from("direct:in").doTry().to("netty:tcp://...?sync=true&correlationManager=#corr").doCatch(Exception.class)...end()...}}:
>  the exchange completes once as expected ({{doCatch}}, the next step, 
> {{onComplete}}, {{ExchangeCompleted}}), and then again: the on-completions of 
> the exchange run a second time as {{onFailure}}, an {{ExchangeFailed}} event 
> follows the {{ExchangeCompleted}} event, and the inflight repository removes 
> the exchange twice, so its size is -1 after one such request. A negative 
> inflight count can make a graceful shutdown stop waiting for exchanges that 
> are really inflight (not tested). A consumer's own on-completions (commit, 
> rollback) are registered the same way (not tested).
> h3. Reproduction
> A real Netty server on localhost, the producer with {{sync=true}} and a 
> correlation manager extending {{TimeoutCorrelationManagerSupport}} 
> ({{timeout=300}}, {{timeoutChecker=50}}), called as an {{AsyncProcessor}}; 
> the callbacks are counted.
> * write fails first: codecs {{LengthFieldPrepender(2)}} and 
> {{StringEncoder}}, a 70003 character request. The callback is called after 1 
> to 35 ms with {{EncoderException}} (NettyClientTCPWorker thread) and again at 
> about 325 ms with {{ExchangeTimedOutException}} (NettyTimeoutWorkerPool 
> thread), 3 of 3 runs.
> * timeout first: the server accepts but does not read, a 16 MB textline 
> request; the timeout completes the exchange at about 318 ms, then the server 
> resets the connection and the write failure completes it again 1 to 7 ms 
> later with {{IOException}}, 3 of 3 runs.
> * route level (timeout first, the route above): {{doCatch}} with 
> {{ExchangeTimedOutException}}, the step after it once, {{onComplete}} and 
> {{ExchangeCompleted}}, then {{onFailure}} and {{ExchangeFailed}}, inflight 
> -1, 3 of 3 runs.
> * controls: a reply (1 callback, no exception), and a small request that is 
> read but not answered, then the connection closed (1 callback, 
> {{ExchangeTimedOutException}}).
> A TLA+ model of the write, the write listener, the eviction and the reply 
> checks that the callback is called at most once and the exchange is not 
> written after it completed. Both orders violate it on the current code; with 
> the fix below the properties hold and every exchange completes. The control 
> without write failures (reply against timeout only) holds on the current code.
> h3. Proposed fix
> Whoever completes the exchange claims it first, with the flag the state 
> already has:
> * {{NettyCamelState}}: a {{markDone()}} that does 
> {{callbackCalled.compareAndSet(false, true)}}; {{callbackDoneOnce}} uses it; 
> a new {{onExceptionCaughtOnce(doneSync, cause)}} sets the cause (or the 
> generic {{IOException}}) only after it has claimed the exchange. The existing 
> {{onExceptionCaughtOnce(doneSync)}} is kept.
> * {{NettyProducer}}: pass the cause of the failed write to 
> {{onExceptionCaughtOnce}} instead of setting it on the exchange before.
> * {{TimeoutCorrelationManagerSupport.onEviction}}: only if 
> {{value.markDone()}} set the timeout exception or body and call the callback.
> The correlation entry of a failed write still goes at the timeout, as today; 
> the claim makes that a no-op. Only managers extending 
> {{TimeoutCorrelationManagerSupport}} (the documented recommendation for 
> custom correlation managers) were affected; 
> {{DefaultNettyCamelStateCorrelationManager}} completes everything on the 
> channel's event loop through {{callbackDoneOnce}} and behaves the same.
> Tests: both orders as unit tests with a {{TimeoutCorrelationManagerSupport}} 
> subclass and an outbound handler that fails the write, either at once or when 
> the test fails the pending write after the timeout was processed (a worker 
> pool that counts the processed timeouts, no sleeps).
> Affected: 3.11.0 and later (CAMEL-16718 added the write-failure completion; 
> before that a failed write was left to the timeout). The same code is at the 
> camel-4.0.0 and 4.18.0 tags.
> Duplicate check (2026-09-29): JIRA text "TimeoutCorrelationManagerSupport" 
> (CAMEL-18533, thread pools), "correlationManager" with "netty" (none), 
> component camel-netty since June 2025 (13 issues: SSL, converters, codecs), 
> "netty" with "callback" and "twice" (CAMEL-7500, 2014, retry). GitHub pull 
> requests "TimeoutCorrelationManagerSupport", "netty correlation", "netty 
> producer callback": #2794 (2019, TimeoutMap listener) only.
> _Filed with Claude Code on behalf of allthingssecurity._



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

Reply via email to