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");
+ }
}