shashank created CAMEL-25128:
--------------------------------
Summary: 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
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)