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

gnodet pushed a commit to branch camel-4.18.x
in repository https://gitbox.apache.org/repos/asf/camel.git


The following commit(s) were added to refs/heads/camel-4.18.x by this push:
     new cab8265a0b02 [backport camel-4.18.x] CAMEL-24764: 
camel-huaweicloud-obs - close the downloaded object stream (connection leak) 
and fix putObject/charset/proxy issues (#26540)
cab8265a0b02 is described below

commit cab8265a0b02e7169b0dbc1f3e7d28114beade98
Author: Guillaume Nodet - AI Bot <[email protected]>
AuthorDate: Thu Sep 17 13:01:38 2026 +0200

    [backport camel-4.18.x] CAMEL-24764: camel-huaweicloud-obs - close the 
downloaded object stream (connection leak) and fix putObject/charset/proxy 
issues (#26540)
    
    Co-authored-by: gnodet-bot <[email protected]>
---
 .../component/huaweicloud/obs/OBSEndpoint.java     |  2 +-
 .../component/huaweicloud/obs/OBSProducer.java     | 16 ++++++----
 .../camel/component/huaweicloud/obs/OBSUtils.java  |  7 +++--
 .../component/huaweicloud/obs/GetObjectTest.java   | 34 ++++++++++++++++++++++
 4 files changed, 50 insertions(+), 9 deletions(-)

diff --git 
a/components/camel-huawei/camel-huaweicloud-obs/src/main/java/org/apache/camel/component/huaweicloud/obs/OBSEndpoint.java
 
b/components/camel-huawei/camel-huaweicloud-obs/src/main/java/org/apache/camel/component/huaweicloud/obs/OBSEndpoint.java
index 02068d7c519e..5d1eab570560 100644
--- 
a/components/camel-huawei/camel-huaweicloud-obs/src/main/java/org/apache/camel/component/huaweicloud/obs/OBSEndpoint.java
+++ 
b/components/camel-huawei/camel-huaweicloud-obs/src/main/java/org/apache/camel/component/huaweicloud/obs/OBSEndpoint.java
@@ -378,7 +378,7 @@ public class OBSEndpoint extends ScheduledPollEndpoint {
 
         // setup proxy information (if needed)
         if (ObjectHelper.isNotEmpty(getProxyHost())
-                && ObjectHelper.isNotEmpty(getProxyPort())) {
+                && getProxyPort() > 0) {
             HttpProxyConfiguration httpConfig = new HttpProxyConfiguration();
             httpConfig.setProxyAddr(getProxyHost());
             httpConfig.setProxyPort(getProxyPort());
diff --git 
a/components/camel-huawei/camel-huaweicloud-obs/src/main/java/org/apache/camel/component/huaweicloud/obs/OBSProducer.java
 
b/components/camel-huawei/camel-huaweicloud-obs/src/main/java/org/apache/camel/component/huaweicloud/obs/OBSProducer.java
index 575ccb7999c4..ef051494e040 100644
--- 
a/components/camel-huawei/camel-huaweicloud-obs/src/main/java/org/apache/camel/component/huaweicloud/obs/OBSProducer.java
+++ 
b/components/camel-huawei/camel-huaweicloud-obs/src/main/java/org/apache/camel/component/huaweicloud/obs/OBSProducer.java
@@ -19,6 +19,7 @@ package org.apache.camel.component.huaweicloud.obs;
 import java.io.ByteArrayInputStream;
 import java.io.File;
 import java.io.InputStream;
+import java.nio.charset.StandardCharsets;
 import java.util.ArrayList;
 import java.util.List;
 
@@ -127,6 +128,11 @@ public class OBSProducer extends DefaultProducer {
         LOG.trace("Checking if bucket {} exists", 
clientConfigurations.getBucketName());
         if (!obsClient.headBucket(clientConfigurations.getBucketName())) {
             LOG.warn("No bucket found with name {}. Attempting to create", 
clientConfigurations.getBucketName());
+            // bucket location is optional to create a new bucket; default it, 
mirroring the createBucket operation
+            if 
(ObjectHelper.isEmpty(clientConfigurations.getBucketLocation())) {
+                LOG.warn("No bucket location given, defaulting to '{}'", 
OBSConstants.DEFAULT_LOCATION);
+                
clientConfigurations.setBucketLocation(OBSConstants.DEFAULT_LOCATION);
+            }
             
OBSRegion.checkValidRegion(clientConfigurations.getBucketLocation());
             CreateBucketRequest request = new CreateBucketRequest(
                     clientConfigurations.getBucketName(),
@@ -153,10 +159,10 @@ public class OBSProducer extends DefaultProducer {
         } else if (body instanceof String) {
             // the string content will be stored in the remote object
             LOG.trace("Writing text body into an object");
-            InputStream stream = new ByteArrayInputStream(((String) 
body).getBytes());
-            putObjectResult = 
obsClient.putObject(clientConfigurations.getBucketName(),
-                    clientConfigurations.getObjectName(), stream);
-            stream.close();
+            try (InputStream stream = new ByteArrayInputStream(((String) 
body).getBytes(StandardCharsets.UTF_8))) {
+                putObjectResult = 
obsClient.putObject(clientConfigurations.getBucketName(),
+                        clientConfigurations.getObjectName(), stream);
+            }
 
         } else if (body instanceof InputStream) {
             // this covers miscellaneous file types
@@ -186,7 +192,7 @@ public class OBSProducer extends DefaultProducer {
         }
 
         LOG.debug("Downloading remote obs object {} from bucket {}", 
clientConfigurations.getObjectName(),
-                clientConfigurations.getBucketLocation());
+                clientConfigurations.getBucketName());
 
         ObsObject obsObject = obsClient
                 .getObject(clientConfigurations.getBucketName(), 
clientConfigurations.getObjectName());
diff --git 
a/components/camel-huawei/camel-huaweicloud-obs/src/main/java/org/apache/camel/component/huaweicloud/obs/OBSUtils.java
 
b/components/camel-huawei/camel-huaweicloud-obs/src/main/java/org/apache/camel/component/huaweicloud/obs/OBSUtils.java
index 244d09b4589e..9f0c65b0b91a 100644
--- 
a/components/camel-huawei/camel-huaweicloud-obs/src/main/java/org/apache/camel/component/huaweicloud/obs/OBSUtils.java
+++ 
b/components/camel-huawei/camel-huaweicloud-obs/src/main/java/org/apache/camel/component/huaweicloud/obs/OBSUtils.java
@@ -54,9 +54,10 @@ public final class OBSUtils {
     public static void mapObsObject(Exchange exchange, ObsObject obsObject) {
         Message message = exchange.getIn();
 
-        // set exchange body to a byte array of object contents
-        try {
-            message.setBody(OBSUtils.toBytes(obsObject.getObjectContent()));
+        // set exchange body to a byte array of object contents. The SDK 
object-content stream holds a pooled
+        // HTTP connection, so it must be closed once read - otherwise every 
downloaded object leaks a connection.
+        try (InputStream content = obsObject.getObjectContent()) {
+            message.setBody(OBSUtils.toBytes(content));
         } catch (IOException e) {
             throw new RuntimeCamelException(e);
         }
diff --git 
a/components/camel-huawei/camel-huaweicloud-obs/src/test/java/org/apache/camel/component/huaweicloud/obs/GetObjectTest.java
 
b/components/camel-huawei/camel-huaweicloud-obs/src/test/java/org/apache/camel/component/huaweicloud/obs/GetObjectTest.java
index 97f50489a9c4..7956c0dbe2b1 100644
--- 
a/components/camel-huawei/camel-huaweicloud-obs/src/test/java/org/apache/camel/component/huaweicloud/obs/GetObjectTest.java
+++ 
b/components/camel-huawei/camel-huaweicloud-obs/src/test/java/org/apache/camel/component/huaweicloud/obs/GetObjectTest.java
@@ -16,10 +16,14 @@
  */
 package org.apache.camel.component.huaweicloud.obs;
 
+import java.io.ByteArrayInputStream;
 import java.io.File;
 import java.io.FileInputStream;
+import java.io.IOException;
 import java.io.InputStream;
+import java.nio.charset.StandardCharsets;
 import java.util.Date;
+import java.util.concurrent.atomic.AtomicBoolean;
 
 import com.obs.services.ObsClient;
 import com.obs.services.model.ObjectMetadata;
@@ -38,6 +42,7 @@ import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
 
 public class GetObjectTest extends CamelTestSupport {
 
@@ -116,4 +121,33 @@ public class GetObjectTest extends CamelTestSupport {
         assertEquals(objectName, 
responseExchange.getIn().getHeader(Exchange.FILE_NAME));
 
     }
+
+    @Test
+    public void testGetObjectClosesContentStream() throws Exception {
+        ObsObject response = new ObsObject();
+        response.setBucketName(bucketName);
+        response.setObjectKey(objectName);
+
+        // the SDK content stream holds a pooled HTTP connection; the producer 
must close it after reading
+        AtomicBoolean closed = new AtomicBoolean(false);
+        InputStream trackedStream = new 
ByteArrayInputStream("hello".getBytes(StandardCharsets.UTF_8)) {
+            @Override
+            public void close() throws IOException {
+                closed.set(true);
+                super.close();
+            }
+        };
+        response.setObjectContent(trackedStream);
+        response.setMetadata(new ObjectMetadata());
+
+        Mockito.when(mockClient.getObject(bucketName, 
objectName)).thenReturn(response);
+
+        MockEndpoint mock = getMockEndpoint("mock:get_object_result");
+        mock.expectedMinimumMessageCount(1);
+        template.sendBody("direct:get_object", "dummy");
+        mock.assertIsSatisfied();
+
+        assertTrue(closed.get(),
+                "the OBS object content stream must be closed after download 
to release the pooled connection");
+    }
 }

Reply via email to