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

Andrea Cosentino commented on CAMEL-25171:
------------------------------------------

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

Fix and test included; each fix in the PR was revert-checked individually.

_Claude Code on behalf of oscerd_

> camel-couchbase: the consumer discards route failures, and with 
> consumerProcessedStrategy=delete the document is already gone
> -----------------------------------------------------------------------------------------------------------------------------
>
>                 Key: CAMEL-25171
>                 URL: https://issues.apache.org/jira/browse/CAMEL-25171
>             Project: Camel
>          Issue Type: Bug
>          Components: camel-couchbase
>            Reporter: Andrea Cosentino
>            Assignee: Andrea Cosentino
>            Priority: Major
>
> {{CouchbaseConsumer.processBatch}} hands each exchange to the route and then 
> ignores what happened to it:
> {code:java}
> for (int index = 0; index < total && this.isBatchAllowed(); ++index) {
>     Exchange exchange = (Exchange) exchanges.poll();
>     exchange.setProperty(ExchangePropertyKey.BATCH_INDEX, index);
>     exchange.setProperty(ExchangePropertyKey.BATCH_SIZE, total);
>     exchange.setProperty(ExchangePropertyKey.BATCH_COMPLETE, index == total - 
> 1);
>     this.pendingExchanges = total - index - 1;
>     getProcessor().process(exchange);
> }
> {code}
> A failing route does not throw out of {{process()}} - the consumer processor 
> is async, so the failure is left on {{exchange.getException()}}. Nothing here 
> reads it. {{getExceptionHandler()}} does not appear anywhere in 
> {{CouchbaseConsumer}} at all, so {{bridgeErrorHandler}} is inert for this 
> component and a failed exchange produces no log line whatsoever.
> This is the same defect as CAMEL-25024 in camel-mongodb.
> h2. Why it is worse here: the document has already been deleted
> With {{consumerProcessedStrategy=delete}}, the removal happens in the polling 
> loop - in {{pollWithSqlQuery}} and {{pollWithView}}, while the exchanges are 
> being built - and therefore *before* {{processBatch}} runs:
> {code:java}
> if ("delete".equalsIgnoreCase(consumerProcessedStrategy)) {
>     CouchbaseCollectionOperation.removeDocument(collection, id, 
> endpoint.getWriteQueryTimeout(),
>             endpoint.getConsumerRetryPause());
> }
> {code}
> So when the route fails, the document is already gone from Couchbase, the 
> failure is swallowed, and nothing is logged. The message is lost silently.
> h2. The batch bookkeeping was copied without its safety net
> {{processBatch}} is a copy of {{GenericFileConsumer.processBatch}} with the 
> recovery machinery removed. The original keeps a {{notStarted}} queue, 
> decrements {{answer}} for exchanges that did not start, and drains *both* 
> {{notStarted}} and the leftover {{exchanges}} queue before returning. The 
> couchbase copy keeps the {{int answer = total;}} shape but never adjusts it, 
> and never drains the remainder - so exchanges still queued when 
> {{isBatchAllowed()}} flips to false at shutdown are neither processed nor 
> released. With a pooled exchange factory those are taken from the pool and 
> never returned.
> h2. Proposed fix
> Introduce the {{processExchange}} shape used in camel-mongodb: check 
> {{exchange.getException()}} after {{process()}}, report it through 
> {{getExceptionHandler()}}, and stop treating a failed exchange as delivered. 
> Release the exchanges left in the queue when the batch is cut short.



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

Reply via email to