Andrea Cosentino created CAMEL-25171:
----------------------------------------

             Summary: 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


{{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