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)

Reply via email to