[ 
https://issues.apache.org/jira/browse/CAMEL-24933?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18118339#comment-18118339
 ] 

Andrea Cosentino commented on CAMEL-24933:
------------------------------------------

PR opened: https://github.com/apache/camel/pull/26779

Three defects in the same two files, one commit each. Each was verified by 
reverting it and re-running - the three reverts produce exactly the three 
matching test failures.

For CAMEL-16073 note the redelivery-timing change called out in the PR and the 
upgrade guide: a negative acknowledgement removes the message from the client 
unacked tracker, so redelivery follows negativeAckRedeliveryDelayMicros (60s) 
instead of ackTimeoutMillis (10s in this component).

----
_Claude Code on behalf of oscerd (Andrea Cosentino)._

> camel-pulsar - the consumer leaks a pooled exchange for every message
> ---------------------------------------------------------------------
>
>                 Key: CAMEL-24933
>                 URL: https://issues.apache.org/jira/browse/CAMEL-24933
>             Project: Camel
>          Issue Type: Bug
>          Components: camel-pulsar
>            Reporter: Andrea Cosentino
>            Assignee: Andrea Cosentino
>            Priority: Major
>
> h3. Summary
> The consumer creates an {{Exchange}} through the pooled exchange factory and 
> then throws it away for a
> copy. With {{camel.main.exchange-factory=pooled}} every consumed message 
> leaks one pooled exchange and
> makes {{PooledExchangeFactory.release}} throw a {{ClassCastException}} that 
> is swallowed at DEBUG.
> h3. Details
> {{PulsarMessageListener.received()}}:
> {code:java}
> final Exchange exchange = PulsarMessageUtils.updateExchange(message, 
> pulsarConsumer.createExchange(false));
> {code}
> and {{PulsarMessageUtils.updateExchange}} starts with:
> {code:java}
> final Exchange output = input.copy();
> {code}
> {{DefaultPooledExchange.newCopy()}} returns {{new DefaultExchange(this)}}, 
> i.e. a copy is never pooled.
> So with pooling enabled:
> # {{createExchange(false)}} takes a {{DefaultPooledExchange}} out of the 
> factory;
> # {{copy()}} discards it in favour of a plain {{DefaultExchange}} - the 
> pooled one is never {{done()}}
>   and never returned to the pool;
> # the callback calls {{releaseExchange(copy, false)}}, and since the copy is 
> not a {{PooledExchange}}
>   the {{pooledExchange.done()}} branch in {{DefaultConsumer.releaseExchange}} 
> is skipped;
> # {{PooledExchangeFactory.release}} opens with {{PooledExchange ee = 
> (PooledExchange) exchange;}}, so it
>   throws {{ClassCastException}}, catches it, logs {{"Error resetting 
> exchange: ..."}} at DEBUG and
>   increments the {{discarded}} statistic.
> The result is one orphaned pooled exchange and one thrown-and-swallowed 
> {{ClassCastException}} per
> consumed message. The pool never refills, so pooling brings this component no 
> benefit and some cost.
> Without pooling the copy is "only" a second exchange allocation per message.
> h3. Proposed fix
> Populate the exchange the consumer created instead of copying it. 
> {{updateExchange}} has exactly one
> caller, so the change is contained. Note also that {{output.setIn(msg)}} in 
> that method sets back the
> message it just read, and that the sibling {{updateExchangeWithException}} is 
> dead code - nothing calls
> it.
> ----
> _Reported by Claude Code on behalf of oscerd (Andrea Cosentino)._



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

Reply via email to