From 91fb887d6494fdc2ab4df6c619f16c7eafe73ade Mon Sep 17 00:00:00 2001 From: Vijay-Dwivedi Date: Wed, 26 Aug 2026 20:04:05 +0530 Subject: [PATCH 1/2] OBS04O-106 | Close request entity InputStream after executeAndClose --- .../com/emc/object/AbstractJerseyClient.java | 23 ++++- .../com/emc/object/s3/S3JerseyClientTest.java | 97 +++++++++++++++++++ 2 files changed, 118 insertions(+), 2 deletions(-) diff --git a/src/main/java/com/emc/object/AbstractJerseyClient.java b/src/main/java/com/emc/object/AbstractJerseyClient.java index 83fea3a4..0190820f 100644 --- a/src/main/java/com/emc/object/AbstractJerseyClient.java +++ b/src/main/java/com/emc/object/AbstractJerseyClient.java @@ -26,6 +26,8 @@ */ package com.emc.object; +import java.io.IOException; +import java.io.InputStream; import java.net.URI; import java.util.Map; @@ -56,8 +58,25 @@ protected AbstractJerseyClient(ObjectConfig objectConfig) { } protected Response executeAndClose(Client client, ObjectRequest request) { - Response response = executeRequest(client, request); - response.close(); + Response response; + try { + response = executeRequest(client, request); + response.close(); + } finally { + // Jersey 2's Apache connector does not close the request entity input stream after + // consuming it during the request. Close it here so that callers who wrap the stream + // (e.g. with a digest-computing stream) can read the computed result after the request. + if (request instanceof EntityRequest) { + Object entity = ((EntityRequest) request).getEntity(); + if (entity instanceof InputStream) { + try { + ((InputStream) entity).close(); + } catch (IOException e) { + log.warn("could not close request entity stream", e); + } + } + } + } return response; } diff --git a/src/test/java/com/emc/object/s3/S3JerseyClientTest.java b/src/test/java/com/emc/object/s3/S3JerseyClientTest.java index 9d6177e7..a83159d1 100644 --- a/src/test/java/com/emc/object/s3/S3JerseyClientTest.java +++ b/src/test/java/com/emc/object/s3/S3JerseyClientTest.java @@ -61,6 +61,7 @@ import java.net.URL; import java.net.URLEncoder; import java.nio.charset.StandardCharsets; +import java.security.DigestInputStream; import java.security.MessageDigest; import java.util.*; import java.util.concurrent.*; @@ -3649,4 +3650,100 @@ private String getContentMD5(Object obj) { return contentMD5; } + @Test + public void testPutObjectClosesInputStream() throws Exception { + String key = "stream-close-test"; + byte[] data = "Hello Stream Close!".getBytes(StandardCharsets.UTF_8); + + // wrap in a stream that tracks whether close() was called + CloseTrackingInputStream trackingStream = new CloseTrackingInputStream(new ByteArrayInputStream(data)); + + PutObjectRequest request = new PutObjectRequest(getTestBucket(), key, trackingStream); + request.setObjectMetadata(new S3ObjectMetadata().withContentLength((long) data.length)); + client.putObject(request); + + // verify the stream was closed by executeAndClose + Assert.assertTrue("putObject should close the request entity InputStream", trackingStream.isClosed()); + + // verify the object was written correctly + String result = client.readObject(getTestBucket(), key, String.class); + Assert.assertEquals("Hello Stream Close!", result); + } + + @Test + public void testPutObjectStreamDigestAccessible() throws Exception { + String key = "stream-digest-test"; + byte[] data = new byte[1024]; + new Random().nextBytes(data); + String expectedMd5 = Hex.encodeHexString(DigestUtils.md5(data)); + + // wrap in a digest-computing stream (simulates what ecs-sync does) + DigestInputStream digestStream = new DigestInputStream( + new ByteArrayInputStream(data), MessageDigest.getInstance("MD5")); + + PutObjectRequest request = new PutObjectRequest(getTestBucket(), key, digestStream); + request.setObjectMetadata(new S3ObjectMetadata().withContentLength((long) data.length)); + client.putObject(request); + + // after putObject, the stream should be closed and we should be able to read the digest + String actualMd5 = Hex.encodeHexString(digestStream.getMessageDigest().digest()); + Assert.assertEquals("MD5 digest should be accessible after putObject", expectedMd5, actualMd5); + } + + @Test + public void testUploadPartClosesInputStream() throws Exception { + String key = "mpu-stream-close-test"; + byte[] data = new byte[5 * 1024 * 1024]; // 5 MB minimum part size + new Random().nextBytes(data); + + String uploadId = client.initiateMultipartUpload(getTestBucket(), key); + + try { + CloseTrackingInputStream trackingStream = new CloseTrackingInputStream(new ByteArrayInputStream(data)); + + UploadPartRequest request = new UploadPartRequest(getTestBucket(), key, uploadId, 1, trackingStream); + request.setContentLength((long) data.length); + + client.uploadPart(request); + + Assert.assertTrue("uploadPart should close the request entity InputStream", trackingStream.isClosed()); + } finally { + try { + client.abortMultipartUpload(new AbortMultipartUploadRequest(getTestBucket(), key, uploadId)); + } catch (Exception ignored) { + } + } + } + + /** + * Simple InputStream wrapper that tracks whether close() has been called. + */ + private static class CloseTrackingInputStream extends InputStream { + private final InputStream delegate; + private boolean closed = false; + + CloseTrackingInputStream(InputStream delegate) { + this.delegate = delegate; + } + + @Override + public int read() throws IOException { + return delegate.read(); + } + + @Override + public int read(byte[] b, int off, int len) throws IOException { + return delegate.read(b, off, len); + } + + @Override + public void close() throws IOException { + closed = true; + delegate.close(); + } + + boolean isClosed() { + return closed; + } + } } From 07b92db8d6e924884174fa9a049c11b7f33c3197 Mon Sep 17 00:00:00 2001 From: Vijay-Dwivedi Date: Thu, 27 Aug 2026 09:41:17 +0530 Subject: [PATCH 2/2] updated test case --- src/test/java/com/emc/object/s3/S3JerseyClientTest.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/src/test/java/com/emc/object/s3/S3JerseyClientTest.java b/src/test/java/com/emc/object/s3/S3JerseyClientTest.java index a83159d1..b6a596f2 100644 --- a/src/test/java/com/emc/object/s3/S3JerseyClientTest.java +++ b/src/test/java/com/emc/object/s3/S3JerseyClientTest.java @@ -33,6 +33,7 @@ import com.emc.object.s3.bean.*; import com.emc.object.s3.bean.BucketPolicyStatement.Effect; import com.emc.object.s3.jersey.FaultInjectionFilter; +import com.emc.object.s3.jersey.S3EncryptionClient; import com.emc.object.s3.jersey.S3JerseyClient; import com.emc.object.s3.request.*; import com.emc.object.util.RestUtil; @@ -3692,6 +3693,7 @@ public void testPutObjectStreamDigestAccessible() throws Exception { @Test public void testUploadPartClosesInputStream() throws Exception { + Assume.assumeFalse("S3EncryptionClient does not support MPU", client instanceof S3EncryptionClient); String key = "mpu-stream-close-test"; byte[] data = new byte[5 * 1024 * 1024]; // 5 MB minimum part size new Random().nextBytes(data);