This is an automated email from the ASF dual-hosted git repository.
ChenSammi pushed a commit to branch HDDS-13323-sts
in repository https://gitbox.apache.org/repos/asf/ozone.git
The following commit(s) were added to refs/heads/HDDS-13323-sts by this push:
new 9f233336da2 HDDS-15182. Avoid extra read for modification time on
CopyObject/CopyPart (#10203)
9f233336da2 is described below
commit 9f233336da2de800b4ab6a29b195627ff5a21370
Author: fmorg-git <[email protected]>
AuthorDate: Wed Jun 24 01:21:52 2026 -0700
HDDS-15182. Avoid extra read for modification time on CopyObject/CopyPart
(#10203)
---
.../client/io/BlockDataStreamOutputEntryPool.java | 8 +-
.../client/io/BlockOutputStreamEntryPool.java | 8 +-
.../ozone/client/io/KeyDataStreamOutput.java | 4 +
.../hadoop/ozone/client/io/KeyOutputStream.java | 4 +
.../ozone/client/io/OzoneDataStreamOutput.java | 10 +-
.../hadoop/ozone/client/io/OzoneOutputStream.java | 10 +-
.../helpers/OmMultipartCommitUploadPartInfo.java | 9 +-
.../ozone/om/protocol/OzoneManagerProtocol.java | 3 +-
...OzoneManagerProtocolClientSideTranslatorPB.java | 20 ++-
.../hadoop/ozone/s3/awssdk/S3SDKTestUtils.java | 33 +++++
.../hadoop/ozone/s3/awssdk/package-info.java} | 25 +---
.../ozone/s3/awssdk/v1/AbstractS3SDKV1Tests.java | 67 +++++++++
.../ozone/s3/awssdk/v2/AbstractS3SDKV2Tests.java | 70 ++++++++++
.../hadoop/ozone/s3/awssdk/v2/package-info.java} | 25 +---
.../src/main/proto/OmClientProtocol.proto | 4 +
.../ozone/om/request/key/OMKeyCommitRequest.java | 5 +
.../om/request/key/OMKeyCommitRequestWithFSO.java | 5 +
.../S3MultipartUploadCommitPartRequest.java | 1 +
.../hadoop/ozone/s3/endpoint/CopyPartResult.java | 4 +-
.../hadoop/ozone/s3/endpoint/CopyResult.java} | 25 ++--
.../hadoop/ozone/s3/endpoint/ObjectEndpoint.java | 58 +++++---
.../ozone/s3/endpoint/ObjectEndpointStreaming.java | 15 ++-
.../hadoop/ozone/client/OzoneBucketStub.java | 51 ++++---
.../ozone/client/OzoneDataStreamOutputStub.java | 18 ++-
.../hadoop/ozone/client/OzoneOutputStreamStub.java | 20 ++-
...estCopyObjectAndUploadPartCopyLastModified.java | 150 +++++++++++++++++++++
26 files changed, 539 insertions(+), 113 deletions(-)
diff --git
a/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/BlockDataStreamOutputEntryPool.java
b/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/BlockDataStreamOutputEntryPool.java
index ff054c35095..be05a294dcf 100644
---
a/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/BlockDataStreamOutputEntryPool.java
+++
b/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/BlockDataStreamOutputEntryPool.java
@@ -60,6 +60,7 @@ public class BlockDataStreamOutputEntryPool implements
KeyMetadataAware {
private final Map<String, String> metadata = new HashMap<>();
private final XceiverClientFactory xceiverClientFactory;
private OmMultipartCommitUploadPartInfo commitUploadPartInfo;
+ private long modificationTime;
private final long openID;
private final ExcludeList excludeList;
private List<StreamBuffer> bufferList;
@@ -254,8 +255,9 @@ void commitKey(long offset) throws IOException {
if (keyArgs.getIsMultipartKey()) {
commitUploadPartInfo =
omClient.commitMultipartUploadPart(buildKeyArgs(), openID);
+ modificationTime = commitUploadPartInfo.getModificationTime();
} else {
- omClient.commitKey(buildKeyArgs(), openID);
+ modificationTime = omClient.commitKey(buildKeyArgs(), openID);
}
} else {
LOG.warn("Closing KeyDataStreamOutput, but key args is null");
@@ -304,6 +306,10 @@ public OmMultipartCommitUploadPartInfo
getCommitUploadPartInfo() {
return commitUploadPartInfo;
}
+ public long getModificationTime() {
+ return modificationTime;
+ }
+
public ExcludeList getExcludeList() {
return excludeList;
}
diff --git
a/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/BlockOutputStreamEntryPool.java
b/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/BlockOutputStreamEntryPool.java
index f3b98626f09..cef5a7ed056 100644
---
a/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/BlockOutputStreamEntryPool.java
+++
b/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/BlockOutputStreamEntryPool.java
@@ -83,6 +83,7 @@ public class BlockOutputStreamEntryPool implements
KeyMetadataAware {
*/
private final BufferPool bufferPool;
private OmMultipartCommitUploadPartInfo commitUploadPartInfo;
+ private long modificationTime;
private final long openID;
private final ExcludeList excludeList;
private final ContainerClientMetrics clientMetrics;
@@ -329,8 +330,9 @@ void commitKey(long offset) throws IOException {
if (keyArgs.getIsMultipartKey()) {
commitUploadPartInfo =
omClient.commitMultipartUploadPart(buildKeyArgs(), openID);
+ modificationTime = commitUploadPartInfo.getModificationTime();
} else {
- omClient.commitKey(buildKeyArgs(), openID);
+ modificationTime = omClient.commitKey(buildKeyArgs(), openID);
}
} else {
LOG.warn("Closing KeyOutputStream, but key args is null");
@@ -422,6 +424,10 @@ public OmMultipartCommitUploadPartInfo
getCommitUploadPartInfo() {
return commitUploadPartInfo;
}
+ public long getModificationTime() {
+ return modificationTime;
+ }
+
public ExcludeList getExcludeList() {
return excludeList;
}
diff --git
a/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/KeyDataStreamOutput.java
b/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/KeyDataStreamOutput.java
index ceacd624e93..2d4bc9486cf 100644
---
a/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/KeyDataStreamOutput.java
+++
b/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/KeyDataStreamOutput.java
@@ -476,6 +476,10 @@ public OmMultipartCommitUploadPartInfo
getCommitUploadPartInfo() {
return blockDataStreamOutputEntryPool.getCommitUploadPartInfo();
}
+ public long getModificationTime() {
+ return blockDataStreamOutputEntryPool.getModificationTime();
+ }
+
@VisibleForTesting
public ExcludeList getExcludeList() {
return blockDataStreamOutputEntryPool.getExcludeList();
diff --git
a/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/KeyOutputStream.java
b/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/KeyOutputStream.java
index 2f9edfa94ea..8c805607883 100644
---
a/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/KeyOutputStream.java
+++
b/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/KeyOutputStream.java
@@ -676,6 +676,10 @@ private void closeInternal() throws IOException {
return blockOutputStreamEntryPool.getCommitUploadPartInfo();
}
+ public long getModificationTime() {
+ return blockOutputStreamEntryPool.getModificationTime();
+ }
+
@VisibleForTesting
public ExcludeList getExcludeList() {
return blockOutputStreamEntryPool.getExcludeList();
diff --git
a/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/OzoneDataStreamOutput.java
b/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/OzoneDataStreamOutput.java
index 7ce3f71b375..1351c90ad33 100644
---
a/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/OzoneDataStreamOutput.java
+++
b/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/OzoneDataStreamOutput.java
@@ -100,7 +100,7 @@ public synchronized void close() throws IOException {
}
public OmMultipartCommitUploadPartInfo getCommitUploadPartInfo() {
- KeyDataStreamOutput keyDataStreamOutput = getKeyDataStreamOutput();
+ final KeyDataStreamOutput keyDataStreamOutput = getKeyDataStreamOutput();
if (keyDataStreamOutput != null) {
return keyDataStreamOutput.getCommitUploadPartInfo();
}
@@ -108,6 +108,14 @@ public OmMultipartCommitUploadPartInfo
getCommitUploadPartInfo() {
return null;
}
+ public long getModificationTime() {
+ final KeyDataStreamOutput keyDataStreamOutput = getKeyDataStreamOutput();
+ if (keyDataStreamOutput != null) {
+ return keyDataStreamOutput.getModificationTime();
+ }
+ throw new IllegalStateException("OutputStream is not a
KeyDataStreamOutput: " + byteBufferStreamOutput.getClass());
+ }
+
public KeyDataStreamOutput getKeyDataStreamOutput() {
if (byteBufferStreamOutput instanceof OzoneOutputStream) {
OutputStream outputStream =
diff --git
a/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/OzoneOutputStream.java
b/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/OzoneOutputStream.java
index c0e14b089ef..1e01998bc45 100644
---
a/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/OzoneOutputStream.java
+++
b/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/OzoneOutputStream.java
@@ -128,7 +128,7 @@ public void hsync() throws IOException {
}
public OmMultipartCommitUploadPartInfo getCommitUploadPartInfo() {
- KeyOutputStream keyOutputStream = getKeyOutputStream();
+ final KeyOutputStream keyOutputStream = getKeyOutputStream();
if (keyOutputStream != null) {
return keyOutputStream.getCommitUploadPartInfo();
}
@@ -136,6 +136,14 @@ public OmMultipartCommitUploadPartInfo
getCommitUploadPartInfo() {
return null;
}
+ public long getModificationTime() {
+ final KeyOutputStream keyOutputStream = getKeyOutputStream();
+ if (keyOutputStream != null) {
+ return keyOutputStream.getModificationTime();
+ }
+ throw new IllegalStateException("OutputStream is not a KeyOutputStream: "
+ outputStream.getClass());
+ }
+
public OutputStream getOutputStream() {
return outputStream;
}
diff --git
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartCommitUploadPartInfo.java
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartCommitUploadPartInfo.java
index 93774a82dc5..0c072b65277 100644
---
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartCommitUploadPartInfo.java
+++
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartCommitUploadPartInfo.java
@@ -27,9 +27,12 @@ public class OmMultipartCommitUploadPartInfo {
private final String eTag;
- public OmMultipartCommitUploadPartInfo(String partName, String eTag) {
+ private final long modificationTime;
+
+ public OmMultipartCommitUploadPartInfo(String partName, String eTag, long
modificationTime) {
this.partName = partName;
this.eTag = eTag;
+ this.modificationTime = modificationTime;
}
public String getETag() {
@@ -39,4 +42,8 @@ public String getETag() {
public String getPartName() {
return partName;
}
+
+ public long getModificationTime() {
+ return modificationTime;
+ }
}
diff --git
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocol/OzoneManagerProtocol.java
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocol/OzoneManagerProtocol.java
index dd884fdf29c..16f9e3dc39f 100644
---
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocol/OzoneManagerProtocol.java
+++
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocol/OzoneManagerProtocol.java
@@ -244,9 +244,10 @@ default OpenKeySession openKey(OmKeyArgs args) throws
IOException {
*
* @param args the key to commit
* @param clientID the client identification
+ * @return the modification time of the committed key in epoch milliseconds
* @throws IOException
*/
- default void commitKey(OmKeyArgs args, long clientID)
+ default long commitKey(OmKeyArgs args, long clientID)
throws IOException {
throw new UnsupportedOperationException("OzoneManager does not require " +
"this to be implemented, as write requests use a new approach.");
diff --git
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocolPB/OzoneManagerProtocolClientSideTranslatorPB.java
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocolPB/OzoneManagerProtocolClientSideTranslatorPB.java
index 8fdf9712c08..a074d7b6228 100644
---
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocolPB/OzoneManagerProtocolClientSideTranslatorPB.java
+++
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocolPB/OzoneManagerProtocolClientSideTranslatorPB.java
@@ -253,6 +253,7 @@
import org.apache.hadoop.ozone.upgrade.UpgradeFinalization.StatusAndMessages;
import org.apache.hadoop.ozone.util.ProtobufUtils;
import org.apache.hadoop.security.token.Token;
+import org.apache.hadoop.util.Time;
/**
* The client side implementation of OzoneManagerProtocol.
@@ -823,9 +824,9 @@ public void hsyncKey(OmKeyArgs args, long clientId)
}
@Override
- public void commitKey(OmKeyArgs args, long clientId)
+ public long commitKey(OmKeyArgs args, long clientId)
throws IOException {
- updateKey(args, clientId, false, false);
+ return updateKeyAndGetModificationTime(args, clientId, false, false);
}
@Override
@@ -849,6 +850,12 @@ public static void setReplicationConfig(ReplicationConfig
replication,
private void updateKey(OmKeyArgs args, long clientId, boolean hsync, boolean
recovery)
throws IOException {
+ // Preserve legacy behavior (ignore response payload).
+ updateKeyAndGetModificationTime(args, clientId, hsync, recovery);
+ }
+
+ private long updateKeyAndGetModificationTime(OmKeyArgs args, long clientId,
boolean hsync, boolean recovery)
+ throws IOException {
CommitKeyRequest.Builder req = CommitKeyRequest.newBuilder();
List<OmKeyLocationInfo> locationInfoList = args.getLocationInfoList();
Objects.requireNonNull(locationInfoList, "locationInfoList == null");
@@ -874,7 +881,11 @@ private void updateKey(OmKeyArgs args, long clientId,
boolean hsync, boolean rec
.setCommitKeyRequest(req)
.build();
- handleError(submitRequest(omRequest));
+ final OMResponse resp = handleError(submitRequest(omRequest));
+ if (resp.hasCommitKeyResponse() &&
resp.getCommitKeyResponse().hasModificationTime()) {
+ return resp.getCommitKeyResponse().getModificationTime();
+ }
+ return Time.now();
}
@Override
@@ -1791,7 +1802,8 @@ public OmMultipartCommitUploadPartInfo
commitMultipartUploadPart(
handleError(submitRequest(omRequest))
.getCommitMultiPartUploadResponse();
- return new OmMultipartCommitUploadPartInfo(response.getPartName(),
response.getETag());
+ final long modificationTime = response.hasModificationTime() ?
response.getModificationTime() : Time.now();
+ return new OmMultipartCommitUploadPartInfo(response.getPartName(),
response.getETag(), modificationTime);
}
@Override
diff --git
a/hadoop-ozone/integration-test-s3/src/test/java/org/apache/hadoop/ozone/s3/awssdk/S3SDKTestUtils.java
b/hadoop-ozone/integration-test-s3/src/test/java/org/apache/hadoop/ozone/s3/awssdk/S3SDKTestUtils.java
index ec42a0d7b4f..7c5a12dce95 100644
---
a/hadoop-ozone/integration-test-s3/src/test/java/org/apache/hadoop/ozone/s3/awssdk/S3SDKTestUtils.java
+++
b/hadoop-ozone/integration-test-s3/src/test/java/org/apache/hadoop/ozone/s3/awssdk/S3SDKTestUtils.java
@@ -31,6 +31,10 @@
import java.util.regex.Pattern;
import org.apache.commons.io.IOUtils;
import org.apache.commons.lang3.RandomUtils;
+import org.apache.hadoop.ozone.MiniOzoneCluster;
+import org.apache.hadoop.ozone.client.OzoneBucket;
+import org.apache.hadoop.ozone.client.OzoneClient;
+import org.apache.hadoop.ozone.client.OzoneMultipartUploadPartListParts;
import org.apache.ozone.test.InputSubstream;
/**
@@ -40,9 +44,38 @@ public final class S3SDKTestUtils {
public static final Pattern UPLOAD_ID_PATTERN =
Pattern.compile("<UploadId>(.+?)</UploadId>");
+ private static final int DEFAULT_LIST_PARTS_MAX = 100;
+
private S3SDKTestUtils() {
}
+ /**
+ * Returns the modification time of a key stored in OM, in epoch
milliseconds.
+ */
+ public static long getOmKeyModificationTime(MiniOzoneCluster cluster, String
bucketName, String keyName)
+ throws IOException {
+ try (OzoneClient ozoneClient = cluster.newClient()) {
+ final OzoneBucket bucket =
ozoneClient.getObjectStore().getS3Volume().getBucket(bucketName);
+ return bucket.getKey(keyName).getModificationTime().toEpochMilli();
+ }
+ }
+
+ /**
+ * Returns the modification time of a multipart upload part stored in OM, in
epoch milliseconds.
+ */
+ public static long getOmPartModificationTime(MiniOzoneCluster cluster,
String bucketName, String keyName,
+ String uploadId, int partNumber) throws IOException {
+ try (OzoneClient ozoneClient = cluster.newClient()) {
+ final OzoneBucket bucket =
ozoneClient.getObjectStore().getS3Volume().getBucket(bucketName);
+ final OzoneMultipartUploadPartListParts parts =
bucket.listParts(keyName, uploadId, 0, DEFAULT_LIST_PARTS_MAX);
+ return parts.getPartInfoList().stream()
+ .filter(part -> part.getPartNumber() == partNumber)
+ .findFirst()
+ .orElseThrow(() -> new IllegalStateException("Part " + partNumber +
" not found for upload " + uploadId))
+ .getModificationTime();
+ }
+ }
+
/**
* Calculate the MD5 digest from an input stream from a specific offset and
length.
* @param inputStream The input stream where the digest will be read from.
diff --git
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartCommitUploadPartInfo.java
b/hadoop-ozone/integration-test-s3/src/test/java/org/apache/hadoop/ozone/s3/awssdk/package-info.java
similarity index 62%
copy from
hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartCommitUploadPartInfo.java
copy to
hadoop-ozone/integration-test-s3/src/test/java/org/apache/hadoop/ozone/s3/awssdk/package-info.java
index 93774a82dc5..3ecbb322cb4 100644
---
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartCommitUploadPartInfo.java
+++
b/hadoop-ozone/integration-test-s3/src/test/java/org/apache/hadoop/ozone/s3/awssdk/package-info.java
@@ -15,28 +15,7 @@
* limitations under the License.
*/
-package org.apache.hadoop.ozone.om.helpers;
-
/**
- * This class holds information about the response from commit multipart
- * upload part request.
+ * Shared utilities and integration tests for the AWS S3 SDK.
*/
-public class OmMultipartCommitUploadPartInfo {
-
- private final String partName;
-
- private final String eTag;
-
- public OmMultipartCommitUploadPartInfo(String partName, String eTag) {
- this.partName = partName;
- this.eTag = eTag;
- }
-
- public String getETag() {
- return eTag;
- }
-
- public String getPartName() {
- return partName;
- }
-}
+package org.apache.hadoop.ozone.s3.awssdk;
diff --git
a/hadoop-ozone/integration-test-s3/src/test/java/org/apache/hadoop/ozone/s3/awssdk/v1/AbstractS3SDKV1Tests.java
b/hadoop-ozone/integration-test-s3/src/test/java/org/apache/hadoop/ozone/s3/awssdk/v1/AbstractS3SDKV1Tests.java
index 8238ade3921..4769bbf27fb 100644
---
a/hadoop-ozone/integration-test-s3/src/test/java/org/apache/hadoop/ozone/s3/awssdk/v1/AbstractS3SDKV1Tests.java
+++
b/hadoop-ozone/integration-test-s3/src/test/java/org/apache/hadoop/ozone/s3/awssdk/v1/AbstractS3SDKV1Tests.java
@@ -20,6 +20,8 @@
import static org.apache.hadoop.ozone.OzoneConsts.MB;
import static org.apache.hadoop.ozone.s3.awssdk.S3SDKTestUtils.calculateDigest;
import static org.apache.hadoop.ozone.s3.awssdk.S3SDKTestUtils.createFile;
+import static
org.apache.hadoop.ozone.s3.awssdk.S3SDKTestUtils.getOmKeyModificationTime;
+import static
org.apache.hadoop.ozone.s3.awssdk.S3SDKTestUtils.getOmPartModificationTime;
import static
org.apache.hadoop.ozone.s3.util.S3Consts.CUSTOM_METADATA_HEADER_PREFIX;
import static org.apache.hadoop.ozone.s3.util.S3Utils.stripQuotes;
import static org.assertj.core.api.Assertions.assertThat;
@@ -43,6 +45,8 @@
import com.amazonaws.services.s3.model.CompleteMultipartUploadResult;
import com.amazonaws.services.s3.model.CopyObjectRequest;
import com.amazonaws.services.s3.model.CopyObjectResult;
+import com.amazonaws.services.s3.model.CopyPartRequest;
+import com.amazonaws.services.s3.model.CopyPartResult;
import com.amazonaws.services.s3.model.CreateBucketRequest;
import com.amazonaws.services.s3.model.GeneratePresignedUrlRequest;
import com.amazonaws.services.s3.model.GetObjectRequest;
@@ -514,6 +518,69 @@ public void testCopyObject() {
assertEquals("37b51d194a7513e45b56f6524f2d51f2", copyResult.getETag());
}
+ @Test
+ public void testCopyObjectLastModifiedMatchesOmDb() throws Exception {
+ final String sourceBucketName = getBucketName();
+ final String destBucketName = getBucketName();
+ final String sourceKey = getKeyName();
+ final String destKey = getKeyName();
+ final String content = "copy-last-modified-test-content";
+ s3Client.createBucket(sourceBucketName);
+ s3Client.createBucket(destBucketName);
+
+ final InputStream inputStream = new
ByteArrayInputStream(content.getBytes(StandardCharsets.UTF_8));
+ s3Client.putObject(sourceBucketName, sourceKey, inputStream, new
ObjectMetadata());
+
+ final long sourceModificationTime = getOmKeyModificationTime(cluster,
sourceBucketName, sourceKey);
+
+ Thread.sleep(100);
+
+ final CopyObjectResult copyResult = s3Client.copyObject(sourceBucketName,
sourceKey, destBucketName, destKey);
+
+ assertNotNull(copyResult.getLastModifiedDate());
+ final long responseModificationTime =
copyResult.getLastModifiedDate().getTime();
+ final long destModificationTime = getOmKeyModificationTime(cluster,
destBucketName, destKey);
+
+ assertEquals(destModificationTime, responseModificationTime);
+ assertNotEquals(sourceModificationTime, responseModificationTime);
+ }
+
+ @Test
+ public void testUploadPartCopyLastModifiedMatchesOmDb() throws Exception {
+ final String bucketName = getBucketName();
+ final String sourceKey = getKeyName();
+ final String destKey = getKeyName();
+ final String content = "copy-last-modified-test-content";
+ s3Client.createBucket(bucketName);
+
+ final InputStream inputStream = new
ByteArrayInputStream(content.getBytes(StandardCharsets.UTF_8));
+ s3Client.putObject(bucketName, sourceKey, inputStream, new
ObjectMetadata());
+
+ final long sourceModificationTime = getOmKeyModificationTime(cluster,
bucketName, sourceKey);
+
+ Thread.sleep(100);
+
+ final InitiateMultipartUploadResult initResponse =
s3Client.initiateMultipartUpload(
+ new InitiateMultipartUploadRequest(bucketName, destKey));
+ final String uploadId = initResponse.getUploadId();
+
+ final CopyPartResult copyPartResult = s3Client.copyPart(
+ new CopyPartRequest()
+ .withSourceBucketName(bucketName)
+ .withSourceKey(sourceKey)
+ .withDestinationBucketName(bucketName)
+ .withDestinationKey(destKey)
+ .withUploadId(uploadId)
+ .withPartNumber(1));
+
+ assertNotNull(copyPartResult.getLastModifiedDate());
+ final long responseModificationTime =
copyPartResult.getLastModifiedDate().getTime();
+ final long destPartModificationTime = getOmPartModificationTime(cluster,
bucketName, destKey, uploadId, 1);
+
+ assertEquals(destPartModificationTime, responseModificationTime);
+ assertNotEquals(sourceModificationTime, responseModificationTime);
+ }
+
@Test
public void testCopyObjectWithSourceIfMatch() {
final String sourceBucketName = getBucketName("source");
diff --git
a/hadoop-ozone/integration-test-s3/src/test/java/org/apache/hadoop/ozone/s3/awssdk/v2/AbstractS3SDKV2Tests.java
b/hadoop-ozone/integration-test-s3/src/test/java/org/apache/hadoop/ozone/s3/awssdk/v2/AbstractS3SDKV2Tests.java
index 28d8cbf1f61..97c40fa1567 100644
---
a/hadoop-ozone/integration-test-s3/src/test/java/org/apache/hadoop/ozone/s3/awssdk/v2/AbstractS3SDKV2Tests.java
+++
b/hadoop-ozone/integration-test-s3/src/test/java/org/apache/hadoop/ozone/s3/awssdk/v2/AbstractS3SDKV2Tests.java
@@ -20,6 +20,8 @@
import static org.apache.hadoop.ozone.OzoneConsts.MB;
import static org.apache.hadoop.ozone.s3.awssdk.S3SDKTestUtils.calculateDigest;
import static org.apache.hadoop.ozone.s3.awssdk.S3SDKTestUtils.createFile;
+import static
org.apache.hadoop.ozone.s3.awssdk.S3SDKTestUtils.getOmKeyModificationTime;
+import static
org.apache.hadoop.ozone.s3.awssdk.S3SDKTestUtils.getOmPartModificationTime;
import static org.apache.hadoop.ozone.s3.util.S3Utils.stripQuotes;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
@@ -141,6 +143,7 @@
import software.amazon.awssdk.services.s3.model.Tag;
import software.amazon.awssdk.services.s3.model.Tagging;
import software.amazon.awssdk.services.s3.model.UploadPartCopyRequest;
+import software.amazon.awssdk.services.s3.model.UploadPartCopyResponse;
import software.amazon.awssdk.services.s3.model.UploadPartRequest;
import software.amazon.awssdk.services.s3.model.UploadPartResponse;
import software.amazon.awssdk.services.s3.presigner.S3Presigner;
@@ -992,6 +995,73 @@ public void testCopyObject() {
assertEquals("\"37b51d194a7513e45b56f6524f2d51f2\"",
copyObjectResponse.copyObjectResult().eTag());
}
+ @Test
+ public void testCopyObjectLastModifiedMatchesOmDb() throws Exception {
+ final String sourceBucketName = getBucketName();
+ final String destBucketName = getBucketName();
+ final String sourceKey = getKeyName();
+ final String destKey = getKeyName();
+ final String content = "copy-last-modified-test-content";
+ s3Client.createBucket(b -> b.bucket(sourceBucketName));
+ s3Client.createBucket(b -> b.bucket(destBucketName));
+
+ s3Client.putObject(b -> b.bucket(sourceBucketName).key(sourceKey),
RequestBody.fromString(content));
+
+ final long sourceModificationTime = getOmKeyModificationTime(cluster,
sourceBucketName, sourceKey);
+
+ Thread.sleep(100);
+
+ final CopyObjectResponse copyObjectResponse = s3Client.copyObject(
+ CopyObjectRequest.builder()
+ .sourceBucket(sourceBucketName)
+ .sourceKey(sourceKey)
+ .destinationBucket(destBucketName)
+ .destinationKey(destKey)
+ .build());
+
+ assertNotNull(copyObjectResponse.copyObjectResult().lastModified());
+ final long responseModificationTime =
copyObjectResponse.copyObjectResult().lastModified().toEpochMilli();
+ final long destModificationTime = getOmKeyModificationTime(cluster,
destBucketName, destKey);
+
+ assertEquals(destModificationTime, responseModificationTime);
+ assertNotEquals(sourceModificationTime, responseModificationTime);
+ }
+
+ @Test
+ public void testUploadPartCopyLastModifiedMatchesOmDb() throws Exception {
+ final String bucketName = getBucketName();
+ final String sourceKey = getKeyName();
+ final String destKey = getKeyName();
+ final String content = "copy-last-modified-test-content";
+ s3Client.createBucket(b -> b.bucket(bucketName));
+
+ s3Client.putObject(b -> b.bucket(bucketName).key(sourceKey),
RequestBody.fromString(content));
+
+ final long sourceModificationTime = getOmKeyModificationTime(cluster,
bucketName, sourceKey);
+
+ Thread.sleep(100);
+
+ final CreateMultipartUploadResponse initResponse =
s3Client.createMultipartUpload(
+ b -> b.bucket(bucketName).key(destKey));
+ final String uploadId = initResponse.uploadId();
+
+ final UploadPartCopyResponse copyPartResponse = s3Client.uploadPartCopy(
+ b -> b
+ .sourceBucket(bucketName)
+ .sourceKey(sourceKey)
+ .destinationBucket(bucketName)
+ .destinationKey(destKey)
+ .uploadId(uploadId)
+ .partNumber(1));
+
+ assertNotNull(copyPartResponse.copyPartResult().lastModified());
+ final long responseModificationTime =
copyPartResponse.copyPartResult().lastModified().toEpochMilli();
+ final long destPartModificationTime = getOmPartModificationTime(cluster,
bucketName, destKey, uploadId, 1);
+
+ assertEquals(destPartModificationTime, responseModificationTime);
+ assertNotEquals(sourceModificationTime, responseModificationTime);
+ }
+
@Test
public void testCopyObjectWithSourceIfMatch() {
final String sourceBucketName = getBucketName("source");
diff --git
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartCommitUploadPartInfo.java
b/hadoop-ozone/integration-test-s3/src/test/java/org/apache/hadoop/ozone/s3/awssdk/v2/package-info.java
similarity index 62%
copy from
hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartCommitUploadPartInfo.java
copy to
hadoop-ozone/integration-test-s3/src/test/java/org/apache/hadoop/ozone/s3/awssdk/v2/package-info.java
index 93774a82dc5..9e921b64575 100644
---
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartCommitUploadPartInfo.java
+++
b/hadoop-ozone/integration-test-s3/src/test/java/org/apache/hadoop/ozone/s3/awssdk/v2/package-info.java
@@ -15,28 +15,7 @@
* limitations under the License.
*/
-package org.apache.hadoop.ozone.om.helpers;
-
/**
- * This class holds information about the response from commit multipart
- * upload part request.
+ * Integration tests for the AWS Java SDK v2 S3 client.
*/
-public class OmMultipartCommitUploadPartInfo {
-
- private final String partName;
-
- private final String eTag;
-
- public OmMultipartCommitUploadPartInfo(String partName, String eTag) {
- this.partName = partName;
- this.eTag = eTag;
- }
-
- public String getETag() {
- return eTag;
- }
-
- public String getPartName() {
- return partName;
- }
-}
+package org.apache.hadoop.ozone.s3.awssdk.v2;
diff --git
a/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto
b/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto
index 7674d85ca92..16a7e844ab1 100644
--- a/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto
+++ b/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto
@@ -1586,6 +1586,8 @@ message CommitKeyRequest {
message CommitKeyResponse {
+ // Modification time of the committed key, set by OM during preExecute.
+ optional uint64 modificationTime = 1;
}
message AllocateBlockRequest {
@@ -1790,6 +1792,8 @@ message MultipartCommitUploadPartResponse {
optional string partName = 1;
// This one is returned as Etag for S3.
optional string eTag = 2;
+ // Modification time of the committed part key, set by OM during
preExecute.
+ optional uint64 modificationTime = 3;
}
message MultipartUploadCompleteRequest {
diff --git
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCommitRequest.java
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCommitRequest.java
index 3d8bf932e09..4184b3908ca 100644
---
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCommitRequest.java
+++
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCommitRequest.java
@@ -62,6 +62,7 @@
import org.apache.hadoop.ozone.om.response.key.OMKeyCommitResponse;
import org.apache.hadoop.ozone.om.upgrade.OMLayoutFeature;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.CommitKeyRequest;
+import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.CommitKeyResponse;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.KeyArgs;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.KeyLocation;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.OMRequest;
@@ -409,6 +410,10 @@ public OMClientResponse
validateAndUpdateCache(OzoneManager ozoneManager, Execut
omBucketInfo.incrUsedBytes(correctedSpace);
+ omResponse.setCommitKeyResponse(CommitKeyResponse.newBuilder()
+ .setModificationTime(commitKeyArgs.getModificationTime())
+ .build());
+
omClientResponse = new OMKeyCommitResponse(omResponse.build(),
omKeyInfo, dbOzoneKey, dbOpenKey, omBucketInfo.copyObject(),
oldKeyVersionsToDeleteMap, isHSync, newOpenKeyInfo,
dbOpenKeyToDeleteKey, openKeyToDelete);
diff --git
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCommitRequestWithFSO.java
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCommitRequestWithFSO.java
index 25b5a4b15d4..b3d06b3173b 100644
---
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCommitRequestWithFSO.java
+++
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCommitRequestWithFSO.java
@@ -52,6 +52,7 @@
import org.apache.hadoop.ozone.om.response.OMClientResponse;
import org.apache.hadoop.ozone.om.response.key.OMKeyCommitResponseWithFSO;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.CommitKeyRequest;
+import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.CommitKeyResponse;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.KeyArgs;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.OMRequest;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.OMResponse;
@@ -347,6 +348,10 @@ public OMClientResponse
validateAndUpdateCache(OzoneManager ozoneManager, Execut
omBucketInfo.incrUsedBytes(correctedSpace);
+ omResponse.setCommitKeyResponse(CommitKeyResponse.newBuilder()
+ .setModificationTime(commitKeyArgs.getModificationTime())
+ .build());
+
omClientResponse = new OMKeyCommitResponseWithFSO(omResponse.build(),
omKeyInfo, dbFileKey, dbOpenFileKey, omBucketInfo.copyObject(),
oldKeyVersionsToDeleteMap, volumeId, isHSync, newOpenKeyInfo,
dbOpenKeyToDeleteKey, openKeyToDelete);
diff --git
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadCommitPartRequest.java
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadCommitPartRequest.java
index ac123ff680a..fd9eead76d3 100644
---
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadCommitPartRequest.java
+++
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadCommitPartRequest.java
@@ -268,6 +268,7 @@ public OMClientResponse validateAndUpdateCache(OzoneManager
ozoneManager, Execut
if (eTag != null) {
commitResponseBuilder.setETag(eTag);
}
+ commitResponseBuilder.setModificationTime(keyArgs.getModificationTime());
omResponse.setCommitMultiPartUploadResponse(commitResponseBuilder);
omClientResponse =
getOmClientResponse(ozoneManager, keyVersionsToDeleteMap, openKey,
diff --git
a/hadoop-ozone/s3gateway/src/main/java/org/apache/hadoop/ozone/s3/endpoint/CopyPartResult.java
b/hadoop-ozone/s3gateway/src/main/java/org/apache/hadoop/ozone/s3/endpoint/CopyPartResult.java
index f3b8b6e60e6..7f652bad331 100644
---
a/hadoop-ozone/s3gateway/src/main/java/org/apache/hadoop/ozone/s3/endpoint/CopyPartResult.java
+++
b/hadoop-ozone/s3gateway/src/main/java/org/apache/hadoop/ozone/s3/endpoint/CopyPartResult.java
@@ -45,9 +45,9 @@ public class CopyPartResult {
public CopyPartResult() {
}
- public CopyPartResult(String eTag) {
+ public CopyPartResult(String eTag, Instant lastModified) {
this.eTag = eTag;
- this.lastModified = Instant.now();
+ this.lastModified = lastModified;
}
public Instant getLastModified() {
diff --git
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartCommitUploadPartInfo.java
b/hadoop-ozone/s3gateway/src/main/java/org/apache/hadoop/ozone/s3/endpoint/CopyResult.java
similarity index 68%
copy from
hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartCommitUploadPartInfo.java
copy to
hadoop-ozone/s3gateway/src/main/java/org/apache/hadoop/ozone/s3/endpoint/CopyResult.java
index 93774a82dc5..8b67ddbcbaa 100644
---
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartCommitUploadPartInfo.java
+++
b/hadoop-ozone/s3gateway/src/main/java/org/apache/hadoop/ozone/s3/endpoint/CopyResult.java
@@ -15,28 +15,31 @@
* limitations under the License.
*/
-package org.apache.hadoop.ozone.om.helpers;
+package org.apache.hadoop.ozone.s3.endpoint;
/**
- * This class holds information about the response from commit multipart
- * upload part request.
+ * Result of a copy operation.
*/
-public class OmMultipartCommitUploadPartInfo {
-
- private final String partName;
-
+public class CopyResult {
private final String eTag;
+ private final long size;
+ private final long modificationTime;
- public OmMultipartCommitUploadPartInfo(String partName, String eTag) {
- this.partName = partName;
+ public CopyResult(String eTag, long size, long modificationTime) {
this.eTag = eTag;
+ this.size = size;
+ this.modificationTime = modificationTime;
}
public String getETag() {
return eTag;
}
- public String getPartName() {
- return partName;
+ public long getSize() {
+ return size;
+ }
+
+ public long getModificationTime() {
+ return modificationTime;
}
}
diff --git
a/hadoop-ozone/s3gateway/src/main/java/org/apache/hadoop/ozone/s3/endpoint/ObjectEndpoint.java
b/hadoop-ozone/s3gateway/src/main/java/org/apache/hadoop/ozone/s3/endpoint/ObjectEndpoint.java
index bfa3c3d5c79..9ab04c8502e 100644
---
a/hadoop-ozone/s3gateway/src/main/java/org/apache/hadoop/ozone/s3/endpoint/ObjectEndpoint.java
+++
b/hadoop-ozone/s3gateway/src/main/java/org/apache/hadoop/ozone/s3/endpoint/ObjectEndpoint.java
@@ -18,7 +18,6 @@
package org.apache.hadoop.ozone.s3.endpoint;
import static
org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationType.EC;
-import static
org.apache.hadoop.ozone.audit.AuditLogger.PerformanceStringBuilder;
import static
org.apache.hadoop.ozone.s3.S3GatewayConfigKeys.OZONE_S3G_FSO_DIRECTORY_CREATION_ENABLED;
import static
org.apache.hadoop.ozone.s3.S3GatewayConfigKeys.OZONE_S3G_FSO_DIRECTORY_CREATION_ENABLED_DEFAULT;
import static
org.apache.hadoop.ozone.s3.exception.S3ErrorTable.INVALID_ARGUMENT;
@@ -81,7 +80,9 @@
import org.apache.commons.lang3.tuple.Pair;
import org.apache.hadoop.hdds.client.ECReplicationConfig;
import org.apache.hadoop.hdds.client.ReplicationConfig;
+import org.apache.hadoop.hdds.scm.client.HddsClientUtils;
import org.apache.hadoop.ozone.OzoneConsts;
+import org.apache.hadoop.ozone.audit.AuditLogger.PerformanceStringBuilder;
import org.apache.hadoop.ozone.audit.S3GAction;
import org.apache.hadoop.ozone.client.OzoneBucket;
import org.apache.hadoop.ozone.client.OzoneKey;
@@ -955,7 +956,8 @@ private Response createMultipartKey(OzoneVolume volume,
OzoneBucket ozoneBucket,
if (copyHeader != null) {
getMetrics().updateCopyObjectSuccessStats(startNanos);
- return Response.ok(new CopyPartResult(eTag)).build();
+ final Instant lastModified =
Instant.ofEpochMilli(omMultipartCommitUploadPartInfo.getModificationTime());
+ return Response.ok(new CopyPartResult(eTag, lastModified)).build();
} else {
getMetrics().updateCreateMultipartKeySuccessStats(startNanos);
return Response.ok().header(HttpHeaders.ETAG, eTag).build();
@@ -976,6 +978,22 @@ private Response createMultipartKey(OzoneVolume volume,
OzoneBucket ozoneBucket,
throw os3Exception;
}
throw newError(bucketName, key, ex);
+ } catch (IOException ex) {
+ // Ensure we handle permission failures - these can surface as
IOException wrapping OMException.
+ if (copyHeader != null) {
+ getMetrics().updateCopyObjectFailureStats(startNanos);
+ } else {
+ getMetrics().updateCreateMultipartKeyFailureStats(startNanos);
+ }
+ final OMException omEx = (OMException)
HddsClientUtils.containsException(ex, OMException.class);
+ if (omEx != null) {
+ if (omEx.getResult() == ResultCodes.NO_SUCH_MULTIPART_UPLOAD_ERROR) {
+ throw newError(NO_SUCH_UPLOAD, uploadID, omEx);
+ } else {
+ throw newError(bucketName, key, omEx);
+ }
+ }
+ throw ex;
} finally {
// Reset the thread-local message digest instance in case of exception
// and MessageDigest#digest is never called
@@ -986,7 +1004,7 @@ private Response createMultipartKey(OzoneVolume volume,
OzoneBucket ozoneBucket,
}
@SuppressWarnings("checkstyle:ParameterNumber")
- void copy(OzoneVolume volume, DigestInputStream src, long srcKeyLen,
+ CopyResult copy(OzoneVolume volume, DigestInputStream src, long srcKeyLen,
String destKey, String destBucket,
ReplicationConfig replication,
Map<String, String> metadata,
@@ -995,29 +1013,36 @@ void copy(OzoneVolume volume, DigestInputStream src,
long srcKeyLen,
S3ConditionalRequest.WriteConditions writeConditions)
throws IOException {
long copyLength;
-
+ final String eTag;
+ final long modificationTime;
if (isDatastreamEnabled() && !(replication != null &&
replication.getReplicationType() == EC) &&
srcKeyLen > getDatastreamMinLength()) {
perf.appendStreamMode();
- copyLength = ObjectEndpointStreaming
+ final CopyResult copyResult = ObjectEndpointStreaming
.copyKeyWithStream(volume.getBucket(destBucket), destKey, srcKeyLen,
getChunkSize(), replication, metadata, src, perf, startNanos,
tags,
writeConditions);
+ eTag = copyResult.getETag();
+ copyLength = copyResult.getSize();
+ modificationTime = copyResult.getModificationTime();
} else {
- try (OzoneOutputStream dest = openKeyForPut(
+ final OzoneOutputStream destStream = openKeyForPut(
volume.getName(), destBucket, destKey, srcKeyLen,
- replication, metadata, tags, writeConditions)) {
+ replication, metadata, tags, writeConditions);
+ try (OzoneOutputStream dest = destStream) {
long metadataLatencyNs =
getMetrics().updateCopyKeyMetadataStats(startNanos);
perf.appendMetaLatencyNanos(metadataLatencyNs);
copyLength = IOUtils.copyLarge(src, dest, 0, srcKeyLen, new
byte[getIOBufferSize(srcKeyLen)]);
- final String md5Hash =
DatatypeConverter.printHexBinary(src.getMessageDigest().digest()).toLowerCase();
- dest.getMetadata().put(OzoneConsts.ETAG, md5Hash);
+ eTag =
DatatypeConverter.printHexBinary(src.getMessageDigest().digest()).toLowerCase();
+ destStream.getMetadata().put(OzoneConsts.ETAG, eTag);
}
+ modificationTime = destStream.getModificationTime();
}
getMetrics().incCopyObjectSuccessLength(copyLength);
perf.appendSizeBytes(copyLength);
+ return new CopyResult(eTag, copyLength, modificationTime);
}
private CopyObjectResponse copyObject(OzoneVolume volume,
@@ -1116,19 +1141,14 @@ private CopyObjectResponse copyObject(OzoneVolume
volume,
"GetObject", () -> getClientProtocol().getKey(volume.getName(),
sourceBucket, sourceKey));
DigestInputStream sourceDigestInputStream = new
DigestInputStream(src, md5Digest)) {
getMetrics().updateCopyKeyMetadataStats(startNanos);
- runWithS3ActionString("PutObject", () -> {
- copy(volume, sourceDigestInputStream, sourceKeyLen, destkey,
destBucket,
- replicationConfig, customMetadata, perf, startNanos, tags,
writeConditions);
- return null;
- });
-
- final OzoneKeyDetails destKeyDetails =
getClientProtocol().getKeyDetails(
- volume.getName(), destBucket, destkey);
+ final CopyResult copyResult = runWithS3ActionString("PutObject", () ->
+ copy(volume, sourceDigestInputStream, sourceKeyLen, destkey,
destBucket,
+ replicationConfig, customMetadata, perf, startNanos, tags,
writeConditions));
getMetrics().updateCopyObjectSuccessStats(startNanos);
CopyObjectResponse copyObjectResponse = new CopyObjectResponse();
-
copyObjectResponse.setETag(wrapInQuotes(destKeyDetails.getMetadata().get(OzoneConsts.ETAG)));
-
copyObjectResponse.setLastModified(destKeyDetails.getModificationTime());
+ copyObjectResponse.setETag(wrapInQuotes(copyResult.getETag()));
+
copyObjectResponse.setLastModified(Instant.ofEpochMilli(copyResult.getModificationTime()));
return copyObjectResponse;
}
} catch (OMException ex) {
diff --git
a/hadoop-ozone/s3gateway/src/main/java/org/apache/hadoop/ozone/s3/endpoint/ObjectEndpointStreaming.java
b/hadoop-ozone/s3gateway/src/main/java/org/apache/hadoop/ozone/s3/endpoint/ObjectEndpointStreaming.java
index 9688c84c67b..cd018ef09f3 100644
---
a/hadoop-ozone/s3gateway/src/main/java/org/apache/hadoop/ozone/s3/endpoint/ObjectEndpointStreaming.java
+++
b/hadoop-ozone/s3gateway/src/main/java/org/apache/hadoop/ozone/s3/endpoint/ObjectEndpointStreaming.java
@@ -180,7 +180,7 @@ private static OzoneDataStreamOutput
openStreamKeyForPut(OzoneBucket bucket,
}
@SuppressWarnings("checkstyle:ParameterNumber")
- public static long copyKeyWithStream(
+ public static CopyResult copyKeyWithStream(
OzoneBucket bucket,
String keyPath,
long length,
@@ -192,18 +192,21 @@ public static long copyKeyWithStream(
S3ConditionalRequest.WriteConditions writeConditions)
throws IOException {
long writeLen;
- try (OzoneDataStreamOutput streamOutput = openStreamKeyForPut(bucket,
+ String eTag;
+ final OzoneDataStreamOutput streamOutput = openStreamKeyForPut(bucket,
keyPath, length, replicationConfig, keyMetadata, tags,
- writeConditions)) {
+ writeConditions);
+ try (OzoneDataStreamOutput stream = streamOutput) {
long metadataLatencyNs =
METRICS.updateCopyKeyMetadataStats(startNanos);
writeLen = writeToStreamOutput(streamOutput, body, bufferSize, length);
- String eTag =
DatatypeConverter.printHexBinary(body.getMessageDigest().digest())
+ eTag = DatatypeConverter.printHexBinary(body.getMessageDigest().digest())
.toLowerCase();
perf.appendMetaLatencyNanos(metadataLatencyNs);
- ((KeyMetadataAware)streamOutput).getMetadata().put(OzoneConsts.ETAG,
eTag);
+ streamOutput.getMetadata().put(OzoneConsts.ETAG, eTag);
}
- return writeLen;
+
+ return new CopyResult(eTag, writeLen, streamOutput.getModificationTime());
}
private static long writeToStreamOutput(OzoneDataStreamOutput streamOutput,
diff --git
a/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/client/OzoneBucketStub.java
b/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/client/OzoneBucketStub.java
index 8a8864c398b..af853b318e6 100644
---
a/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/client/OzoneBucketStub.java
+++
b/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/client/OzoneBucketStub.java
@@ -145,20 +145,21 @@ public OzoneOutputStream createKey(String key, long size,
new KeyMetadataAwareOutputStream(metadata) {
@Override
public void close() throws IOException {
+ super.close();
keyContents.put(key, toByteArray());
+ final long mtime = getModificationTime();
keyDetails.put(key, new OzoneKeyDetails(
getVolumeName(),
getName(),
key,
size,
- System.currentTimeMillis(),
- System.currentTimeMillis(),
+ mtime,
+ mtime,
new ArrayList<>(), finalReplicationCon, getMetadata(), null,
() -> readKey(key), true,
UserGroupInformation.getCurrentUser().getShortUserName(),
tags
));
- super.close();
}
};
@@ -179,18 +180,19 @@ public OzoneOutputStream rewriteKey(String keyName, long
size, long existingKeyG
new KeyMetadataAwareOutputStream(metadata) {
@Override
public void close() throws IOException {
+ super.close();
keyContents.put(keyName, toByteArray());
+ final long mtime = getModificationTime();
keyDetails.put(keyName, new OzoneKeyDetails(
getVolumeName(),
getName(),
keyName,
size,
- System.currentTimeMillis(),
- System.currentTimeMillis(),
+ mtime,
+ mtime,
new ArrayList<>(), finalReplicationCon, metadata, null,
() -> readKey(keyName), true, null, null
));
- super.close();
}
};
@@ -254,13 +256,14 @@ public void close() throws IOException {
Map<String, String> objectMetadata = keyMetadata == null ?
new HashMap<>() : keyMetadata;
+ final long mtime = getModificationTime();
keyDetails.put(key, new OzoneKeyDetails(
getVolumeName(),
getName(),
key,
size,
- System.currentTimeMillis(),
- System.currentTimeMillis(),
+ mtime,
+ mtime,
new ArrayList<>(), rConfig, objectMetadata, null,
null, false,
UserGroupInformation.getCurrentUser().getShortUserName(),
@@ -340,7 +343,7 @@ public void close() throws IOException {
buffer.get(bytes);
Part part = new Part(key + size, bytes,
- getMetadata().get(ETAG));
+ getMetadata().get(ETAG), getModificationTime());
if (partList.get(key) == null) {
Map<Integer, Part> parts = new TreeMap<>();
parts.put(partNumber, part);
@@ -518,8 +521,9 @@ public OzoneOutputStream createMultipartKey(String key,
long size,
new KeyMetadataAwareOutputStream((int) size, new HashMap<>()) {
@Override
public void close() throws IOException {
+ super.close();
Part part = new Part(key + size,
- toByteArray(), getMetadata().get(ETAG));
+ toByteArray(), getMetadata().get(ETAG),
getModificationTime());
if (partList.get(key) == null) {
Map<Integer, Part> parts = new TreeMap<>();
parts.put(partNumber, part);
@@ -527,7 +531,6 @@ public void close() throws IOException {
} else {
partList.get(key).put(partNumber, part);
}
- super.close();
}
};
return new OzoneOutputStreamStub(keyOutputStream, key + size);
@@ -649,7 +652,7 @@ public OzoneMultipartUploadPartListParts listParts(String
key,
if (partEntry.getKey() > partNumberMarker) {
PartInfo partInfo = new PartInfo(partEntry.getKey(),
partEntry.getValue().getPartName(),
- Time.now(), partEntry.getValue().getContent().length,
+ partEntry.getValue().getModificationTime(),
partEntry.getValue().getContent().length,
DatatypeConverter.printHexBinary(eTagProvider.digest(partEntry
.getValue().getContent())).toLowerCase());
partInfoList.add(partInfo);
@@ -732,13 +735,18 @@ public void deleteObjectTagging(String keyName) throws
IOException {
public static class Part {
private String partName;
private byte[] content;
-
private String eTag;
+ private long modificationTime;
- public Part(String name, byte[] data, String eTag) {
+ public Part(String name, byte[] data, String eTag, long modificationTime) {
this.partName = name;
this.content = data.clone();
this.eTag = eTag;
+ this.modificationTime = modificationTime;
+ }
+
+ public long getModificationTime() {
+ return modificationTime;
}
public String getPartName() {
@@ -798,6 +806,7 @@ public static class KeyMetadataAwareOutputStream extends
KeyOutputStream impleme
private final ByteArrayOutputStream buffer = new ByteArrayOutputStream();
private final Map<String, String> metadata;
private List<CheckedRunnable<IOException>> preCommits =
Collections.emptyList();
+ private long modificationTime;
public KeyMetadataAwareOutputStream(Map<String, String> metadata) {
super(null, null);
@@ -831,9 +840,15 @@ public void close() throws IOException {
for (CheckedRunnable<IOException> preCommit : preCommits) {
preCommit.run();
}
+ modificationTime = Time.now();
buffer.close();
}
+ @Override
+ public long getModificationTime() {
+ return modificationTime;
+ }
+
@Override
public void setPreCommits(List<CheckedRunnable<IOException>> preCommits) {
this.preCommits = preCommits != null ? preCommits :
Collections.emptyList();
@@ -859,6 +874,7 @@ public static class KeyMetadataAwareByteBufferStreamOutput
private final Map<String, String> metadata;
private List<CheckedRunnable<IOException>> preCommits =
Collections.emptyList();
+ private long modificationTime;
public KeyMetadataAwareByteBufferStreamOutput(
Map<String, String> metadata) {
@@ -878,10 +894,15 @@ public void flush() throws IOException {
@Override
public void close() throws IOException {
-
for (CheckedRunnable<IOException> preCommit : preCommits) {
preCommit.run();
}
+ modificationTime = Time.now();
+ }
+
+ @Override
+ public long getModificationTime() {
+ return modificationTime;
}
@Override
diff --git
a/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/client/OzoneDataStreamOutputStub.java
b/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/client/OzoneDataStreamOutputStub.java
index c19312692b4..28ebbe04bf8 100644
---
a/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/client/OzoneDataStreamOutputStub.java
+++
b/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/client/OzoneDataStreamOutputStub.java
@@ -63,8 +63,22 @@ public synchronized void close() throws IOException {
@Override
public OmMultipartCommitUploadPartInfo getCommitUploadPartInfo() {
- return closed ? new OmMultipartCommitUploadPartInfo(partName,
- getMetadata().get(OzoneConsts.ETAG)) : null;
+ final KeyDataStreamOutput keyDataStreamOutput = getKeyDataStreamOutput();
+ if (closed && keyDataStreamOutput != null) {
+ return new OmMultipartCommitUploadPartInfo(
+ partName, getMetadata().get(OzoneConsts.ETAG),
keyDataStreamOutput.getModificationTime());
+ }
+ return null;
+ }
+
+ @Override
+ public long getModificationTime() {
+ final KeyDataStreamOutput keyDataStreamOutput = getKeyDataStreamOutput();
+ if (keyDataStreamOutput != null) {
+ return keyDataStreamOutput.getModificationTime();
+ }
+ throw new IllegalStateException(
+ "OutputStream is not a KeyDataStreamOutput: " +
getByteBufStreamOutput().getClass());
}
@Override
diff --git
a/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/client/OzoneOutputStreamStub.java
b/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/client/OzoneOutputStreamStub.java
index ac3dc8b4da1..2661e98edc4 100644
---
a/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/client/OzoneOutputStreamStub.java
+++
b/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/client/OzoneOutputStreamStub.java
@@ -21,6 +21,7 @@
import java.io.OutputStream;
import org.apache.hadoop.ozone.OzoneConsts;
import org.apache.hadoop.ozone.client.io.KeyMetadataAware;
+import org.apache.hadoop.ozone.client.io.KeyOutputStream;
import org.apache.hadoop.ozone.client.io.OzoneOutputStream;
import org.apache.hadoop.ozone.om.helpers.OmMultipartCommitUploadPartInfo;
@@ -69,7 +70,22 @@ public synchronized void close() throws IOException {
@Override
public OmMultipartCommitUploadPartInfo getCommitUploadPartInfo() {
- return closed ? new OmMultipartCommitUploadPartInfo(partName,
-
((KeyMetadataAware)getOutputStream()).getMetadata().get(OzoneConsts.ETAG)) :
null;
+ final KeyOutputStream keyOutputStream = getKeyOutputStream();
+ if (closed && keyOutputStream != null) {
+ return new OmMultipartCommitUploadPartInfo(
+ partName, ((KeyMetadataAware)
getOutputStream()).getMetadata().get(OzoneConsts.ETAG),
+ keyOutputStream.getModificationTime());
+ }
+ return null;
+ }
+
+ @Override
+ public long getModificationTime() {
+ final KeyOutputStream keyOutputStream = getKeyOutputStream();
+ if (keyOutputStream != null) {
+ return keyOutputStream.getModificationTime();
+ }
+ throw new IllegalStateException(
+ "OutputStream is not a KeyOutputStream: " +
getOutputStream().getClass());
}
}
diff --git
a/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/s3/endpoint/TestCopyObjectAndUploadPartCopyLastModified.java
b/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/s3/endpoint/TestCopyObjectAndUploadPartCopyLastModified.java
new file mode 100644
index 00000000000..8de1e246904
--- /dev/null
+++
b/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/s3/endpoint/TestCopyObjectAndUploadPartCopyLastModified.java
@@ -0,0 +1,150 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hadoop.ozone.s3.endpoint;
+
+import static java.nio.charset.StandardCharsets.UTF_8;
+import static
org.apache.hadoop.ozone.s3.endpoint.EndpointTestUtils.initiateMultipartUpload;
+import static org.apache.hadoop.ozone.s3.endpoint.EndpointTestUtils.put;
+import static org.apache.hadoop.ozone.s3.util.S3Consts.COPY_SOURCE_HEADER;
+import static org.apache.hadoop.ozone.s3.util.S3Consts.STORAGE_CLASS_HEADER;
+import static org.apache.hadoop.ozone.s3.util.S3Consts.X_AMZ_CONTENT_SHA256;
+import static org.apache.hadoop.ozone.s3.util.S3Utils.urlEncode;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+import java.io.OutputStream;
+import java.util.HashMap;
+import javax.ws.rs.core.HttpHeaders;
+import javax.ws.rs.core.Response;
+import org.apache.hadoop.hdds.client.ReplicationConfig;
+import org.apache.hadoop.hdds.client.ReplicationFactor;
+import org.apache.hadoop.hdds.client.ReplicationType;
+import org.apache.hadoop.ozone.client.OzoneBucket;
+import org.apache.hadoop.ozone.client.OzoneClient;
+import org.apache.hadoop.ozone.client.OzoneMultipartUploadPartListParts;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+/**
+ * Verifies CopyObject and UploadPartCopy return LastModified from the commit
+ * response rather than the source key time or a synthetic gateway timestamp.
+ */
+public class TestCopyObjectAndUploadPartCopyLastModified {
+
+ private static final String SOURCE_BUCKET = "source-bucket";
+ private static final String DEST_BUCKET = "dest-bucket";
+ private static final String SOURCE_KEY = "source-key";
+ private static final String DEST_KEY = "dest-key";
+ private static final String MPU_KEY = "mpu-key";
+ private static final String CONTENT = "copy-last-modified-test-content";
+ private static final long SOURCE_AGE_MS = 100L;
+
+ private ObjectEndpoint endpoint;
+ private HttpHeaders headers;
+ private OzoneBucket sourceBucket;
+ private OzoneBucket destBucket;
+
+ @BeforeEach
+ void setUp() throws Exception {
+ headers = mock(HttpHeaders.class);
+
when(headers.getHeaderString(X_AMZ_CONTENT_SHA256)).thenReturn("UNSIGNED-PAYLOAD");
+ when(headers.getHeaderString(STORAGE_CLASS_HEADER)).thenReturn("STANDARD");
+
+ final OzoneClient client = EndpointBuilder.newObjectEndpointBuilder()
+ .setHeaders(headers)
+ .build()
+ .getClient();
+ client.getObjectStore().createS3Bucket(SOURCE_BUCKET);
+ client.getObjectStore().createS3Bucket(DEST_BUCKET);
+ sourceBucket = client.getObjectStore().getS3Bucket(SOURCE_BUCKET);
+ destBucket = client.getObjectStore().getS3Bucket(DEST_BUCKET);
+
+ createSourceKey();
+
+ endpoint = EndpointBuilder.newObjectEndpointBuilder()
+ .setHeaders(headers)
+ .setClient(client)
+ .build();
+ }
+
+ @Test
+ void testCopyObjectLastModifiedReflectsDestCommitTime() throws Exception {
+ final long sourceModificationTime = sourceBucket.getKey(SOURCE_KEY)
+ .getModificationTime().toEpochMilli();
+ Thread.sleep(SOURCE_AGE_MS);
+
+ when(headers.getHeaderString(COPY_SOURCE_HEADER)).thenReturn(
+ SOURCE_BUCKET + "/" + urlEncode(SOURCE_KEY));
+
+ try (Response response = put(endpoint, DEST_BUCKET, DEST_KEY, CONTENT)) {
+ assertEquals(200, response.getStatus());
+
+ final CopyObjectResponse copyObjectResponse = (CopyObjectResponse)
response.getEntity();
+ assertNotNull(copyObjectResponse.getLastModified());
+ assertNotNull(copyObjectResponse.getETag());
+
+ final long responseModificationTime =
copyObjectResponse.getLastModified().toEpochMilli();
+ final long destModificationTime =
destBucket.getKey(DEST_KEY).getModificationTime().toEpochMilli();
+
+ assertEquals(destModificationTime, responseModificationTime);
+ assertNotEquals(sourceModificationTime, responseModificationTime);
+ }
+ }
+
+ @Test
+ void testUploadPartCopyLastModifiedReflectsPartCommitTime() throws Exception
{
+ final long sourceModificationTime = sourceBucket.getKey(SOURCE_KEY)
+ .getModificationTime().toEpochMilli();
+ Thread.sleep(SOURCE_AGE_MS);
+
+ final String uploadId = initiateMultipartUpload(endpoint, DEST_BUCKET,
MPU_KEY);
+
+ when(headers.getHeaderString(COPY_SOURCE_HEADER)).thenReturn(SOURCE_BUCKET
+ "/" + urlEncode(SOURCE_KEY));
+
+ try (Response response = put(endpoint, DEST_BUCKET, MPU_KEY, 1, uploadId,
"")) {
+ assertEquals(200, response.getStatus());
+
+ final CopyPartResult copyPartResult = (CopyPartResult)
response.getEntity();
+ assertNotNull(copyPartResult.getETag());
+ assertNotNull(copyPartResult.getLastModified());
+
+ final long responseModificationTime =
copyPartResult.getLastModified().toEpochMilli();
+
+ // Retrieve the actual modification time of the newly uploaded part from
the bucket
+ final OzoneMultipartUploadPartListParts parts = destBucket.listParts(
+ MPU_KEY, uploadId, 0, 100);
+ final long destPartModificationTime =
parts.getPartInfoList().get(0).getModificationTime();
+
+ // Assert that the response precisely matches the part's actual commit
time
+ assertEquals(destPartModificationTime, responseModificationTime);
+ assertNotEquals(sourceModificationTime, responseModificationTime);
+ }
+ }
+
+ private void createSourceKey() throws Exception {
+ try (OutputStream stream = sourceBucket.createKey(
+ SOURCE_KEY, CONTENT.length(), ReplicationConfig.fromTypeAndFactor(
+ ReplicationType.RATIS, ReplicationFactor.THREE), new HashMap<>()))
{
+ stream.write(CONTENT.getBytes(UTF_8));
+ }
+ Thread.sleep(SOURCE_AGE_MS);
+ }
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]