[ 
https://issues.apache.org/jira/browse/CAMEL-25519?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Work on CAMEL-25519 started by shashank.
----------------------------------------
> camel-google-storage - consumer routes an object again while its first 
> exchange is still in flight (asynchronous routes)
> ------------------------------------------------------------------------------------------------------------------------
>
>                 Key: CAMEL-25519
>                 URL: https://issues.apache.org/jira/browse/CAMEL-25519
>             Project: Camel
>          Issue Type: Bug
>          Components: camel-google-storage
>            Reporter: shashank
>            Assignee: shashank
>            Priority: Major
>
> {{GoogleCloudStorageConsumer.processBatch}} hands each exchange to the route 
> with {{getAsyncProcessor().process(exchange, EmptyAsyncCallback.get())}} 
> (line 229 at main 5af71ab258b8) and returns without waiting; the object is 
> deleted ({{deleteAfterRead=true}}, the default) or moved only by the 
> exchange's on completion ({{processCommit}}, lines 212-225 and 240-271). When 
> the route continues asynchronously, for example {{to("seda:...")}} (seda 
> hands the on completion over to its consumer), an asynchronous producer such 
> as kafka, or {{threads()}}, the poll thread returns, and the next poll (500 
> ms later by default) lists the same object again, downloads it again and 
> routes it a second time. If the first exchange deletes the object after the 
> next poll listed it, the download of that poll ({{createExchange}}, line 337) 
> fails with a 404 {{StorageException}}: the whole poll ends with an exception 
> (logged as a WARN) and none of the objects it listed is routed by it.
> This is the defect fixed for aws2-s3 in CAMEL-17110 ("What the s3 consumer 
> lacks is an in-progress repository") and for minio in CAMEL-25512; the 
> google-storage consumer has no in-progress check. The same happens with 
> {{objectName}} set (one object per poll). With {{deleteAfterRead=false}} 
> every poll routes every object anyway, also while it is in flight.
> h3. Reproduction
> New {{GoogleCloudStorageConsumerInProgressTest}} with the module's in-memory 
> {{LocalStorageHelper}} storage: one object, route 
> {{google-storage://myCamelBucket?delay=10}} to {{seda:work}}; the seda route 
> holds the exchange until the consumer has completed a second poll (counted 
> with a {{pollStrategy}}). Two cases: listing the bucket, and 
> {{objectName=a.txt}}. On main both fail:
> {noformat}
> [a.txt was consumed again by a later poll while its first exchange was still 
> in flight]
> expected: 1
>  but was: 2
> {noformat}
> New {{GoogleCloudStorageConsumerDeletedObjectTest}}: three objects, the 
> storage deletes {{b.txt}} right after the listing (as another consumer 
> would); on main the poll fails ({{[the poll failed because b.txt was deleted 
> between the listing and the download] Expecting AtomicInteger(1) to have 
> value: 0}}) and {{a.txt}} and {{c.txt}} are only routed by a later poll.
> h3. Proposed fix
> The consumer keeps the names of the objects whose exchanges are in progress 
> (a set in the consumer, as the minio fix): a poll skips them, in both listing 
> and {{objectName}} mode; a name is added before the download and released 
> when its exchange completes or fails (in the on completion, after the 
> delete/move), when the poll does not process it (download failure, consumer 
> stopping). A listed object that no longer exists when its content is 
> downloaded (404) is skipped and the rest of the poll goes on, the same as 
> with {{objectName}} for a missing object; as the exchanges are created before 
> the batch is routed, {{CamelBatchSize}} only counts the objects that are 
> routed. No new option: aws2-s3 exposes a pluggable {{inProgressRepository}}; 
> like its default memory repository, the set is not cleared when the route is 
> stopped (names of exchanges still in flight are released when they complete), 
> and an exchange that never completes keeps its object reserved. 4.23 
> upgrade-guide note (in-flight objects are skipped, also with 
> {{deleteAfterRead=false}}; a deleted object no longer fails the poll). New 
> {{GoogleCloudStorageConsumerInProgressFailureTest}} checks that a failed 
> object is consumed again, and 
> {{GoogleCloudStorageConsumerInProgressTest.notProcessedObjectIsConsumedAgain}}
>  that an object whose exchange a stopping consumer did not process is not 
> left in progress (both pass on main too, which has no in-progress set; they 
> guard the releases on failure and on stop, and the second one fails when the 
> release on stop is removed). camel-google-storage unit tests: 58, 0 failures.
> Not covered: if the first exchange completes (deletes and releases the 
> object) after the next poll listed it but before that poll downloads it, the 
> download sees a 404; with the fix that object is skipped instead of failing 
> the poll.
> Found with a TLA+ model of the poll (list, download per listed object, 
> dispatch per exchange, stop of the batch) and asynchronous completion 
> (delete): "with deleteAfterRead every object is routed once" is violated in 
> 10 steps (poll, download, dispatch, poll again, download, dispatch again) and 
> "a poll does not fail because of an object another exchange of this consumer 
> deleted" in 10 steps; both hold with a synchronous route. With the fix both 
> hold, and no name stays in progress when nothing is in flight (also after a 
> stopped batch). Then confirmed with the real consumer as above.
> Affected: main, camel-4.22.x, camel-4.18.x, camel-4.14.x (same code).
> Duplicate check (2026-10-10): JIRA component camel-google-storage since 2024 
> (7 issues: download path containment, producer metadata, health check, 
> prefix; none about consumption), text "GoogleCloudStorageConsumer" (6, none 
> about in-flight objects), "google-storage" + "progress" (CAMEL-25512 is the 
> minio sibling); GitHub pull requests "google-storage consumer": only the 
> minio PR #27645 (merged); no open PR touches camel-google-storage.
> _Filed with Claude Code on behalf of allthingssecurity._



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

Reply via email to