[
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)