[
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)