This is an automated email from the ASF dual-hosted git repository.
spacemonkd pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ozone.git
The following commit(s) were added to refs/heads/master by this push:
new 4417f6b7de2 HDDS-14665. Add upgrade handling to multipart requests
(#10062)
4417f6b7de2 is described below
commit 4417f6b7de28dfda53330bdce4c51958d24aa037
Author: Abhishek Pal <[email protected]>
AuthorDate: Thu Jul 23 08:58:09 2026 +0530
HDDS-14665. Add upgrade handling to multipart requests (#10062)
---
.../ozone/om/helpers/OmMultipartKeyInfo.java | 30 +++--
.../ozone/om/helpers/OmMultipartPartInfo.java | 32 +++---
.../ozone/om/helpers/TestOmMultipartKeyInfo.java | 22 +++-
.../AbstractTestStorageDistributionEndpoint.java | 2 +
.../rpc/TestOzoneClientMultipartUploadWithFSO.java | 29 +++--
.../hadoop/ozone/om/snapshot/OmSnapshotTests.java | 21 ++++
.../src/main/proto/OmClientProtocol.proto | 1 +
.../S3InitiateMultipartUploadRequest.java | 39 ++++++-
.../S3InitiateMultipartUploadRequestWithFSO.java | 3 +
.../S3MultipartUploadCommitPartRequest.java | 18 +--
.../hadoop/ozone/om/upgrade/OMLayoutFeature.java | 3 +-
.../apache/hadoop/ozone/om/TestKeyManagerUnit.java | 6 +-
.../s3/multipart/S3MultipartRequestTests.java | 62 ++++++++++
.../TestS3InitiateMultipartUploadRequest.java | 108 ++++++++++++++++++
.../TestS3MultipartUploadAbortRequest.java | 88 +++++++++++++++
.../TestS3MultipartUploadCommitPartRequest.java | 125 +++++++++++++++++++--
.../TestS3MultipartUploadCompleteRequest.java | 77 +++++++++++++
.../hadoop/ozone/s3/endpoint/ObjectEndpoint.java | 26 +++++
18 files changed, 630 insertions(+), 62 deletions(-)
diff --git
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartKeyInfo.java
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartKeyInfo.java
index cfa92fa4103..7e170d5c073 100644
---
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartKeyInfo.java
+++
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartKeyInfo.java
@@ -41,8 +41,11 @@
* upload part information of the key.
*/
public final class OmMultipartKeyInfo extends WithObjectID implements
CopyObject<OmMultipartKeyInfo> {
- public static final byte LEGACY_SCHEMA_VERSION = 0;
- public static final byte SPLIT_PARTS_TABLE_SCHEMA_VERSION = 1;
+ // This stores the schema version of the multipart key.
+ // 0 - Legacy Schema -> Uses the same table to store the multipart part info
+ // 1 - New Schema -> Uses a separate table to store the multipart part info
+ public static final int LEGACY_SCHEMA_VERSION = 0;
+ public static final int SPLIT_PARTS_TABLE_SCHEMA_VERSION = 1;
private static final Codec<OmMultipartKeyInfo> CODEC = new DelegatedCodec<>(
Proto2Codec.get(MultipartKeyInfo.getDefaultInstance()),
@@ -87,10 +90,7 @@ public final class OmMultipartKeyInfo extends WithObjectID
implements CopyObject
*/
private final long parentID;
- // This stores the schema version of the multipart key.
- // 0 - Legacy Schema -> Uses the same table to store the multipart part info
- // 1 - New Schema -> Uses a separate table to store the multipart part info
- private final byte schemaVersion;
+ private final int schemaVersion;
public static Codec<OmMultipartKeyInfo> getCodec() {
return CODEC;
@@ -275,7 +275,7 @@ public ReplicationConfig getReplicationConfig() {
return replicationConfig;
}
- public byte getSchemaVersion() {
+ public int getSchemaVersion() {
return schemaVersion;
}
@@ -297,7 +297,7 @@ public static class Builder extends
WithObjectID.Builder<OmMultipartKeyInfo> {
private final AclListBuilder acls;
private final TreeMap<Integer, PartKeyInfo> partKeyInfoList;
private long parentID;
- private byte schemaVersion;
+ private int schemaVersion;
public Builder() {
this.acls = AclListBuilder.empty();
@@ -410,8 +410,8 @@ public Builder setParentID(long parentObjId) {
return this;
}
- public Builder setSchemaVersion(byte schemaVersion) {
- this.schemaVersion = schemaVersion;
+ public Builder setSchemaVersion(int schemaVersion) {
+ this.schemaVersion = validateAndConvertSchemaVersion(schemaVersion);
return this;
}
@@ -457,7 +457,7 @@ public static Builder builderFromProto(
.setObjectID(multipartKeyInfo.getObjectID())
.setUpdateID(multipartKeyInfo.getUpdateID())
.setParentID(multipartKeyInfo.getParentID())
- .setSchemaVersion((byte) multipartKeyInfo.getSchemaVersion());
+
.setSchemaVersion(validateAndConvertSchemaVersion(multipartKeyInfo.getSchemaVersion()));
}
/**
@@ -516,6 +516,14 @@ public MultipartKeyInfo getProto() {
return builder.build();
}
+ private static int validateAndConvertSchemaVersion(int schemaVersion) {
+ if (schemaVersion != LEGACY_SCHEMA_VERSION && schemaVersion !=
SPLIT_PARTS_TABLE_SCHEMA_VERSION) {
+ throw new IllegalArgumentException("Unsupported schemaVersion: "
+ + schemaVersion + ". Expected one of [0, 1].");
+ }
+ return schemaVersion;
+ }
+
@Override
public String getObjectInfo() {
return getProto().toString();
diff --git
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartPartInfo.java
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartPartInfo.java
index 19d8aa96845..9a00f0893b9 100644
---
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartPartInfo.java
+++
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartPartInfo.java
@@ -20,7 +20,6 @@
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
-import java.util.Objects;
import org.apache.commons.lang3.StringUtils;
import org.apache.hadoop.fs.FileChecksum;
import org.apache.hadoop.fs.FileEncryptionInfo;
@@ -70,6 +69,10 @@ private OmMultipartPartInfo(Builder b) {
if (b.partNumber <= 0) {
throw new IllegalArgumentException("partNumber is required and > 0");
}
+ // An ETag is MANDATORY for every multipart part stored in the split
parts-table schema,
+ // for ALL clients. The S3 gateway already computes the MD5 ETag, any
other client must also supply one.
+ // This mirrors the AWS S3 contract where UploadPart always yields an ETag
that CompleteMultipartUpload requires,
+ // and lets the ETag be used for part validation and returned by listParts.
if (StringUtils.isBlank(b.eTag)) {
throw new IllegalArgumentException("eTag is required");
}
@@ -158,8 +161,7 @@ public Builder setETag(String eTagValue) {
return this;
}
- public Builder setKeyLocationInfos(
- List<OmKeyLocationInfoGroup> keyLocationInfos) {
+ public Builder setKeyLocationInfos(List<OmKeyLocationInfoGroup>
keyLocationInfos) {
this.keyLocationInfos = new ArrayList<>(keyLocationInfos);
return this;
}
@@ -179,38 +181,33 @@ public OmMultipartPartInfo build() {
}
}
- public static OmMultipartPartInfo getFromProto(
- MultipartPartInfo multipartPartInfo) {
+ public static OmMultipartPartInfo getFromProto(MultipartPartInfo
multipartPartInfo) {
validateRequiredProtoFields(multipartPartInfo);
Builder builder = new Builder()
.setPartName(multipartPartInfo.getPartName())
.setPartNumber(multipartPartInfo.getPartNumber())
.setDataSize(multipartPartInfo.getDataSize())
.setModificationTime(multipartPartInfo.getModificationTime())
- .setETag(multipartPartInfo.getETag())
.setKeyLocationInfos(getKeyLocationInfosFromProto(multipartPartInfo))
+ .setETag(multipartPartInfo.getETag())
.setEncInfo(null);
if (!multipartPartInfo.hasObjectID()) {
- LOG.warn("MultipartPartInfo missing objectID for part {}",
- multipartPartInfo.getPartNumber());
+ LOG.warn("MultipartPartInfo missing objectID for part {}",
multipartPartInfo.getPartNumber());
}
builder.setObjectID(multipartPartInfo.getObjectID());
if (!multipartPartInfo.hasUpdateID()) {
- LOG.warn("MultipartPartInfo missing updateID for part {}",
- multipartPartInfo.getPartNumber());
+ LOG.warn("MultipartPartInfo missing updateID for part {}",
multipartPartInfo.getPartNumber());
}
builder.setUpdateID(multipartPartInfo.getUpdateID());
if (multipartPartInfo.hasFileEncryptionInfo()) {
- builder.setEncInfo(
- OMPBHelper.convert(multipartPartInfo.getFileEncryptionInfo()));
+
builder.setEncInfo(OMPBHelper.convert(multipartPartInfo.getFileEncryptionInfo()));
}
if (multipartPartInfo.hasFileChecksum()) {
- builder.setFileChecksum(
- OMPBHelper.convert(multipartPartInfo.getFileChecksum()));
+
builder.setFileChecksum(OMPBHelper.convert(multipartPartInfo.getFileChecksum()));
}
return builder.build();
@@ -232,6 +229,10 @@ public MultipartPartInfo getProto() {
if (keyLocationInfos == null || keyLocationInfos.isEmpty()) {
throw new IllegalArgumentException("keyLocationList is required");
}
+ if (StringUtils.isBlank(eTag)) {
+ throw new IllegalArgumentException("eTag is required");
+ }
+
MultipartPartInfo.Builder builder = MultipartPartInfo.newBuilder()
.setPartName(partName)
.setPartNumber(partNumber)
@@ -240,7 +241,7 @@ public MultipartPartInfo getProto() {
.setModificationTime(modificationTime)
.setObjectID(objectID)
.setUpdateID(updateID)
- .setETag(Objects.requireNonNull(eTag, "eTag is required"));
+ .setETag(eTag);
if (encInfo != null) {
builder.setFileEncryptionInfo(OMPBHelper.convert(encInfo));
@@ -355,6 +356,7 @@ private static void
validateRequiredProtoFields(MultipartPartInfo partInfo) {
if (!partInfo.hasPartNumber()) {
throw new IllegalArgumentException("MultipartPartInfo missing
partNumber");
}
+
if (!partInfo.hasETag() || StringUtils.isBlank(partInfo.getETag())) {
throw new IllegalArgumentException("MultipartPartInfo missing eTag");
}
diff --git
a/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/helpers/TestOmMultipartKeyInfo.java
b/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/helpers/TestOmMultipartKeyInfo.java
index 3a5658dd5bd..4033d94f455 100644
---
a/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/helpers/TestOmMultipartKeyInfo.java
+++
b/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/helpers/TestOmMultipartKeyInfo.java
@@ -120,7 +120,7 @@ public void distinctListOfParts() {
@Test
public void addPartKeyInfoRejectsSchemaVersionOne() {
OmMultipartKeyInfo subject = createSubject()
- .setSchemaVersion((byte) 1)
+ .setSchemaVersion(1)
.build();
assertThrows(IllegalStateException.class,
@@ -128,7 +128,7 @@ public void addPartKeyInfoRejectsSchemaVersionOne() {
}
@Test
- public void getProtoRejectsLegacyPartListForSchemaVersionOne() {
+ public void getProtoRejectsPartListForSchemaVersionOne() {
PartKeyInfo part = createPart(createKeyInfo()).build();
TreeMap<Integer, PartKeyInfo> legacyMap = new TreeMap<>();
legacyMap.put(part.getPartNumber(), part);
@@ -136,7 +136,7 @@ public void
getProtoRejectsLegacyPartListForSchemaVersionOne() {
OmMultipartKeyInfo subject = new OmMultipartKeyInfo.Builder()
.setUploadID(UUID.randomUUID().toString())
.setCreationTime(Time.now())
- .setSchemaVersion((byte) 1)
+ .setSchemaVersion(1)
.setReplicationConfig(StandaloneReplicationConfig.getInstance(
HddsProtos.ReplicationFactor.ONE))
.setPartKeyInfoList(legacyMap)
@@ -145,6 +145,22 @@ public void
getProtoRejectsLegacyPartListForSchemaVersionOne() {
assertThrows(IllegalStateException.class, subject::getProto);
}
+ @Test
+ public void builderFromProtoRejectsUnsupportedSchemaVersion() {
+ OmMultipartKeyInfo subject = createSubject()
+ .setReplicationConfig(StandaloneReplicationConfig.getInstance(
+ HddsProtos.ReplicationFactor.ONE))
+ .build();
+
+ OzoneManagerProtocolProtos.MultipartKeyInfo invalidProto =
subject.getProto()
+ .toBuilder()
+ .setSchemaVersion(256)
+ .build();
+
+ assertThrows(IllegalArgumentException.class,
+ () -> OmMultipartKeyInfo.getFromProto(invalidProto));
+ }
+
private static OmMultipartKeyInfo.Builder createSubject() {
return new OmMultipartKeyInfo.Builder()
.setUploadID(UUID.randomUUID().toString())
diff --git
a/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/AbstractTestStorageDistributionEndpoint.java
b/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/AbstractTestStorageDistributionEndpoint.java
index ee63904be59..8ea2d8c2125 100644
---
a/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/AbstractTestStorageDistributionEndpoint.java
+++
b/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/AbstractTestStorageDistributionEndpoint.java
@@ -50,6 +50,7 @@
import org.apache.hadoop.hdds.scm.server.StorageContainerManager;
import org.apache.hadoop.hdds.utils.IOUtils;
import org.apache.hadoop.ozone.MiniOzoneCluster;
+import org.apache.hadoop.ozone.OzoneConsts;
import org.apache.hadoop.ozone.client.BucketArgs;
import org.apache.hadoop.ozone.client.ObjectStore;
import org.apache.hadoop.ozone.client.OzoneBucket;
@@ -207,6 +208,7 @@ protected void createOpenKeysAndMultipartKeys(String
volumeName,
.createMultipartKey(volumeName, bucketName, "mpukey1",
100L, 1, multipartInfo.getUploadID());
partStream.write(new byte[100]);
+ partStream.getMetadata().put(OzoneConsts.ETAG, "mpukey1-part1-etag");
partStream.close();
}
diff --git
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestOzoneClientMultipartUploadWithFSO.java
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestOzoneClientMultipartUploadWithFSO.java
index 53ff4c1b9f5..ff0c6909461 100644
---
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestOzoneClientMultipartUploadWithFSO.java
+++
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestOzoneClientMultipartUploadWithFSO.java
@@ -82,9 +82,11 @@
import org.apache.hadoop.ozone.om.helpers.OmMultipartCommitUploadPartInfo;
import org.apache.hadoop.ozone.om.helpers.OmMultipartInfo;
import org.apache.hadoop.ozone.om.helpers.OmMultipartKeyInfo;
+import org.apache.hadoop.ozone.om.helpers.OmMultipartPartInfo;
import org.apache.hadoop.ozone.om.helpers.OmMultipartUploadCompleteInfo;
import org.apache.hadoop.ozone.om.helpers.OzoneFSUtils;
import org.apache.hadoop.ozone.om.helpers.QuotaUtil;
+import org.apache.hadoop.ozone.om.request.util.OMMultipartUploadUtils;
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos;
import org.apache.ozone.test.NonHATests;
import org.junit.jupiter.api.AfterAll;
@@ -657,14 +659,27 @@ private void verifyPartNamesInDB(Map<Integer, String>
partsMap,
metadataMgr.getMultipartInfoTable().get(multipartKey);
assertNotNull(omMultipartKeyInfo);
- for (OzoneManagerProtocolProtos.PartKeyInfo partKeyInfo :
- omMultipartKeyInfo.getPartKeyInfoMap()) {
- String partKeyName = partKeyInfo.getPartName();
-
- // reconstruct full part name with volume, bucket, partKeyName
- String fullKeyPartName =
- metadataMgr.getOzoneKey(volumeName, bucketName, keyName);
+ // Collect the part names as stored in the DB. For the split parts-table
+ // schema the parts live in the multipart parts table rather than inline
+ // in the multipart info table.
+ List<String> dbPartNames = new ArrayList<>();
+ if (omMultipartKeyInfo.getSchemaVersion()
+ == OmMultipartKeyInfo.SPLIT_PARTS_TABLE_SCHEMA_VERSION) {
+ for (OmMultipartPartInfo partInfo :
+ OMMultipartUploadUtils.scanParts(metadataMgr, uploadID).values()) {
+ dbPartNames.add(partInfo.getPartName());
+ }
+ } else {
+ for (OzoneManagerProtocolProtos.PartKeyInfo partKeyInfo :
+ omMultipartKeyInfo.getPartKeyInfoMap()) {
+ dbPartNames.add(partKeyInfo.getPartName());
+ }
+ }
+ // reconstruct full part name with volume, bucket, partKeyName
+ String fullKeyPartName =
+ metadataMgr.getOzoneKey(volumeName, bucketName, keyName);
+ for (String partKeyName : dbPartNames) {
// partKeyName format in DB - partKeyName + ClientID
assertTrue(partKeyName.startsWith(fullKeyPartName),
"Invalid partKeyName format in DB: " + partKeyName
diff --git
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/snapshot/OmSnapshotTests.java
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/snapshot/OmSnapshotTests.java
index 44131d81c88..d6115220ff9 100644
---
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/snapshot/OmSnapshotTests.java
+++
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/snapshot/OmSnapshotTests.java
@@ -88,6 +88,7 @@
import java.util.regex.Matcher;
import java.util.regex.Pattern;
import java.util.stream.Collectors;
+import org.apache.commons.codec.digest.DigestUtils;
import org.apache.commons.io.IOUtils;
import org.apache.commons.lang3.RandomStringUtils;
import org.apache.hadoop.fs.FSDataOutputStream;
@@ -3261,18 +3262,21 @@ public void testSnapshotDiffWithCreateMultipartKeys()
throws Exception {
try (OzoneOutputStream stream = bucket.createMultipartKey(
regularPartsKey, regularPart.length, 1, regularMpuInfo.getUploadID()))
{
stream.write(regularPart);
+ stream.getMetadata().put(OzoneConsts.ETAG,
DigestUtils.md5Hex(regularPart));
}
byte[] streamPart = "stream data".getBytes(UTF_8);
try (OzoneDataStreamOutput streamOut = bucket.createMultipartStreamKey(
streamPartsKey, streamPart.length, 1, streamMpuInfo.getUploadID())) {
streamOut.write(streamPart);
+ streamOut.getMetadata().put(OzoneConsts.ETAG,
DigestUtils.md5Hex(streamPart));
}
byte[] mixedPart = "mixed data".getBytes(UTF_8);
try (OzoneOutputStream mixedStream = bucket.createMultipartKey(
mixedPartsKey, mixedPart.length, 1, mixedMpuInfo.getUploadID())) {
mixedStream.write(mixedPart);
+ mixedStream.getMetadata().put(OzoneConsts.ETAG,
DigestUtils.md5Hex(mixedPart));
}
assertEquals(1,
@@ -3322,14 +3326,17 @@ public void testSnapshotDiffWithAbortMultipartUpload()
throws Exception {
try (OzoneOutputStream part1Stream = bucket.createMultipartKey(
partialAbortKey, part1Data.length, 1, partialInfo.getUploadID())) {
part1Stream.write(part1Data);
+ part1Stream.getMetadata().put(OzoneConsts.ETAG,
DigestUtils.md5Hex(part1Data));
}
try (OzoneOutputStream part2Stream = bucket.createMultipartKey(
partialAbortKey, part2Data.length, 2, partialInfo.getUploadID())) {
part2Stream.write(part2Data);
+ part2Stream.getMetadata().put(OzoneConsts.ETAG,
DigestUtils.md5Hex(part2Data));
}
try (OzoneDataStreamOutput part3Stream = bucket.createMultipartStreamKey(
partialAbortKey, part3Data.length, 3, partialInfo.getUploadID())) {
part3Stream.write(part3Data);
+ part3Stream.getMetadata().put(OzoneConsts.ETAG,
DigestUtils.md5Hex(part3Data));
}
OzoneMultipartUploadPartListParts partsList = bucket.listParts(
@@ -3348,10 +3355,12 @@ public void testSnapshotDiffWithAbortMultipartUpload()
throws Exception {
try (OzoneOutputStream stream = bucket.createMultipartKey(
multiAbortKey1, part1Data.length, 1, multiInfo1.getUploadID())) {
stream.write(part1Data);
+ stream.getMetadata().put(OzoneConsts.ETAG,
DigestUtils.md5Hex(part1Data));
}
try (OzoneDataStreamOutput stream = bucket.createMultipartStreamKey(
multiAbortKey2, part2Data.length, 1, multiInfo2.getUploadID())) {
stream.write(part2Data);
+ stream.getMetadata().put(OzoneConsts.ETAG,
DigestUtils.md5Hex(part2Data));
}
bucket.abortMultipartUpload(multiAbortKey1, multiInfo1.getUploadID());
@@ -3398,10 +3407,12 @@ public void
testSnapshotDiffWithCompleteInvisibleMPULifecycle() throws Exception
try (OzoneOutputStream stream = bucket.createMultipartKey(
mpuKey1, regularData1.length, 1, mpuInfo1.getUploadID())) {
stream.write(regularData1);
+ stream.getMetadata().put(OzoneConsts.ETAG,
DigestUtils.md5Hex(regularData1));
}
try (OzoneOutputStream stream = bucket.createMultipartKey(
mpuKey1, regularData2.length, 2, mpuInfo1.getUploadID())) {
stream.write(regularData2);
+ stream.getMetadata().put(OzoneConsts.ETAG,
DigestUtils.md5Hex(regularData2));
}
byte[] streamData1 = "Stream multipart data 1".getBytes(UTF_8);
@@ -3410,10 +3421,12 @@ public void
testSnapshotDiffWithCompleteInvisibleMPULifecycle() throws Exception
try (OzoneDataStreamOutput stream = bucket.createMultipartStreamKey(
mpuKey2, streamData1.length, 1, mpuInfo2.getUploadID())) {
stream.write(streamData1);
+ stream.getMetadata().put(OzoneConsts.ETAG,
DigestUtils.md5Hex(streamData1));
}
try (OzoneDataStreamOutput stream = bucket.createMultipartStreamKey(
mpuKey2, streamData2.length, 2, mpuInfo2.getUploadID())) {
stream.write(streamData2);
+ stream.getMetadata().put(OzoneConsts.ETAG,
DigestUtils.md5Hex(streamData2));
}
@@ -3423,10 +3436,12 @@ public void
testSnapshotDiffWithCompleteInvisibleMPULifecycle() throws Exception
try (OzoneOutputStream stream = bucket.createMultipartKey(
mpuKey3, mixedRegular.length, 1, mpuInfo3.getUploadID())) {
stream.write(mixedRegular);
+ stream.getMetadata().put(OzoneConsts.ETAG,
DigestUtils.md5Hex(mixedRegular));
}
try (OzoneDataStreamOutput stream = bucket.createMultipartStreamKey(
mpuKey3, mixedStream.length, 2, mpuInfo3.getUploadID())) {
stream.write(mixedStream);
+ stream.getMetadata().put(OzoneConsts.ETAG,
DigestUtils.md5Hex(mixedStream));
}
assertEquals(2,
@@ -3470,6 +3485,7 @@ private void completeSinglePartMPU(OzoneBucket bucket,
String keyName, String da
byte[] partData = createLargePartData(data, MIN_PART_SIZE);
OzoneOutputStream partStream = bucket.createMultipartKey(keyName,
partData.length, 1, uploadId);
partStream.write(partData);
+ partStream.getMetadata().put(OzoneConsts.ETAG,
DigestUtils.md5Hex(partData));
partStream.close();
OzoneMultipartUploadPartListParts partsList = bucket.listParts(keyName,
uploadId, 0, 100);
@@ -3491,6 +3507,7 @@ private void completeMultiplePartMPU(
try (OzoneOutputStream partStream = bucket.createMultipartKey(
keyName, partData.length, partNum, uploadId)) {
partStream.write(partData);
+ partStream.getMetadata().put(OzoneConsts.ETAG,
DigestUtils.md5Hex(partData));
}
}
@@ -3514,12 +3531,14 @@ private void completeMixedPartMPU(
try (OzoneOutputStream partStream = bucket.createMultipartKey(
keyName, part1Data.length, 1, uploadId)) {
partStream.write(part1Data);
+ partStream.getMetadata().put(OzoneConsts.ETAG,
DigestUtils.md5Hex(part1Data));
}
byte[] part2Data = createLargePartData(streamData, MIN_PART_SIZE);
try (OzoneDataStreamOutput partStream = bucket.createMultipartStreamKey(
keyName, part2Data.length, 2, uploadId)) {
partStream.write(part2Data);
+ partStream.getMetadata().put(OzoneConsts.ETAG,
DigestUtils.md5Hex(part2Data));
}
OzoneMultipartUploadPartListParts partsList = bucket.listParts(keyName,
uploadId, 0, 2);
@@ -3547,6 +3566,7 @@ private void completeMPUWithReplication(
try (OzoneOutputStream partStream = bucket.createMultipartKey(
keyName, partData.length, 1, uploadId)) {
partStream.write(partData);
+ partStream.getMetadata().put(OzoneConsts.ETAG,
DigestUtils.md5Hex(partData));
}
OzoneMultipartUploadPartListParts partsList = bucket.listParts(keyName,
uploadId, 0, 1);
@@ -3566,6 +3586,7 @@ private void completeMPUWithMetadata(OzoneBucket bucket,
String keyName,
byte[] partData = createLargePartData("MPU with metadata and tags",
MIN_PART_SIZE);
OzoneOutputStream partStream = bucket.createMultipartKey(keyName,
partData.length, 1, uploadId);
partStream.write(partData);
+ partStream.getMetadata().put(OzoneConsts.ETAG,
DigestUtils.md5Hex(partData));
partStream.close();
OzoneMultipartUploadPartListParts partsList = bucket.listParts(keyName,
uploadId, 0, 100);
diff --git
a/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto
b/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto
index f09020b7975..bccc6d11f11 100644
--- a/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto
+++ b/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto
@@ -1770,6 +1770,7 @@ message ServiceInfo {
message MultipartInfoInitiateRequest {
required KeyArgs keyArgs = 1;
+ optional uint32 schemaVersion = 2;
}
diff --git
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3InitiateMultipartUploadRequest.java
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3InitiateMultipartUploadRequest.java
index 22f470c80c7..72637c1a377 100644
---
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3InitiateMultipartUploadRequest.java
+++
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3InitiateMultipartUploadRequest.java
@@ -100,10 +100,14 @@ public OMRequest preExecute(OzoneManager ozoneManager)
throws IOException {
KeyArgs resolvedArgs = resolveBucketAndCheckKeyAcls(newKeyArgs.build(),
ozoneManager, ACLType.CREATE);
+ int schemaVersion = resolveMultipartSchemaVersion(ozoneManager);
+ MultipartInfoInitiateRequest.Builder requestBuilder =
+ multipartInfoInitiateRequest.toBuilder()
+ .setKeyArgs(resolvedArgs)
+ .setSchemaVersion(schemaVersion);
return getOmRequest().toBuilder()
.setUserInfo(getUserInfo())
- .setInitiateMultiPartUploadRequest(
- multipartInfoInitiateRequest.toBuilder().setKeyArgs(resolvedArgs))
+ .setInitiateMultiPartUploadRequest(requestBuilder)
.build();
}
@@ -195,6 +199,9 @@ public OMClientResponse validateAndUpdateCache(OzoneManager
ozoneManager, Execut
replicationConfig)
.setObjectID(objectID)
.setUpdateID(transactionLogIndex)
+ // Source of truth is the value stamped onto the proto in preExecute
+ // (before Ratis). Never re-check MLV here in the replicated apply
path.
+ .setSchemaVersion(multipartInfoInitiateRequest.getSchemaVersion())
.build();
omKeyInfo = new OmKeyInfo.Builder()
@@ -284,6 +291,34 @@ protected void logResult(OzoneManager ozoneManager,
}
}
+ /**
+ * Resolve the schema version stamped onto a newly initiated multipart
upload.
+ * <p>
+ * This is a server-authoritative decision made once, in {@code preExecute}
+ * (i.e. on the leader, before the request is submitted to Ratis). Any
+ * client-supplied {@code schemaVersion} on the request is intentionally
+ * ignored so that a client can never force the OM to persist an on-disk
+ * format the cluster is not ready for.
+ * <p>
+ * The split parts-table on-disk format is gated on the
+ * {@link OMLayoutFeature#MPU_PARTS_TABLE_SPLIT} layout feature:
+ * <ul>
+ * <li>pre-finalized (or mixed-binary rolling upgrade) → legacy
schema,
+ * so no split-table rows are written on a cluster that may still be
+ * downgraded;</li>
+ * <li>finalized → split parts-table schema.</li>
+ * </ul>
+ * Because finalization is replicated through Ratis, the leader's view here
is
+ * consistent across the quorum, and the stamped value (not the live layout
+ * version) is what all subsequent processing obeys.
+ */
+ protected int resolveMultipartSchemaVersion(OzoneManager ozoneManager) {
+ return ozoneManager.getVersionManager()
+ .isAllowed(OMLayoutFeature.MPU_PARTS_TABLE_SPLIT)
+ ? OmMultipartKeyInfo.SPLIT_PARTS_TABLE_SCHEMA_VERSION
+ : OmMultipartKeyInfo.LEGACY_SCHEMA_VERSION;
+ }
+
@RequestFeatureValidator(
conditions = ValidationCondition.CLUSTER_NEEDS_FINALIZATION,
processingPhase = RequestProcessingPhase.PRE_PROCESS,
diff --git
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3InitiateMultipartUploadRequestWithFSO.java
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3InitiateMultipartUploadRequestWithFSO.java
index 919491d7049..5596be4ff39 100644
---
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3InitiateMultipartUploadRequestWithFSO.java
+++
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3InitiateMultipartUploadRequestWithFSO.java
@@ -168,6 +168,9 @@ public OMClientResponse validateAndUpdateCache(OzoneManager
ozoneManager, Execut
.setObjectID(pathInfoFSO.getLeafNodeObjectId())
.setUpdateID(transactionLogIndex)
.setParentID(pathInfoFSO.getLastKnownParentId())
+ // Source of truth is the value stamped onto the proto in preExecute
+ // (before Ratis). Never re-check MLV here in the replicated apply
path.
+ .setSchemaVersion(multipartInfoInitiateRequest.getSchemaVersion())
.build();
omKeyInfo = new OmKeyInfo.Builder()
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 78f1af96cfb..24fe698336a 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
@@ -223,7 +223,15 @@ public OMClientResponse
validateAndUpdateCache(OzoneManager ozoneManager, Execut
// Add this part information in to multipartKeyInfo.
multipartKeyInfo.addPartKeyInfo(partKeyInfo.build());
} else {
- validateSplitPartInfo(omKeyInfo, partNumber);
+ // an ETag is MANDATORY for every committed part in the split
parts-table schema,
+ // enforced server-side for ALL clients (S3 gateway and native Ozone
client alike).
+ // The S3 gateway computes the MD5 ETag on upload; any other client
must also supply one.
+ // Reject the commit early with a clear INVALID_REQUEST if it is
missing,
+ if
(StringUtils.isBlank(omKeyInfo.getMetadata().get(OzoneConsts.ETAG))) {
+ throw new OMException(
+ "Missing ETag for multipart upload part " + partNumber,
+ OMException.ResultCodes.INVALID_REQUEST);
+ }
multipartPartInfo = OmMultipartPartInfo.from(partName, partNumber,
omKeyInfo);
omMetadataManager.getMultipartPartsTable().addCacheEntry(
new CacheKey<>(multipartPartKey),
@@ -410,14 +418,6 @@ private String getMultipartKey(String volumeName, String
bucketName,
keyName, uploadID);
}
- private void validateSplitPartInfo(OmKeyInfo omKeyInfo, int partNumber)
- throws OMException {
- if (StringUtils.isBlank(omKeyInfo.getMetadata().get(OzoneConsts.ETAG))) {
- throw new OMException("Missing ETag for multipart upload part "
- + partNumber, OMException.ResultCodes.INVALID_REQUEST);
- }
- }
-
@RequestFeatureValidator(
conditions = ValidationCondition.CLUSTER_NEEDS_FINALIZATION,
processingPhase = RequestProcessingPhase.PRE_PROCESS,
diff --git
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/upgrade/OMLayoutFeature.java
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/upgrade/OMLayoutFeature.java
index 790ab5845a5..189bb3f3e1d 100644
---
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/upgrade/OMLayoutFeature.java
+++
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/upgrade/OMLayoutFeature.java
@@ -45,7 +45,8 @@ public enum OMLayoutFeature implements LayoutFeature {
HBASE_SUPPORT(7, "Full support of hsync, lease recovery and listOpenFiles
APIs for HBase"),
DELEGATION_TOKEN_SYMMETRIC_SIGN(8, "Delegation token signed by symmetric
key"),
SNAPSHOT_DEFRAG(9, "Supporting defragmentation of snapshot"),
- S3_LIFECYCLE_SUPPORT(10, "S3 bucket lifecycle configuration support");
+ S3_LIFECYCLE_SUPPORT(10, "S3 bucket lifecycle configuration support"),
+ MPU_PARTS_TABLE_SPLIT(11, "Split multipart table into separate table for
parts and key");
/////////////////////////////// /////////////////////////////
diff --git
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/TestKeyManagerUnit.java
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/TestKeyManagerUnit.java
index 9b184421207..be307c06f95 100644
---
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/TestKeyManagerUnit.java
+++
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/TestKeyManagerUnit.java
@@ -159,8 +159,7 @@ public void listMultipartUploadPartsWithZeroUpload() throws
IOException {
}
@Test
- public void listMultipartUploadPartsWithoutEtagField() throws IOException {
- // For backward compatibility reasons
+ public void listMultipartUploadPartsWithEtagField() throws IOException {
final String volume = volumeName();
final String bucket = "bucketForEtag";
final String key = "dir/key1";
@@ -169,7 +168,7 @@ public void listMultipartUploadPartsWithoutEtagField()
throws IOException {
initMultipartUpload(writeClient, volume, bucket, key);
- // Commit some MPU parts without eTag field
+ // Commit some MPU parts, each carrying its (now mandatory) eTag.
for (int i = 1; i <= 5; i++) {
OmKeyArgs partKeyArgs =
new OmKeyArgs.Builder()
@@ -199,6 +198,7 @@ public void listMultipartUploadPartsWithoutEtagField()
throws IOException {
.setReplicationConfig(
RatisReplicationConfig.getInstance(ReplicationFactor.THREE))
.setLocationInfoList(Collections.emptyList())
+ .addMetadata(OzoneConsts.ETAG, "etag-" + i)
.build();
writeClient.commitMultipartUploadPart(commitPartKeyArgs,
openKey.getId());
diff --git
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartRequestTests.java
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartRequestTests.java
index e23b84f5293..08e8487ee08 100644
---
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartRequestTests.java
+++
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartRequestTests.java
@@ -34,6 +34,8 @@
import java.util.Map;
import org.apache.hadoop.hdds.client.ReplicationConfig;
import org.apache.hadoop.hdds.conf.OzoneConfiguration;
+import org.apache.hadoop.hdds.utils.db.cache.CacheKey;
+import org.apache.hadoop.hdds.utils.db.cache.CacheValue;
import org.apache.hadoop.ozone.audit.AuditLogger;
import org.apache.hadoop.ozone.audit.AuditMessage;
import org.apache.hadoop.ozone.om.IOmMetadataReader;
@@ -47,8 +49,11 @@
import org.apache.hadoop.ozone.om.ResolvedBucket;
import org.apache.hadoop.ozone.om.helpers.BucketLayout;
import org.apache.hadoop.ozone.om.helpers.KeyValueUtil;
+import org.apache.hadoop.ozone.om.helpers.OmMultipartKeyInfo;
import org.apache.hadoop.ozone.om.request.OMClientRequest;
import org.apache.hadoop.ozone.om.request.OMRequestTestUtils;
+import org.apache.hadoop.ozone.om.response.OMClientResponse;
+import org.apache.hadoop.ozone.om.upgrade.OMLayoutFeature;
import org.apache.hadoop.ozone.om.upgrade.OMLayoutVersionManager;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.KeyArgs;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.KeyLocation;
@@ -108,6 +113,11 @@ public void setup() throws Exception {
args.getVolumeName(), args.getBucketName(),
"owner", BucketLayout.DEFAULT);
});
+ // MPU request tests default to a pre-finalized layout version, i.e. a
+ // cluster that has not finalized MPU_PARTS_TABLE_SPLIT. Newly initiated
+ // uploads therefore use the legacy (schema 0) inline parts layout. Tests
+ // that need the finalized behaviour (split parts table, schema 1) opt in
+ // by calling finalizeMpuPartsTableSplit().
OMLayoutVersionManager lvm = mock(OMLayoutVersionManager.class);
when(lvm.getMetadataLayoutVersion()).thenReturn(0);
when(ozoneManager.getVersionManager()).thenReturn(lvm);
@@ -115,6 +125,18 @@ public void setup() throws Exception {
when(ozoneManager.getConfig()).thenReturn(ozoneConfiguration.getObject(OmConfig.class));
}
+ /**
+ * Simulate a cluster that has finalized the multipart parts-table split
+ * layout feature, so newly initiated uploads resolve to the split (schema 1)
+ * parts-table layout. This stubs the exact signal the request path checks
+ * ({@code isAllowed(MPU_PARTS_TABLE_SPLIT)}) rather than the metadata layout
+ * version, which the MPU schema gate does not read directly.
+ */
+ protected void finalizeMpuPartsTableSplit() {
+ when(ozoneManager.getVersionManager()
+ .isAllowed(OMLayoutFeature.MPU_PARTS_TABLE_SPLIT)).thenReturn(true);
+ }
+
@AfterEach
public void stop() {
omMetrics.unRegister();
@@ -302,6 +324,46 @@ protected OMRequest doPreExecuteCompleteMPU(
}
+ /**
+ * Initiate an MPU and optionally rewrite the stored multipart metadata to a
+ * specific schema version.
+ *
+ * <p>The schema version rewrite lets tests emulate post-finalization MPU
+ * entries without needing the rest of the upgrade pipeline.</p>
+ */
+ protected String initiateMultipartUploadWithSchemaVersion(
+ String volumeName, String bucketName, String keyName,
+ int schemaVersion) throws Exception {
+ OMRequest initiateMPURequest =
+ doPreExecuteInitiateMPU(volumeName, bucketName, keyName);
+
+ S3InitiateMultipartUploadRequest s3InitiateMultipartUploadRequest =
+ getS3InitiateMultipartUploadReq(initiateMPURequest);
+
+ OMClientResponse omClientResponse =
+ s3InitiateMultipartUploadRequest.validateAndUpdateCache(ozoneManager,
+ 1L);
+
+ String multipartUploadID = omClientResponse.getOMResponse()
+ .getInitiateMultiPartUploadResponse().getMultipartUploadID();
+
+ if (schemaVersion != 0) {
+ String multipartKey = omMetadataManager.getMultipartKey(volumeName,
+ bucketName, keyName, multipartUploadID);
+ OmMultipartKeyInfo multipartKeyInfo = omMetadataManager
+ .getMultipartInfoTable().get(multipartKey);
+ assertNotNull(multipartKeyInfo);
+
+ omMetadataManager.getMultipartInfoTable().addCacheEntry(
+ new CacheKey<>(multipartKey),
+ CacheValue.get(2L, multipartKeyInfo.toBuilder()
+ .setSchemaVersion(schemaVersion)
+ .build()));
+ }
+
+ return multipartUploadID;
+ }
+
/**
* Perform preExecute of Initiate Multipart upload request for given
* volume, bucket and key name.
diff --git
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3InitiateMultipartUploadRequest.java
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3InitiateMultipartUploadRequest.java
index 79cf8743e3d..e316fd4e248 100644
---
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3InitiateMultipartUploadRequest.java
+++
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3InitiateMultipartUploadRequest.java
@@ -24,6 +24,7 @@
import static org.junit.jupiter.api.Assertions.assertTrue;
import java.util.ArrayList;
+import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@@ -32,11 +33,14 @@
import org.apache.hadoop.ozone.OzoneAcl;
import org.apache.hadoop.ozone.om.helpers.OmBucketInfo;
import org.apache.hadoop.ozone.om.helpers.OmKeyInfo;
+import org.apache.hadoop.ozone.om.helpers.OmMultipartKeyInfo;
import org.apache.hadoop.ozone.om.request.OMRequestTestUtils;
import org.apache.hadoop.ozone.om.response.OMClientResponse;
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.OMRequest;
import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.CsvSource;
/**
* Tests S3 Initiate Multipart Upload request.
@@ -113,6 +117,110 @@ public void testValidateAndUpdateCache() throws Exception
{
}
+ /**
+ * The schema version is a server-owned decision resolved in {@code
preExecute}
+ * (on the leader) from the finalized layout version; {@code
+ * validateAndUpdateCache} only forwards the stamped value into the persisted
+ * multipart info row. Pre-finalization -> legacy (0); finalized -> split
(1).
+ */
+ @ParameterizedTest
+ @CsvSource({
+ // finalized, expectedSchemaVersion (0 = LEGACY, 1 = SPLIT_PARTS_TABLE)
+ "false, 0",
+ "true, 1",
+ })
+ public void testSchemaVersionStampedInPreExecuteByServer(
+ boolean finalized, int expectedSchemaVersion) throws Exception {
+ String volumeName = UUID.randomUUID().toString();
+ String bucketName = UUID.randomUUID().toString();
+ String keyName = UUID.randomUUID().toString();
+
+ if (finalized) {
+ finalizeMpuPartsTableSplit();
+ }
+
+ OMRequestTestUtils.addVolumeAndBucketToDB(volumeName, bucketName,
+ omMetadataManager, getBucketLayout());
+
+ // preExecute (leader) stamps the schema version onto the request.
+ OMRequest modifiedRequest = doPreExecuteInitiateMPU(volumeName,
+ bucketName, keyName);
+ assertEquals(expectedSchemaVersion,
+
modifiedRequest.getInitiateMultiPartUploadRequest().getSchemaVersion());
+
+ // validateAndUpdateCache only forwards the already-stamped value.
+ OMClientResponse response =
getS3InitiateMultipartUploadReq(modifiedRequest)
+ .validateAndUpdateCache(ozoneManager, 100L);
+ assertEquals(OzoneManagerProtocolProtos.Status.OK,
+ response.getOMResponse().getStatus());
+
+ String multipartKey = getMultipartKey(volumeName, bucketName, keyName,
+ modifiedRequest.getInitiateMultiPartUploadRequest()
+ .getKeyArgs().getMultipartUploadID());
+ OmMultipartKeyInfo multipartKeyInfo = omMetadataManager
+ .getMultipartInfoTable().get(multipartKey);
+ assertNotNull(multipartKeyInfo);
+ assertEquals(expectedSchemaVersion, multipartKeyInfo.getSchemaVersion());
+ }
+
+ /**
+ * The schema version is server-owned: a client-supplied value on the request
+ * is ignored and overwritten in {@code preExecute} with the server decision,
+ * in both directions (client asks for split pre-finalization, and client
asks
+ * for legacy post-finalization).
+ */
+ @ParameterizedTest
+ @CsvSource({
+ // finalized, expectedSchemaVersion (0 = LEGACY, 1 = SPLIT_PARTS_TABLE)
+ "false, 0",
+ "true, 1",
+ })
+ public void testServerIgnoresClientSuppliedSchemaVersion(
+ boolean finalized, int expectedSchemaVersion) throws Exception {
+ String volumeName = UUID.randomUUID().toString();
+ String bucketName = UUID.randomUUID().toString();
+ String keyName = UUID.randomUUID().toString();
+
+ if (finalized) {
+ finalizeMpuPartsTableSplit();
+ }
+
+ OMRequestTestUtils.addVolumeAndBucketToDB(volumeName, bucketName,
+ omMetadataManager, getBucketLayout());
+
+ // Client supplies the opposite of what the server should decide.
+ int clientSuppliedSchemaVersion = finalized
+ ? OmMultipartKeyInfo.LEGACY_SCHEMA_VERSION
+ : OmMultipartKeyInfo.SPLIT_PARTS_TABLE_SCHEMA_VERSION;
+ OMRequest clientRequest = OMRequestTestUtils.createInitiateMPURequest(
+ volumeName, bucketName, keyName, Collections.emptyMap(),
+ Collections.emptyMap());
+ clientRequest = clientRequest.toBuilder()
+ .setInitiateMultiPartUploadRequest(
+ clientRequest.getInitiateMultiPartUploadRequest().toBuilder()
+ .setSchemaVersion(clientSuppliedSchemaVersion))
+ .build();
+
+ // preExecute must overwrite the client value with the server decision.
+ OMRequest modifiedRequest =
+
getS3InitiateMultipartUploadReq(clientRequest).preExecute(ozoneManager);
+ assertEquals(expectedSchemaVersion,
+
modifiedRequest.getInitiateMultiPartUploadRequest().getSchemaVersion());
+
+ OMClientResponse response =
getS3InitiateMultipartUploadReq(modifiedRequest)
+ .validateAndUpdateCache(ozoneManager, 100L);
+ assertEquals(OzoneManagerProtocolProtos.Status.OK,
+ response.getOMResponse().getStatus());
+
+ String multipartKey = getMultipartKey(volumeName, bucketName, keyName,
+ modifiedRequest.getInitiateMultiPartUploadRequest()
+ .getKeyArgs().getMultipartUploadID());
+ OmMultipartKeyInfo multipartKeyInfo = omMetadataManager
+ .getMultipartInfoTable().get(multipartKey);
+ assertNotNull(multipartKeyInfo);
+ assertEquals(expectedSchemaVersion, multipartKeyInfo.getSchemaVersion());
+ }
+
@Test
public void testValidateAndUpdateCacheWithBucketNotFound() throws Exception {
String volumeName = UUID.randomUUID().toString();
diff --git
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3MultipartUploadAbortRequest.java
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3MultipartUploadAbortRequest.java
index 53ae47f5756..3d529ad0aca 100644
---
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3MultipartUploadAbortRequest.java
+++
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3MultipartUploadAbortRequest.java
@@ -25,6 +25,7 @@
import java.util.UUID;
import org.apache.hadoop.hdds.utils.db.cache.CacheKey;
import org.apache.hadoop.hdds.utils.db.cache.CacheValue;
+import org.apache.hadoop.ozone.om.helpers.OmMultipartKeyInfo;
import org.apache.hadoop.ozone.om.request.OMRequestTestUtils;
import org.apache.hadoop.ozone.om.response.OMClientResponse;
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos;
@@ -95,6 +96,93 @@ public void testValidateAndUpdateCache() throws Exception {
}
+ @Test
+ public void
testValidateAndUpdateCacheUsesSchemaVersionOneBeforeFinalization()
+ throws Exception {
+ String volumeName = UUID.randomUUID().toString();
+ String bucketName = UUID.randomUUID().toString();
+ String keyName = getKeyName();
+
+ // The base test fixture is pre-finalized by default.
+ OMRequestTestUtils.addVolumeAndBucketToDB(volumeName, bucketName,
+ omMetadataManager, getBucketLayout());
+
+ createParentPath(volumeName, bucketName);
+
+ String multipartUploadID =
+ initiateMultipartUploadWithSchemaVersion(volumeName, bucketName,
+ keyName, 1);
+
+ String multipartKey = omMetadataManager.getMultipartKey(volumeName,
+ bucketName, keyName, multipartUploadID);
+ OmMultipartKeyInfo multipartKeyInfo = omMetadataManager
+ .getMultipartInfoTable().get(multipartKey);
+ assertNotNull(multipartKeyInfo);
+ assertEquals(1, multipartKeyInfo.getSchemaVersion());
+
+ OMRequest abortMPURequest =
+ doPreExecuteAbortMPU(volumeName, bucketName, keyName,
+ multipartUploadID);
+
+ S3MultipartUploadAbortRequest s3MultipartUploadAbortRequest =
+ getS3MultipartUploadAbortReq(abortMPURequest);
+
+ OMClientResponse omClientResponse =
+ s3MultipartUploadAbortRequest.validateAndUpdateCache(ozoneManager, 2L);
+
+ assertEquals(OzoneManagerProtocolProtos.Status.OK,
+ omClientResponse.getOMResponse().getStatus());
+ }
+
+ @Test
+ public void
testValidateAndUpdateCacheAllowsSchemaVersionZeroAfterFinalization()
+ throws Exception {
+ String volumeName = UUID.randomUUID().toString();
+ String bucketName = UUID.randomUUID().toString();
+ String keyName = getKeyName();
+
+ OMRequestTestUtils.addVolumeAndBucketToDB(volumeName, bucketName,
+ omMetadataManager, getBucketLayout());
+
+ createParentPath(volumeName, bucketName);
+
+ // Upload is initiated on a pre-finalized cluster, so it uses the legacy
+ // (schema 0) inline layout.
+ OMRequest initiateMPURequest = doPreExecuteInitiateMPU(volumeName,
+ bucketName, keyName);
+ S3InitiateMultipartUploadRequest s3InitiateMultipartUploadRequest =
+ getS3InitiateMultipartUploadReq(initiateMPURequest);
+ OMClientResponse initiateResponse =
+ s3InitiateMultipartUploadRequest.validateAndUpdateCache(ozoneManager,
+ 1L);
+ String multipartUploadID = initiateResponse.getOMResponse()
+ .getInitiateMultiPartUploadResponse().getMultipartUploadID();
+
+ String multipartKey = omMetadataManager.getMultipartKey(volumeName,
+ bucketName, keyName, multipartUploadID);
+ OmMultipartKeyInfo multipartKeyInfo = omMetadataManager
+ .getMultipartInfoTable().get(multipartKey);
+ assertNotNull(multipartKeyInfo);
+ assertEquals(0, multipartKeyInfo.getSchemaVersion());
+
+ // Cluster finalizes the split feature, the pre-existing legacy upload must
+ // still be abortable.
+ finalizeMpuPartsTableSplit();
+
+ OMRequest abortMPURequest =
+ doPreExecuteAbortMPU(volumeName, bucketName, keyName,
+ multipartUploadID);
+
+ S3MultipartUploadAbortRequest s3MultipartUploadAbortRequest =
+ getS3MultipartUploadAbortReq(abortMPURequest);
+
+ OMClientResponse omClientResponse =
+ s3MultipartUploadAbortRequest.validateAndUpdateCache(ozoneManager, 2L);
+
+ assertEquals(OzoneManagerProtocolProtos.Status.OK,
+ omClientResponse.getOMResponse().getStatus());
+ }
+
@Test
public void testValidateAndUpdateCacheMultipartNotFound() throws Exception {
String volumeName = UUID.randomUUID().toString();
diff --git
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3MultipartUploadCommitPartRequest.java
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3MultipartUploadCommitPartRequest.java
index b4513008c2e..e9ccbb4d314 100644
---
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3MultipartUploadCommitPartRequest.java
+++
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3MultipartUploadCommitPartRequest.java
@@ -39,7 +39,6 @@
import org.apache.hadoop.hdds.utils.db.cache.CacheKey;
import org.apache.hadoop.hdds.utils.db.cache.CacheValue;
import org.apache.hadoop.ozone.OzoneConsts;
-import org.apache.hadoop.ozone.om.helpers.BucketLayout;
import org.apache.hadoop.ozone.om.helpers.OmKeyInfo;
import org.apache.hadoop.ozone.om.helpers.OmKeyLocationInfo;
import org.apache.hadoop.ozone.om.helpers.OmKeyLocationInfoGroup;
@@ -139,6 +138,100 @@ public void testValidateAndUpdateCacheSuccess() throws
Exception {
.get(partKey));
}
+ @Test
+ public void
testValidateAndUpdateCacheUsesSchemaVersionOneBeforeFinalization()
+ throws Exception {
+ String volumeName = UUID.randomUUID().toString();
+ String bucketName = UUID.randomUUID().toString();
+ String keyName = getKeyName();
+
+ // The base test fixture is pre-finalized by default.
+ OMRequestTestUtils.addVolumeAndBucketToDB(volumeName, bucketName,
+ omMetadataManager, getBucketLayout());
+
+ createParentPath(volumeName, bucketName);
+
+ String multipartUploadID =
+ initiateMultipartUploadWithSchemaVersion(volumeName, bucketName,
+ keyName, OmMultipartKeyInfo.SPLIT_PARTS_TABLE_SCHEMA_VERSION);
+
+ String multipartKey = omMetadataManager.getMultipartKey(volumeName,
+ bucketName, keyName, multipartUploadID);
+ OmMultipartKeyInfo multipartKeyInfo = omMetadataManager
+ .getMultipartInfoTable().get(multipartKey);
+ assertNotNull(multipartKeyInfo);
+ assertEquals(OmMultipartKeyInfo.SPLIT_PARTS_TABLE_SCHEMA_VERSION,
+ multipartKeyInfo.getSchemaVersion());
+
+ long clientID = Time.now();
+ OMRequest commitMultipartRequest = doPreExecuteCommitMPU(volumeName,
+ bucketName, keyName, clientID, multipartUploadID, 1);
+
+ S3MultipartUploadCommitPartRequest s3MultipartUploadCommitPartRequest =
+ getS3MultipartUploadCommitReq(commitMultipartRequest);
+
+ addKeyToOpenKeyTable(volumeName, bucketName, keyName, clientID);
+
+ OMClientResponse omClientResponse =
+ s3MultipartUploadCommitPartRequest.validateAndUpdateCache(ozoneManager,
+ 2L);
+
+ assertEquals(OzoneManagerProtocolProtos.Status.OK,
+ omClientResponse.getOMResponse().getStatus());
+ }
+
+ @Test
+ public void
testValidateAndUpdateCacheAllowsSchemaVersionZeroAfterFinalization()
+ throws Exception {
+ String volumeName = UUID.randomUUID().toString();
+ String bucketName = UUID.randomUUID().toString();
+ String keyName = getKeyName();
+
+ OMRequestTestUtils.addVolumeAndBucketToDB(volumeName, bucketName,
+ omMetadataManager, getBucketLayout());
+
+ createParentPath(volumeName, bucketName);
+
+ // Upload is initiated on a pre-finalized cluster, so it uses the legacy
+ // (schema 0) inline layout.
+ OMRequest initiateMPURequest = doPreExecuteInitiateMPU(volumeName,
+ bucketName, keyName);
+ S3InitiateMultipartUploadRequest s3InitiateMultipartUploadRequest =
+ getS3InitiateMultipartUploadReq(initiateMPURequest);
+ OMClientResponse initiateResponse =
+ s3InitiateMultipartUploadRequest.validateAndUpdateCache(ozoneManager,
+ 1L);
+ String multipartUploadID = initiateResponse.getOMResponse()
+ .getInitiateMultiPartUploadResponse().getMultipartUploadID();
+
+ String multipartKey = omMetadataManager.getMultipartKey(volumeName,
+ bucketName, keyName, multipartUploadID);
+ OmMultipartKeyInfo multipartKeyInfo = omMetadataManager
+ .getMultipartInfoTable().get(multipartKey);
+ assertNotNull(multipartKeyInfo);
+ assertEquals(0, multipartKeyInfo.getSchemaVersion());
+
+ // Cluster finalizes the split feature; the pre-existing legacy part must
+ // still be committable.
+ finalizeMpuPartsTableSplit();
+
+ long clientID = Time.now();
+ OMRequest commitMultipartRequest = doPreExecuteCommitMPU(volumeName,
+ bucketName, keyName, clientID, multipartUploadID, 1);
+
+ S3MultipartUploadCommitPartRequest s3MultipartUploadCommitPartRequest =
+ getS3MultipartUploadCommitReq(commitMultipartRequest);
+
+ addKeyToOpenKeyTable(volumeName, bucketName, keyName, clientID);
+
+ OMClientResponse omClientResponse =
+ s3MultipartUploadCommitPartRequest.validateAndUpdateCache(ozoneManager,
+ 2L);
+
+ assertEquals(OzoneManagerProtocolProtos.Status.OK,
+ omClientResponse.getOMResponse().getStatus());
+ }
+
@Test
public void testValidateAndUpdateCacheMultipartNotFound() throws Exception {
String volumeName = UUID.randomUUID().toString();
@@ -172,7 +265,6 @@ public void testValidateAndUpdateCacheMultipartNotFound()
throws Exception {
bucketName, keyName, multipartUploadID);
assertNull(omMetadataManager.getMultipartInfoTable().get(multipartKey));
-
}
@Test
@@ -184,9 +276,19 @@ public void testValidateAndUpdateCacheKeyNotFound() throws
Exception {
OMRequestTestUtils.addVolumeAndBucketToDB(volumeName, bucketName,
omMetadataManager, getBucketLayout());
+ createParentPath(volumeName, bucketName);
+
+ OMRequest initiateMPURequest = doPreExecuteInitiateMPU(volumeName,
+ bucketName, keyName);
+ S3InitiateMultipartUploadRequest s3InitiateMultipartUploadRequest =
+ getS3InitiateMultipartUploadReq(initiateMPURequest);
+ OMClientResponse initiateResponse =
+ s3InitiateMultipartUploadRequest.validateAndUpdateCache(ozoneManager,
+ 1L);
long clientID = Time.now();
- String multipartUploadID = UUID.randomUUID().toString();
+ String multipartUploadID = initiateResponse.getOMResponse()
+ .getInitiateMultiPartUploadResponse().getMultipartUploadID();
OMRequest commitMultipartRequest = doPreExecuteCommitMPU(volumeName,
bucketName, keyName, clientID, multipartUploadID, 1);
@@ -201,13 +303,8 @@ public void testValidateAndUpdateCacheKeyNotFound() throws
Exception {
OMClientResponse omClientResponse =
s3MultipartUploadCommitPartRequest.validateAndUpdateCache(ozoneManager, 2L);
- if (getBucketLayout() == BucketLayout.FILE_SYSTEM_OPTIMIZED) {
- assertSame(omClientResponse.getOMResponse().getStatus(),
- OzoneManagerProtocolProtos.Status.DIRECTORY_NOT_FOUND);
- } else {
- assertSame(omClientResponse.getOMResponse().getStatus(),
- OzoneManagerProtocolProtos.Status.KEY_NOT_FOUND);
- }
+ assertSame(omClientResponse.getOMResponse().getStatus(),
+ OzoneManagerProtocolProtos.Status.KEY_NOT_FOUND);
}
@@ -742,7 +839,13 @@ public void testSplitSchemaCommitFailsWithoutETag() throws
Exception {
S3MultipartUploadCommitPartRequest request =
getS3MultipartUploadCommitReq(omRequest);
OMClientResponse response = request.validateAndUpdateCache(ozoneManager,
2L);
- assertSame(OzoneManagerProtocolProtos.Status.INVALID_REQUEST,
response.getOMResponse().getStatus());
+ assertSame(OzoneManagerProtocolProtos.Status.INVALID_REQUEST,
+ response.getOMResponse().getStatus());
+
+ // No part row should have been written to the split parts table.
+ SortedMap<Integer, OmMultipartPartInfo> parts =
+ OMMultipartUploadUtils.scanParts(omMetadataManager, uploadId);
+ assertEquals(0, parts.size());
}
@Test
diff --git
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3MultipartUploadCompleteRequest.java
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3MultipartUploadCompleteRequest.java
index 653b0703a67..27b7217b2dc 100644
---
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3MultipartUploadCompleteRequest.java
+++
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3MultipartUploadCompleteRequest.java
@@ -37,6 +37,7 @@
import org.apache.hadoop.ozone.OzoneConsts;
import org.apache.hadoop.ozone.om.helpers.OmBucketInfo;
import org.apache.hadoop.ozone.om.helpers.OmKeyInfo;
+import org.apache.hadoop.ozone.om.helpers.OmMultipartKeyInfo;
import org.apache.hadoop.ozone.om.helpers.RepeatedOmKeyInfo;
import org.apache.hadoop.ozone.om.request.OMRequestTestUtils;
import org.apache.hadoop.ozone.om.response.OMClientResponse;
@@ -206,6 +207,82 @@ private String checkValidateAndUpdateCacheSuccess(String
volumeName,
return multipartUploadID;
}
+ @Test
+ public void
testValidateAndUpdateCacheUsesSchemaVersionOneBeforeFinalization()
+ throws Exception {
+ String volumeName = UUID.randomUUID().toString();
+ String bucketName = UUID.randomUUID().toString();
+ String keyName = getKeyName();
+
+ // The base test fixture is pre-finalized by default.
+ OMRequestTestUtils.addVolumeAndBucketToDB(volumeName, bucketName,
+ omMetadataManager, getBucketLayout());
+
+ String multipartUploadID =
+ initiateMultipartUploadWithSchemaVersion(volumeName, bucketName,
+ keyName, 1);
+
+ OMRequest completeMultipartRequest = doPreExecuteCompleteMPU(volumeName,
+ bucketName, keyName, multipartUploadID, new ArrayList<>());
+
+ S3MultipartUploadCompleteRequest s3MultipartUploadCompleteRequest =
+ getS3MultipartUploadCompleteReq(completeMultipartRequest);
+
+ OMClientResponse omClientResponse =
+ s3MultipartUploadCompleteRequest.validateAndUpdateCache(ozoneManager,
+ 3L);
+
+ assertEquals(OzoneManagerProtocolProtos.Status.INVALID_REQUEST,
+ omClientResponse.getOMResponse().getStatus());
+ }
+
+ @Test
+ public void
testValidateAndUpdateCacheAllowsSchemaVersionZeroAfterFinalization()
+ throws Exception {
+ String volumeName = UUID.randomUUID().toString();
+ String bucketName = UUID.randomUUID().toString();
+ String keyName = getKeyName();
+
+ OMRequestTestUtils.addVolumeAndBucketToDB(volumeName, bucketName,
+ omMetadataManager, getBucketLayout());
+
+ // Upload is initiated on a pre-finalized cluster, so it uses the legacy
+ // (schema 0) inline layout.
+ OMRequest initiateMPURequest = doPreExecuteInitiateMPU(volumeName,
+ bucketName, keyName);
+ S3InitiateMultipartUploadRequest s3InitiateMultipartUploadRequest =
+ getS3InitiateMultipartUploadReq(initiateMPURequest);
+ OMClientResponse initiateResponse =
+ s3InitiateMultipartUploadRequest.validateAndUpdateCache(ozoneManager,
+ 1L);
+ String multipartUploadID = initiateResponse.getOMResponse()
+ .getInitiateMultiPartUploadResponse().getMultipartUploadID();
+
+ String multipartKey = getMultipartKey(volumeName, bucketName, keyName,
+ multipartUploadID);
+ OmMultipartKeyInfo multipartKeyInfo = omMetadataManager
+ .getMultipartInfoTable().get(multipartKey);
+ assertNotNull(multipartKeyInfo);
+ assertEquals(0, multipartKeyInfo.getSchemaVersion());
+
+ // Cluster finalizes the split feature; completing the pre-existing legacy
+ // upload must still behave as before.
+ finalizeMpuPartsTableSplit();
+
+ OMRequest completeMultipartRequest = doPreExecuteCompleteMPU(volumeName,
+ bucketName, keyName, multipartUploadID, new ArrayList<>());
+
+ S3MultipartUploadCompleteRequest s3MultipartUploadCompleteRequest =
+ getS3MultipartUploadCompleteReq(completeMultipartRequest);
+
+ OMClientResponse omClientResponse =
+ s3MultipartUploadCompleteRequest.validateAndUpdateCache(ozoneManager,
+ 3L);
+
+ assertEquals(OzoneManagerProtocolProtos.Status.INVALID_REQUEST,
+ omClientResponse.getOMResponse().getStatus());
+ }
+
protected void addVolumeAndBucket(String volumeName, String bucketName)
throws Exception {
OMRequestTestUtils.addVolumeAndBucketToDB(volumeName, bucketName,
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 3989feec27c..e1687672eaf 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
@@ -544,6 +544,28 @@ static void addTagCountIfAny(
}
}
+ /**
+ * Guarantees that an S3-originated multipart upload part carries an ETag
+ * before it is committed. S3 requires every uploaded part to have an ETag
+ * (see UploadPart and CompleteMultipartUpload), and all S3 gateway
+ * part-write paths stamp the MD5 of the part as its ETag. This pre-commit
+ * check enforces that invariant so no S3 part can be committed without one
+ * (e.g. an UploadPartCopy whose source key has no ETag).
+ *
+ * <p>Note: the Ozone Manager also enforces a mandatory ETag server-side for
+ * every committed part (in the split parts-table schema) for ALL clients,
+ * not just S3. This gateway-side check is an earlier, S3-native failure so
+ * the client gets a clean S3 error instead of an OM INVALID_REQUEST.
+ */
+ private static void requirePartETag(Map<String, String> metadata)
+ throws IOException {
+ if (metadata == null
+ || StringUtils.isBlank(metadata.get(OzoneConsts.ETAG))) {
+ throw new IOException(
+ "S3 multipart upload part cannot be committed without an ETag");
+ }
+ }
+
static void addEntityTagHeader(ResponseBuilder responseBuilder, OzoneKey
key) {
String eTag = key.getMetadata().get(OzoneConsts.ETAG);
if (eTag != null) {
@@ -963,6 +985,8 @@ private Response createMultipartKey(OzoneVolume volume,
OzoneBucket ozoneBucket,
if (raw != null) {
writeGuard.getMetadata().put(OzoneConsts.ETAG, stripQuotes(raw));
}
+ writeGuard.addPreCommit(
+ () -> requirePartETag(writeGuard.getMetadata()));
outputStream = ozoneOutputStream;
}
getMetrics().incCopyObjectSuccessLength(copyLength);
@@ -987,6 +1011,8 @@ private Response createMultipartKey(OzoneVolume volume,
OzoneBucket ozoneBucket,
writeGuard.addPreCommit(checkContentMD5Hook);
}
writeGuard.getMetadata().put(OzoneConsts.ETAG, md5Hash);
+ writeGuard.addPreCommit(
+ () -> requirePartETag(writeGuard.getMetadata()));
outputStream = ozoneOutputStream;
}
getMetrics().incPutKeySuccessLength(putLength);
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]