oscerd opened a new pull request, #26779:
URL: https://github.com/apache/camel/pull/26779

   Three defects in the same two files, so one PR with a commit each.
   
   ## CAMEL-24933 — the consumer orphans a pooled exchange for every message
   
   `PulsarMessageListener.received()` does:
   
   ```java
   final Exchange exchange = PulsarMessageUtils.updateExchange(message, 
pulsarConsumer.createExchange(false));
   ```
   
   and `updateExchange` opened with `input.copy()`. 
`DefaultPooledExchange.newCopy()` returns
   `new DefaultExchange(this)`, so a copy is never pooled. With 
`camel.main.exchange-factory=pooled` that
   means, per consumed message:
   
   * the `DefaultPooledExchange` taken from the factory 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;`, throws
     `ClassCastException`, catches it, logs *"Error resetting exchange"* at 
DEBUG and counts a discard.
   
   So the pool never refills, pooling costs more than it saves here, and a cast 
failure is thrown and
   swallowed on every message. Without pooling it was "only" a second exchange 
allocation per message.
   
   The exchange is now populated in place. `updateExchange` has a single 
caller. I left the unused
   `updateExchangeWithException` alone — it is dead code, but removing a public 
method is an API change this
   fix does not need.
   
   ## CAMEL-24934 — a failed acknowledgement was reported without its cause
   
   ```java
   try {
       acknowledge(consumer, message);
   } catch (Exception e) {
       pulsarConsumer.getExceptionHandler().handleException("Error processing 
exchange", exchange,
               exchange.getException());
   }
   ```
   
   That catch sits in the `else` of `if (exchange.getException() != null)`, so 
`exchange.getException()` is
   `null` there by construction and the caught `e` was never used. A closed 
consumer, an unreachable broker
   or an acknowledgement timeout was logged with no cause and no stack trace, 
and the message came back
   later with nothing in the log to explain why. Now it passes `e`, with its 
own wording.
   
   ## CAMEL-16073 — negatively acknowledge a failed exchange
   
   Open since 2021. On a route failure the message was neither acked nor 
nacked, so it returned only when
   the acknowledgement timeout expired. It is now `negativeAcknowledge`d, 
unless `allowManualAcknowledgement`
   is enabled — in which case the route stays in charge, exactly as it already 
does for `acknowledge`.
   
   **Please read this part before approving.** It changes redelivery timing, 
and not in the direction one
   might assume. Disassembling `ConsumerImpl.negativeAcknowledge` in 
`pulsar-client` 4.2.4 shows it does:
   
   ```java
   negativeAcksTracker.add(messageId);
   unAckedMessageTracker.remove(discardBatch(messageId));
   ```
   
   i.e. a nack **removes the message from the ack-timeout tracker**. 
`camel-pulsar` defaults
   `ackTimeoutMillis` to 10 s (Pulsar's own default is 0, disabled) while 
`negativeAckRedeliveryDelayMicros`
   defaults to 60 s, so with the component defaults a failed message now comes 
back after **60 s instead of
   10 s**. The upgrade guide says so and gives the one-line remedy
   (`negativeAckRedeliveryDelayMicros=10000000`).
   
   The nack is the semantically correct signal — it is explicit, it feeds the 
nack metric and it honours
   `negativeAckRedeliveryBackoff`, which this component already exposes. But if 
you would rather have this
   behind an opt-in option, or have the default delay aligned, say so and I 
will change it.
   
   ## Tests
   
   * `PulsarMessageUtilsTest.testUpdateExchangeUpdatesTheGivenExchangeInPlace` 
— asserts the same instance
     comes back, plus headers and body.
   * `PulsarMessageListenerAcknowledgementTest` — a mocked `PulsarClient` 
(Mockito `RETURNS_SELF` for the
     fluent `ConsumerBuilder`) drives the real listener against a route that 
fails on one body and succeeds
     on another; asserts ack on success, nack on failure, and that an 
acknowledgement failure reaches the
     exception handler carrying its cause.
   
   Each fix was verified by reverting it and re-running: the three reverts 
produce exactly the three
   matching failures and nothing else. Module unit suite green (27 tests).
   
   _Claude Code on behalf of oscerd_
   
   🤖 Generated with [Claude Code](https://claude.com/claude-code)


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to