This is an automated email from the ASF dual-hosted git repository.

davsclaus pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git


The following commit(s) were added to refs/heads/main by this push:
     new d8be2daaa324 CAMEL-24195: camel-aws2-s3 - Fix StreamUploadProducer 
race condition and timeout task exception handling
d8be2daaa324 is described below

commit d8be2daaa32416f2941282f5fc9aa9dad4598cfe
Author: Claus Ibsen <[email protected]>
AuthorDate: Fri Jul 17 22:42:12 2026 +0200

    CAMEL-24195: camel-aws2-s3 - Fix StreamUploadProducer race condition and 
timeout task exception handling
    
    Co-Authored-By: Claude Opus 4.6 <[email protected]>
---
 .../aws2/s3/stream/AWS2S3StreamUploadProducer.java | 29 ++++++++++++++++------
 1 file changed, 21 insertions(+), 8 deletions(-)

diff --git 
a/components/camel-aws/camel-aws2-s3/src/main/java/org/apache/camel/component/aws2/s3/stream/AWS2S3StreamUploadProducer.java
 
b/components/camel-aws/camel-aws2-s3/src/main/java/org/apache/camel/component/aws2/s3/stream/AWS2S3StreamUploadProducer.java
index b75ba975c060..0ed57fd2e99c 100644
--- 
a/components/camel-aws/camel-aws2-s3/src/main/java/org/apache/camel/component/aws2/s3/stream/AWS2S3StreamUploadProducer.java
+++ 
b/components/camel-aws/camel-aws2-s3/src/main/java/org/apache/camel/component/aws2/s3/stream/AWS2S3StreamUploadProducer.java
@@ -137,8 +137,12 @@ public class AWS2S3StreamUploadProducer extends 
DefaultProducer {
             lock.lock();
             try {
                 if (ObjectHelper.isNotEmpty(uploadAggregate)) {
-                    uploadPart(uploadAggregate);
-                    completeUpload(uploadAggregate);
+                    try {
+                        uploadPart(uploadAggregate);
+                        completeUpload(uploadAggregate);
+                    } catch (Exception e) {
+                        LOG.warn("Error during timeout-triggered upload flush 
for {}", uploadAggregate.dynamicKeyName, e);
+                    }
                     uploadAggregate = null;
                 }
 
@@ -148,8 +152,12 @@ public class AWS2S3StreamUploadProducer extends 
DefaultProducer {
                     for (Map.Entry<Long, UploadState> entry : 
timestampBasedUploads.entrySet()) {
                         UploadState state = entry.getValue();
                         if (ObjectHelper.isNotEmpty(state) && 
state.buffer.size() > 0) {
-                            uploadPart(state);
-                            completeUpload(state);
+                            try {
+                                uploadPart(state);
+                                completeUpload(state);
+                            } catch (Exception e) {
+                                LOG.warn("Error during timeout-triggered 
upload flush for {}", state.dynamicKeyName, e);
+                            }
                             keysToRemove.add(entry.getKey());
                         }
                     }
@@ -178,10 +186,15 @@ public class AWS2S3StreamUploadProducer extends 
DefaultProducer {
         byte[] b;
         int maxRead = (getConfiguration().isMultiPartUpload()
                 ? Math.toIntExact(getConfiguration().getPartSize()) : 
getConfiguration().getBufferSize());
-        if (uploadAggregate != null) {
-            uploadAggregate.index++;
-            maxRead -= uploadAggregate.buffer.size();
-            maxRead = Math.max(1, maxRead);
+        lock.lock();
+        try {
+            if (uploadAggregate != null) {
+                uploadAggregate.index++;
+                maxRead -= uploadAggregate.buffer.size();
+                maxRead = Math.max(1, maxRead);
+            }
+        } finally {
+            lock.unlock();
         }
 
         while ((b = AWS2S3Utils.toByteArray(is, maxRead)) != null && b.length 
> 0) {

Reply via email to