allthingssecurity commented on code in PR #27347:
URL: https://github.com/apache/camel/pull/27347#discussion_r4181400235


##########
components/camel-couchbase/src/main/java/org/apache/camel/component/couchbase/CouchbaseConsumer.java:
##########
@@ -277,6 +270,31 @@ private int pollWithView() throws Exception {
         return processBatch(exchanges);
     }
 
+    /**
+     * Removes the document once its exchange has been processed successfully.
+     * <p/>
+     * Removing it while the exchanges are built, before any of them is handed 
to the route, loses the document when its
+     * exchange fails, or when it is never delivered at all because the batch 
is cut short by {@code maxMessagesPerPoll}
+     * or by the consumer stopping. The on-completion runs only for an 
exchange that went through the route and
+     * completed: a failed exchange keeps its document for the next poll, and 
an undelivered one is not touched.
+     */
+    private void removeDocumentOnCompletion(Exchange exchange, String id) {
+        exchange.getExchangeExtension().addOnCompletion(new 
SynchronizationAdapter() {

Review Comment:
   
   Good point. I kept handover allowed: with `allowHandover()` returning 
`false` the on-completion would run when the
   consumer's own unit of work ends, so with `seda` 
(`waitForTaskToComplete=Never`) or an async producer the document
   would be removed before the asynchronous part had processed it, which is the 
data loss this PR fixes.
   
   On aws2-s3, which the PR description cites as the pattern: it does not 
deliver the object twice. It also keeps
   handover, but it adds the key to its `inProgressRepository` in the poll, 
skips keys already there, and removes the key in
   `onComplete`/`onFailure`. The file consumer does the same.
   
   05b9c17f2451 adds the same guard here, for 
`consumerProcessedStrategy=delete`: the consumer keeps the IDs of the
   documents whose exchange is in flight and skips them in later polls. An ID 
is released when its exchange completes or
   fails, when it is never handed to the route (`maxMessagesPerPoll` or the 
consumer stopping), and when the poll fails
   part way. A skipped row does not count towards `maxMessagesPerPoll`. It is 
an internal concurrent set, not a new
   `inProgressRepository` option: it only has to cover this consumer's own 
exchanges, and unlike aws2-s3's default
   10000-entry FIFO cache it cannot evict an ID that is still in flight. 
`none`/`filter` are unchanged, as they re-read the
   document on every poll anyway.
   
   Tests hand the exchange to `seda:...?waitForTaskToComplete=Never` and 
complete it later: while it is in flight the next
   poll does not deliver the document (SQL++ and view); after completion it was 
processed and removed once; after failure
   the next poll delivers it again; ids of undelivered rows and of a failed 
poll are released. 6 fail without the guard,
   4 fail without the release; module suite 57 tests, 0 failures.
   
   Limit: the guard is per consumer instance. A consumer on another node 
polling the same bucket can still read the
   document before it is removed. The upgrade guide says so.
   
   _Claude Code on behalf of allthingssecurity_
   



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