From c3bfabe24d0e4b827a9f87b80e73d10cba83a3f2 Mon Sep 17 00:00:00 2001 From: Chi-Hsuan Huang Date: Mon, 7 Sep 2026 21:36:56 +0800 Subject: [PATCH 01/14] HDDS-16161. Add getBucketInfoForUpdate for callers that mutate the bucket --- .../ozone/om/request/key/OMKeyRequest.java | 25 +++++++++++++++---- 1 file changed, 20 insertions(+), 5 deletions(-) diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyRequest.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyRequest.java index 3d478e461396..244c8c6638e3 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyRequest.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyRequest.java @@ -945,13 +945,11 @@ public static long sumBlockLengths(OmKeyInfo omKeyInfo) { } /** - * Return bucket info for the specified bucket. + * Return bucket info for the specified bucket, for read-only use. *

* The returned {@link OmBucketInfo} is the cached instance, returned by - * reference. A caller that mutates it (for example quota accounting) before a - * point where the request may still fail must first take a - * {@link OmBucketInfo#copyObject()} and publish that copy only on success, - * otherwise a failed request leaks the mutation into the cache. + * reference. Callers that mutate it must use + * {@link #getBucketInfoForUpdate(OMMetadataManager, String, String)} instead. */ @Nullable public static OmBucketInfo getBucketInfo(OMMetadataManager omMetadataManager, @@ -964,6 +962,23 @@ public static OmBucketInfo getBucketInfo(OMMetadataManager omMetadataManager, return value != null ? value.getCacheValue() : null; } + /** + * Return a copy of the cached bucket info for callers that mutate it. + *

+ * Mutations stay invisible until the caller publishes the copy with + * {@code getBucketTable().addCacheEntry(...)}, after all fallible work and + * only on the path that persists the response. Hold the bucket write lock for + * the whole read-modify-publish, not just the publish, or a concurrent writer + * can be lost. + */ + @Nullable + public static OmBucketInfo getBucketInfoForUpdate(OMMetadataManager omMetadataManager, + String volume, String bucket) { + OmBucketInfo omBucketInfo = getBucketInfo(omMetadataManager, volume, bucket); + + return omBucketInfo != null ? omBucketInfo.copyObject() : null; + } + /** * Prepare OmKeyInfo which will be persisted to openKeyTable. * @return OmKeyInfo From 6db5df6236280cd0d0c5823ae2eb6b8727811bdd Mon Sep 17 00:00:00 2001 From: Chi-Hsuan Huang Date: Mon, 7 Sep 2026 21:37:01 +0800 Subject: [PATCH 02/14] HDDS-16161. Publish bucket usage from a copy in directory create requests --- .../file/OMDirectoryCreateRequest.java | 7 +++-- .../file/OMDirectoryCreateRequestWithFSO.java | 6 +++- .../TestOMDirectoryCreateRequestWithFSO.java | 28 +++++++++++++++++++ 3 files changed, 38 insertions(+), 3 deletions(-) diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMDirectoryCreateRequest.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMDirectoryCreateRequest.java index 97e3d612c8ec..6b7051e5d607 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMDirectoryCreateRequest.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMDirectoryCreateRequest.java @@ -177,7 +177,7 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut List missingParents = omPathInfo.getMissingParents(); long baseObjId = ozoneManager.getObjectIdFromTxId(trxnLogIndex); OmBucketInfo omBucketInfo = - getBucketInfo(omMetadataManager, volumeName, bucketName); + getBucketInfoForUpdate(omMetadataManager, volumeName, bucketName); dirKeyInfo = createDirectoryKeyInfoWithACL(keyName, keyArgs, baseObjId, omBucketInfo, omPathInfo, trxnLogIndex, @@ -193,7 +193,10 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut OMFileRequest.addKeyTableCacheEntries(omMetadataManager, volumeName, bucketName, omBucketInfo.getBucketLayout(), dirKeyInfo, missingParentInfos, trxnLogIndex); - + + omMetadataManager.getBucketTable().addCacheEntry( + omMetadataManager.getBucketKey(volumeName, bucketName), omBucketInfo, trxnLogIndex); + result = Result.SUCCESS; omClientResponse = new OMDirectoryCreateResponse(omResponse.build(), dirKeyInfo, missingParentInfos, result, getBucketLayout(), diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMDirectoryCreateRequestWithFSO.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMDirectoryCreateRequestWithFSO.java index dde1eebd02c0..77974fbf3541 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMDirectoryCreateRequestWithFSO.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMDirectoryCreateRequestWithFSO.java @@ -133,7 +133,7 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut omDirectoryResult == NONE) { OmBucketInfo omBucketInfo = - getBucketInfo(omMetadataManager, volumeName, bucketName); + getBucketInfoForUpdate(omMetadataManager, volumeName, bucketName); // prepare all missing parents missingParentInfos = getAllMissingParentDirInfo( ozoneManager, keyArgs, omBucketInfo, omPathInfo, trxnLogIndex); @@ -157,6 +157,10 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut volumeId, bucketId, trxnLogIndex, missingParentInfos, dirInfo); + // Publish only here: createDirectoryInfoWithACL above can still fail with UNAUTHORIZED. + omMetadataManager.getBucketTable().addCacheEntry( + omMetadataManager.getBucketKey(volumeName, bucketName), omBucketInfo, trxnLogIndex); + result = OMDirectoryCreateRequest.Result.SUCCESS; omClientResponse = new OMDirectoryCreateResponseWithFSO(omResponse.build(), diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/file/TestOMDirectoryCreateRequestWithFSO.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/file/TestOMDirectoryCreateRequestWithFSO.java index 8f5b6c807caa..52031b79f2dd 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/file/TestOMDirectoryCreateRequestWithFSO.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/file/TestOMDirectoryCreateRequestWithFSO.java @@ -219,6 +219,34 @@ public void testValidateAndUpdateCacheWithNamespaceQuotaExceeded() OzoneManagerProtocolProtos.Status.QUOTA_EXCEEDED); } + @Test + public void testFailedRequestDoesNotLeakNamespaceQuota() throws Exception { + String volumeName = "vol1"; + String bucketName = "bucket1"; + // A single-level directory has no missing parents, so the first ACL lookup, and with it the + // first chance to fail, happens after the namespace charge is applied. + String keyName = "dir1"; + + OMRequestTestUtils.addVolumeAndBucketToDB(volumeName, bucketName, + omMetadataManager, getBucketLayout()); + String bucketKey = omMetadataManager.getBucketKey(volumeName, bucketName); + long usedNamespaceBefore = omMetadataManager.getBucketTable().get(bucketKey).getUsedNamespace(); + + // Neither preExecute nor setUGI runs, so the request carries no user info and the ACL lookup + // for the leaf directory fails with UNAUTHORIZED. + OMDirectoryCreateRequestWithFSO omDirCreateRequestFSO = + new OMDirectoryCreateRequestWithFSO( + createDirectoryRequest(volumeName, bucketName, keyName), + BucketLayout.FILE_SYSTEM_OPTIMIZED); + + OMClientResponse omClientResponse = + omDirCreateRequestFSO.validateAndUpdateCache(ozoneManager, 100L); + + assertSame(OzoneManagerProtocolProtos.Status.UNAUTHORIZED, + omClientResponse.getOMResponse().getStatus()); + assertEquals(usedNamespaceBefore, omMetadataManager.getBucketTable().get(bucketKey).getUsedNamespace()); + } + @Test public void testValidateAndUpdateCacheWithVolumeNotFound() throws Exception { String volumeName = "vol1"; From ad891243da51bbbf58aeb967c8bcd2abe8c014b8 Mon Sep 17 00:00:00 2001 From: Chi-Hsuan Huang Date: Mon, 7 Sep 2026 21:37:07 +0800 Subject: [PATCH 03/14] HDDS-16161. Publish bucket usage from a copy in file create requests --- .../hadoop/ozone/om/request/file/OMFileCreateRequest.java | 5 ++++- .../ozone/om/request/file/OMFileCreateRequestWithFSO.java | 5 ++++- 2 files changed, 8 insertions(+), 2 deletions(-) diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMFileCreateRequest.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMFileCreateRequest.java index 35e1ac238f7b..b38eeab6fe07 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMFileCreateRequest.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMFileCreateRequest.java @@ -227,7 +227,7 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut // do open key omBucketInfo = - getBucketInfo(omMetadataManager, volumeName, bucketName); + getBucketInfoForUpdate(omMetadataManager, volumeName, bucketName); final ReplicationConfig repConfig = OzoneConfigUtil .resolveReplicationConfigPreference(keyArgs.getType(), keyArgs.getFactor(), keyArgs.getEcReplicationConfig(), @@ -277,6 +277,9 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut bucketName, omBucketInfo.getBucketLayout(), null, missingParentInfos, trxnLogIndex); + omMetadataManager.getBucketTable().addCacheEntry( + omMetadataManager.getBucketKey(volumeName, bucketName), omBucketInfo, trxnLogIndex); + // Prepare response omResponse.setCreateFileResponse(CreateFileResponse.newBuilder() .setKeyInfo(omKeyInfo.getNetworkProtobuf(getOmRequest().getVersion(), diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMFileCreateRequestWithFSO.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMFileCreateRequestWithFSO.java index 6036fe90dbb7..18f506cf728e 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMFileCreateRequestWithFSO.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMFileCreateRequestWithFSO.java @@ -184,7 +184,7 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut .collect(Collectors.toList()); omFileInfo.appendNewBlocks(newLocationList, false); - omBucketInfo = getBucketInfo(omMetadataManager, volumeName, bucketName); + omBucketInfo = getBucketInfoForUpdate(omMetadataManager, volumeName, bucketName); // check bucket and volume quota long preAllocatedSpace = newLocationList.size() * ozoneManager.getScmBlockSize() * repConfig @@ -205,6 +205,9 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut OMFileRequest.addDirectoryTableCacheEntries(omMetadataManager, volumeId, bucketId, trxnLogIndex, missingParentInfos, null); + omMetadataManager.getBucketTable().addCacheEntry( + omMetadataManager.getBucketKey(volumeName, bucketName), omBucketInfo, trxnLogIndex); + // Prepare response. Sets user given full key name in the 'keyName' // attribute in response object. int clientVersion = getOmRequest().getVersion(); From 133b8b8ff9c4238fa697ec5e3b944c0ef51567f9 Mon Sep 17 00:00:00 2001 From: Chi-Hsuan Huang Date: Mon, 7 Sep 2026 21:37:07 +0800 Subject: [PATCH 04/14] HDDS-16161. Publish bucket usage from a copy in OMKeyCreateRequestWithFSO --- .../ozone/om/request/key/OMKeyCreateRequestWithFSO.java | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCreateRequestWithFSO.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCreateRequestWithFSO.java index 3642f96c7fcf..c94480fb72b2 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCreateRequestWithFSO.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCreateRequestWithFSO.java @@ -175,7 +175,7 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut .collect(Collectors.toList()); omFileInfo.appendNewBlocks(newLocationList, false); - omBucketInfo = getBucketInfo(omMetadataManager, volumeName, bucketName); + omBucketInfo = getBucketInfoForUpdate(omMetadataManager, volumeName, bucketName); // check bucket and volume quota long preAllocatedSpace = newLocationList.size() * ozoneManager.getScmBlockSize() * repConfig @@ -201,6 +201,9 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut volumeId, bucketId, trxnLogIndex, missingParentInfos, null); + omMetadataManager.getBucketTable().addCacheEntry( + omMetadataManager.getBucketKey(volumeName, bucketName), omBucketInfo, trxnLogIndex); + // Prepare response. Sets user given full key name in the 'keyName' // attribute in response object. int clientVersion = getOmRequest().getVersion(); From 36a6d1d68c9d890d7571a96b76c5ca262eeed474 Mon Sep 17 00:00:00 2001 From: Chi-Hsuan Huang Date: Mon, 7 Sep 2026 21:37:13 +0800 Subject: [PATCH 05/14] HDDS-16161. Publish bucket usage only when OMKeyCreateRequest charges parent directories --- .../om/request/key/OMKeyCreateRequest.java | 7 +- .../request/key/TestOMKeyCreateRequest.java | 66 +++++++++++++++++++ 2 files changed, 72 insertions(+), 1 deletion(-) diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCreateRequest.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCreateRequest.java index 6059d6c46c75..04bd9bcf832e 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCreateRequest.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCreateRequest.java @@ -247,7 +247,7 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut keyArgs = validateAndRewriteIfMatchAsExpectedGeneration(keyArgs, dbKeyInfo); OmBucketInfo bucketInfo = - getBucketInfo(omMetadataManager, volumeName, bucketName); + getBucketInfoForUpdate(omMetadataManager, volumeName, bucketName); // If FILE_EXISTS we just override like how we used to do for Key Create. if (LOG.isDebugEnabled()) { @@ -333,6 +333,11 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut OMFileRequest.addKeyTableCacheEntries(omMetadataManager, volumeName, bucketName, bucketInfo.getBucketLayout(), null, missingParentInfos, trxnLogIndex); + + // Parent directory creation holds the bucket write lock; key path locking leaves + // numMissingParents at 0. + omMetadataManager.getBucketTable().addCacheEntry( + omMetadataManager.getBucketKey(volumeName, bucketName), bucketInfo, trxnLogIndex); } // Add to cache entry can be done outside of lock for this openKey. diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyCreateRequest.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyCreateRequest.java index f00276040dfa..6fce4602a11d 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyCreateRequest.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyCreateRequest.java @@ -40,6 +40,7 @@ import static org.junit.jupiter.api.Assertions.assertSame; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.junit.jupiter.api.Assumptions.assumeTrue; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyInt; import static org.mockito.ArgumentMatchers.anyLong; @@ -68,6 +69,8 @@ import org.apache.hadoop.hdds.protocol.proto.HddsProtos.KeyValue; import org.apache.hadoop.hdds.scm.container.common.helpers.AllocatedBlock; import org.apache.hadoop.hdds.scm.container.common.helpers.ExcludeList; +import org.apache.hadoop.hdds.utils.db.BatchOperation; +import org.apache.hadoop.hdds.utils.db.cache.CacheKey; import org.apache.hadoop.hdds.scm.pipeline.Pipeline; import org.apache.hadoop.hdds.scm.pipeline.PipelineID; import org.apache.hadoop.ozone.OzoneAcl; @@ -1060,6 +1063,69 @@ private OMRequest createKeyRequestWithExpectedETag(String expectedETag) { .setCreateKeyRequest(createKeyRequest).build(); } + @Test + public void testNamespaceChargeReachesCacheAndDb() throws Exception { + OzoneConfiguration configuration = getOzoneConfiguration(); + configuration.setBoolean(OZONE_OM_ENABLE_FILESYSTEM_PATHS, true); + when(ozoneManager.getConfiguration()).thenReturn(configuration); + when(ozoneManager.getConfig()).thenReturn(configuration.getObject(OmConfig.class)); + when(ozoneManager.getEnableFileSystemPaths()).thenReturn(true); + when(ozoneManager.getOzoneLockProvider()).thenReturn(new OzoneLockProvider(false, true)); + + addVolumeAndBucketToDB(volumeName, bucketName, omMetadataManager, getBucketLayout()); + String bucketKey = omMetadataManager.getBucketKey(volumeName, bucketName); + + // dir1 and dir1/dir2 are the missing parents; the key itself is charged at commit. + OMRequest omRequest = createKeyRequest(false, 0, "dir1/dir2/file1"); + OMKeyCreateRequest omKeyCreateRequest = getOMKeyCreateRequest(omRequest); + omKeyCreateRequest = getOMKeyCreateRequest(omKeyCreateRequest.preExecute(ozoneManager)); + + OMClientResponse omClientResponse = + omKeyCreateRequest.validateAndUpdateCache(ozoneManager, 100L); + assertEquals(OzoneManagerProtocolProtos.Status.OK, + omClientResponse.getOMResponse().getStatus()); + + assertEquals(2, omMetadataManager.getBucketTable().get(bucketKey).getUsedNamespace()); + + try (BatchOperation batchOperation = omMetadataManager.getStore().initBatchOperation()) { + omClientResponse.checkAndUpdateDB(omMetadataManager, batchOperation); + omMetadataManager.getStore().commitBatchOperation(batchOperation); + } + assertEquals(2, omMetadataManager.getBucketTable().getSkipCache(bucketKey).getUsedNamespace()); + } + + @Test + public void testKeyPathLockLeavesBucketCacheEntryAlone() throws Exception { + // Key path locking is only chosen for OBJECT_STORE and for LEGACY without filesystem paths. + assumeTrue(getBucketLayout() == BucketLayout.LEGACY); + + OzoneConfiguration configuration = getOzoneConfiguration(); + configuration.setBoolean(OZONE_OM_ENABLE_FILESYSTEM_PATHS, false); + when(ozoneManager.getConfiguration()).thenReturn(configuration); + when(ozoneManager.getConfig()).thenReturn(configuration.getObject(OmConfig.class)); + when(ozoneManager.getEnableFileSystemPaths()).thenReturn(false); + when(ozoneManager.getOzoneLockProvider()).thenReturn(new OzoneLockProvider(true, false)); + + addVolumeAndBucketToDB(volumeName, bucketName, omMetadataManager, getBucketLayout()); + String bucketKey = omMetadataManager.getBucketKey(volumeName, bucketName); + OmBucketInfo cachedBefore = omMetadataManager.getBucketTable() + .getCacheValue(new CacheKey<>(bucketKey)).getCacheValue(); + + OMRequest omRequest = createKeyRequest(false, 0, "dir1/dir2/file1"); + OMKeyCreateRequest omKeyCreateRequest = getOMKeyCreateRequest(omRequest); + omKeyCreateRequest = getOMKeyCreateRequest(omKeyCreateRequest.preExecute(ozoneManager)); + + OMClientResponse omClientResponse = + omKeyCreateRequest.validateAndUpdateCache(ozoneManager, 100L); + assertEquals(OzoneManagerProtocolProtos.Status.OK, + omClientResponse.getOMResponse().getStatus()); + + // No parent directories are charged here, and only a bucket read lock is held, so the request + // must not replace the cached bucket. + assertSame(cachedBefore, omMetadataManager.getBucketTable() + .getCacheValue(new CacheKey<>(bucketKey)).getCacheValue()); + } + @ParameterizedTest @ValueSource(booleans = {true, false}) public void testKeyCreateWithFileSystemPathsEnabled( From 0ee7e43dbd3f0de3ed3597e77e9927bcad7fbe05 Mon Sep 17 00:00:00 2001 From: Chi-Hsuan Huang Date: Mon, 7 Sep 2026 21:50:35 +0800 Subject: [PATCH 06/14] HDDS-16161. Publish bucket usage from a copy in key commit requests --- .../om/request/key/OMKeyCommitRequest.java | 5 ++- .../key/OMKeyCommitRequestWithFSO.java | 5 ++- .../request/key/TestOMKeyCommitRequest.java | 36 +++++++++++++++++++ 3 files changed, 44 insertions(+), 2 deletions(-) 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 0789e1288507..d8a8bd87c6d6 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 @@ -194,7 +194,7 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut bucketLockAcquired = getOmLockDetails().isLockAcquired(); validateBucketAndVolume(omMetadataManager, volumeName, bucketName); - omBucketInfo = getBucketInfo(omMetadataManager, volumeName, bucketName); + omBucketInfo = getBucketInfoForUpdate(omMetadataManager, volumeName, bucketName); // Check for directory exists with same name, if it exists throw error. if (LOG.isDebugEnabled()) { @@ -406,6 +406,9 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut omBucketInfo.incrUsedBytes(correctedSpace); + omMetadataManager.getBucketTable().addCacheEntry( + omMetadataManager.getBucketKey(volumeName, bucketName), omBucketInfo, trxnLogIndex); + 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 9bc6f7ec0d9d..dce295865218 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 @@ -129,7 +129,7 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut bucketLockAcquired = getOmLockDetails().isLockAcquired(); validateBucketAndVolume(omMetadataManager, volumeName, bucketName); - omBucketInfo = getBucketInfo(omMetadataManager, volumeName, bucketName); + omBucketInfo = getBucketInfoForUpdate(omMetadataManager, volumeName, bucketName); String errMsg = "Cannot create file : " + keyName + " as parent directory doesn't exist"; @@ -350,6 +350,9 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut omBucketInfo.incrUsedBytes(correctedSpace); + omMetadataManager.getBucketTable().addCacheEntry( + omMetadataManager.getBucketKey(volumeName, bucketName), omBucketInfo, trxnLogIndex); + omClientResponse = new OMKeyCommitResponseWithFSO(omResponse.build(), omKeyInfo, dbFileKey, dbOpenFileKey, omBucketInfo.copyObject(), oldKeyVersionsToDeleteMap, volumeId, isHSync, newOpenKeyInfo, dbOpenKeyToDeleteKey, openKeyToDelete); diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyCommitRequest.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyCommitRequest.java index 2a87c8e00a83..fef410988328 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyCommitRequest.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyCommitRequest.java @@ -348,6 +348,42 @@ public void testAtomicCreateIfNotExistsCommitKeyAlreadyExists() throws Exception assertEquals(ATOMIC_WRITE_CONFLICT, omClientResponse.getOMResponse().getStatus()); } + @Test + public void testFailedCommitDoesNotLeakBucketUsage() throws Exception { + // Uncommitted blocks make the request build a pseudo key for deletion, which derives an + // object id from the transaction index. An index above MAX_TRXN_ID makes that step throw + // after the bucket counters have already been changed. + List allocatedKeyLocationList = getKeyLocation(5); + List allocatedBlockList = allocatedKeyLocationList + .stream().map(OmKeyLocationInfo::getFromProtobuf) + .collect(Collectors.toList()); + + OMRequest modifiedOmRequest = doPreExecute(createCommitKeyRequest( + allocatedKeyLocationList.subList(0, 3), false)); + OMKeyCommitRequest omKeyCommitRequest = + getOmKeyCommitRequest(modifiedOmRequest); + + OMRequestTestUtils.addVolumeAndBucketToDB(volumeName, bucketName, + omMetadataManager, omKeyCommitRequest.getBucketLayout()); + addKeyToOpenKeyTable(allocatedBlockList); + + String bucketKey = omMetadataManager.getBucketKey(volumeName, bucketName); + OmBucketInfo cachedBefore = omMetadataManager.getBucketTable().get(bucketKey); + long usedNamespaceBefore = cachedBefore.getUsedNamespace(); + long usedBytesBefore = cachedBefore.getUsedBytes(); + + // The mock returns 0 by default, so let it run the real check the OM performs. + when(ozoneManager.getObjectIdFromTxId(anyLong())).thenAnswer(invocation -> OmUtils + .getObjectIdFromTxId(omMetadataManager.getOmEpoch(), invocation.getArgument(0))); + + assertThrows(IllegalArgumentException.class, () -> omKeyCommitRequest + .validateAndUpdateCache(ozoneManager, OmUtils.MAX_TRXN_ID + 1)); + + OmBucketInfo cachedAfter = omMetadataManager.getBucketTable().get(bucketKey); + assertEquals(usedNamespaceBefore, cachedAfter.getUsedNamespace()); + assertEquals(usedBytesBefore, cachedAfter.getUsedBytes()); + } + @Test public void testValidateAndUpdateCacheWithUncommittedBlocks() throws Exception { From fed6595f021977dff2d5e9c2ce29689dad9d7063 Mon Sep 17 00:00:00 2001 From: Chi-Hsuan Huang Date: Mon, 7 Sep 2026 22:10:29 +0800 Subject: [PATCH 07/14] HDDS-16161. Publish bucket usage from a copy in key delete requests --- .../om/request/key/OMKeyDeleteRequest.java | 5 ++- .../key/OMKeyDeleteRequestWithFSO.java | 5 ++- .../om/request/key/OMKeysDeleteRequest.java | 8 +++- .../request/key/TestOMKeysDeleteRequest.java | 43 +++++++++++++++++++ .../key/TestOMKeysDeleteRequestWithFSO.java | 16 +++++++ 5 files changed, 74 insertions(+), 3 deletions(-) diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyDeleteRequest.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyDeleteRequest.java index 26287ca66d26..bcf2e7a176c8 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyDeleteRequest.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyDeleteRequest.java @@ -160,7 +160,7 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut CacheValue.get(trxnLogIndex)); OmBucketInfo omBucketInfo = - getBucketInfo(omMetadataManager, volumeName, bucketName); + getBucketInfoForUpdate(omMetadataManager, volumeName, bucketName); long quotaReleased = sumBlockLengths(omKeyInfo); // Empty entries won't be added to deleted table so this key shouldn't get added to snapshotUsed space. @@ -186,6 +186,9 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut } } + omMetadataManager.getBucketTable().addCacheEntry( + omMetadataManager.getBucketKey(volumeName, bucketName), omBucketInfo, trxnLogIndex); + omClientResponse = new OMKeyDeleteResponse( omResponse.setDeleteKeyResponse(DeleteKeyResponse.newBuilder()) .build(), omKeyInfo, diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyDeleteRequestWithFSO.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyDeleteRequestWithFSO.java index 769b2e43a5b4..c0c83f5c4f41 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyDeleteRequestWithFSO.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyDeleteRequestWithFSO.java @@ -157,7 +157,7 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut CacheValue.get(trxnLogIndex)); } - omBucketInfo = getBucketInfo(omMetadataManager, volumeName, bucketName); + omBucketInfo = getBucketInfoForUpdate(omMetadataManager, volumeName, bucketName); long quotaReleased = sumBlockLengths(omKeyInfo); // Empty entries won't be added to deleted table so this key shouldn't get added to snapshotUsed space. @@ -188,6 +188,9 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut auditMap.put(OzoneConsts.REPLICATION_CONFIG, omKeyInfo.getReplicationConfig().toString()); } + omMetadataManager.getBucketTable().addCacheEntry( + omMetadataManager.getBucketKey(volumeName, bucketName), omBucketInfo, trxnLogIndex); + omClientResponse = new OMKeyDeleteResponseWithFSO(omResponse .setDeleteKeyResponse(DeleteKeyResponse.newBuilder()).build(), keyName, omKeyInfo, diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeysDeleteRequest.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeysDeleteRequest.java index 3fc15da75ee9..dbccd61c5169 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeysDeleteRequest.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeysDeleteRequest.java @@ -236,7 +236,7 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut } OmBucketInfo omBucketInfo = - getBucketInfo(omMetadataManager, volumeName, bucketName); + getBucketInfoForUpdate(omMetadataManager, volumeName, bucketName); Map openKeyInfoMap = new HashMap<>(); // Mark all keys which can be deleted, in cache as deleted. @@ -259,6 +259,12 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut } final long volumeId = omMetadataManager.getVolumeId(volumeName); + + // A partial delete still persists the accepted keys and the bucket row, so publish here + // rather than only when every key was deleted. + omMetadataManager.getBucketTable().addCacheEntry( + omMetadataManager.getBucketKey(volumeName, bucketName), omBucketInfo, trxnLogIndex); + omClientResponse = getOmClientResponse(ozoneManager, omKeyInfoList, dirList, omResponse, unDeletedKeys, keyToError, deleteStatus, omBucketInfo, volumeId, openKeyInfoMap, state); diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeysDeleteRequest.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeysDeleteRequest.java index 52626c68802b..86a90a40e534 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeysDeleteRequest.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeysDeleteRequest.java @@ -38,6 +38,7 @@ import java.util.List; import java.util.UUID; import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.utils.db.BatchOperation; import org.apache.hadoop.ozone.om.OzoneManager; import org.apache.hadoop.ozone.om.exceptions.OMException; import org.apache.hadoop.ozone.om.request.OMRequestTestUtils; @@ -147,6 +148,44 @@ public void testKeysDeleteRequestFail(RequestSource sourceType) throws Exception checkDeleteKeysResponseForFailure(omKeysDeleteRequest, Status.PARTIAL_DELETE, sourceType); } + @Test + public void testPartialDeletePublishesBucketUsage() throws Exception { + createPreRequisites(); + + String bucketKey = omMetadataManager.getBucketKey(volumeName, bucketName); + long usedNamespaceBefore = omMetadataManager.getBucketTable().get(bucketKey).getUsedNamespace(); + + // One key does not exist, so the request succeeds partially. + omRequest = omRequest.toBuilder() + .setDeleteKeysRequest(DeleteKeysRequest.newBuilder() + .setDeleteKeys(DeleteKeyArgs.newBuilder() + .setBucketName(bucketName).setVolumeName(volumeName) + .addAllKeys(deleteKeyList).addKeys("dummy"))).build(); + + OMClientResponse omClientResponse = getOmKeysDeleteRequest(omRequest) + .validateAndUpdateCache(ozoneManager, 100L); + + assertEquals(Status.PARTIAL_DELETE, omClientResponse.getOMResponse().getStatus()); + + // PARTIAL_DELETE still persists the accepted keys and the bucket row, so the released + // namespace has to reach the cache too. Only the keys that existed are released; "dummy" is + // counted as undeleted. + long usedNamespaceCached = omMetadataManager.getBucketTable().get(bucketKey).getUsedNamespace(); + assertEquals(usedNamespaceBefore - expectedNamespaceReleasedOnPartialDelete(), + usedNamespaceCached); + + try (BatchOperation batchOperation = omMetadataManager.getStore().initBatchOperation()) { + omClientResponse.checkAndUpdateDB(omMetadataManager, batchOperation); + omMetadataManager.getStore().commitBatchOperation(batchOperation); + } + assertEquals(usedNamespaceCached, + omMetadataManager.getBucketTable().getSkipCache(bucketKey).getUsedNamespace()); + } + + protected long expectedNamespaceReleasedOnPartialDelete() { + return KEY_COUNT; + } + @ParameterizedTest @MethodSource("requestSourceType") public void testUpdateIDCountNoMatchKeyCount() throws Exception { @@ -282,6 +321,10 @@ protected void createPreRequisites() throws Exception { createPreRequisites(RequestSource.USER); } + protected OMKeysDeleteRequest getOmKeysDeleteRequest(OMRequest request) { + return new OMKeysDeleteRequest(request, getBucketLayout()); + } + protected void createPreRequisites(RequestSource sourceType) throws Exception { deleteKeyList = new ArrayList<>(); diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeysDeleteRequestWithFSO.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeysDeleteRequestWithFSO.java index 071ef81513a6..8f4ad0750811 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeysDeleteRequestWithFSO.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeysDeleteRequestWithFSO.java @@ -20,6 +20,7 @@ import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor.ONE; import static org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.Status; import static org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.Type.DeleteKeys; +import static org.mockito.Mockito.when; import java.util.ArrayList; import java.util.List; @@ -105,8 +106,23 @@ protected void createPreRequisites() throws Exception { createPreRequisites(RequestSource.USER); } + @Override + protected OMKeysDeleteRequest getOmKeysDeleteRequest(OMRequest request) { + return new OmKeysDeleteRequestWithFSO(request, getBucketLayout()); + } + + @Override + protected long expectedNamespaceReleasedOnPartialDelete() { + // Only the leaf files sit under the shared parent directory created by createPreRequisites. + return 3L; + } + @Override protected void createPreRequisites(RequestSource sourceType) throws Exception { + // Directories are read back through getOMKeyInfoIfExists, which stamps them with the OM + // default; without it they cannot be written to the deleted directory table. + when(ozoneManager.getDefaultReplicationConfig()) + .thenReturn(RatisReplicationConfig.getInstance(ONE)); setDeleteKeyList(new ArrayList<>()); // Add volume, bucket and key entries to OM DB. OMRequestTestUtils From 8f82fe0e3199c36f1b5d9a8c930053c81cd88400 Mon Sep 17 00:00:00 2001 From: Chi-Hsuan Huang Date: Mon, 7 Sep 2026 22:12:32 +0800 Subject: [PATCH 08/14] HDDS-16161. Publish bucket usage from a copy in S3 multipart requests --- .../multipart/S3InitiateMultipartUploadRequestWithFSO.java | 5 ++++- .../request/s3/multipart/S3MultipartUploadAbortRequest.java | 5 ++++- .../s3/multipart/S3MultipartUploadCommitPartRequest.java | 6 +++++- 3 files changed, 13 insertions(+), 3 deletions(-) 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 5596be4ff395..e3c1003314cd 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 @@ -114,7 +114,7 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut // check if the directory already existed in OM checkDirectoryResult(keyName, pathInfoFSO.getDirectoryResult()); - final OmBucketInfo bucketInfo = getBucketInfo(omMetadataManager, + final OmBucketInfo bucketInfo = getBucketInfoForUpdate(omMetadataManager, volumeName, bucketName); // add all missing parents to dir table @@ -214,6 +214,9 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut omMetadataManager.getMultipartInfoTable().addCacheEntry( multipartKey, multipartKeyInfo, transactionLogIndex); + omMetadataManager.getBucketTable().addCacheEntry( + omMetadataManager.getBucketKey(volumeName, bucketName), bucketInfo, transactionLogIndex); + omClientResponse = new S3InitiateMultipartUploadResponseWithFSO( omResponse.setInitiateMultiPartUploadResponse( diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadAbortRequest.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadAbortRequest.java index 7f67856ff8b7..958c65251d79 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadAbortRequest.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadAbortRequest.java @@ -156,7 +156,7 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut OmKeyInfo omKeyInfo = omMetadataManager.getOpenKeyTable(getBucketLayout()) .get(multipartOpenKey); - omBucketInfo = getBucketInfo(omMetadataManager, volumeName, bucketName); + omBucketInfo = getBucketInfoForUpdate(omMetadataManager, volumeName, bucketName); if (omKeyInfo == null) { // In old env, OpenKeycleanupservice may have deleted key from openKeyTable leaving behind @@ -213,6 +213,9 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut .addCacheEntry(new CacheKey<>(multipartKey), CacheValue.get(trxnLogIndex)); + omMetadataManager.getBucketTable().addCacheEntry( + omMetadataManager.getBucketKey(volumeName, bucketName), omBucketInfo, trxnLogIndex); + omClientResponse = getOmClientResponse(ozoneManager, multipartKeyInfo, multipartKey, multipartOpenKey, omResponse, omBucketInfo, partsKeyInfoToDelete, partsTableKeysToDelete); 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 24fe698336a4..6385f0a80d82 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 @@ -260,7 +260,7 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut new CacheKey<>(openKey), CacheValue.get(trxnLogIndex)); - omBucketInfo = getBucketInfo(omMetadataManager, volumeName, bucketName); + omBucketInfo = getBucketInfoForUpdate(omMetadataManager, volumeName, bucketName); // This map should contain maximum of two entries // 1. Overwritten part @@ -313,6 +313,10 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut commitResponseBuilder.setETag(eTag); } omResponse.setCommitMultiPartUploadResponse(commitResponseBuilder); + + omMetadataManager.getBucketTable().addCacheEntry( + omMetadataManager.getBucketKey(volumeName, bucketName), omBucketInfo, trxnLogIndex); + omClientResponse = getOmClientResponse(ozoneManager, keyVersionsToDeleteMap, openKey, omKeyInfo, multipartKey, multipartKeyInfo, multipartPartKey, From d8af8a8ff113c6f6a55181dbd3612d3b9669bfb5 Mon Sep 17 00:00:00 2001 From: Chi-Hsuan Huang Date: Mon, 7 Sep 2026 22:16:11 +0800 Subject: [PATCH 09/14] HDDS-16161. Publish bucket usage from a copy in purge requests --- .../key/OMDirectoriesPurgeRequestWithFSO.java | 27 +++++++- .../om/request/key/OMKeyPurgeRequest.java | 8 ++- .../key/TestOMKeyPurgeRequestAndResponse.java | 65 +++++++++++++++++++ 3 files changed, 94 insertions(+), 6 deletions(-) diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMDirectoriesPurgeRequestWithFSO.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMDirectoriesPurgeRequestWithFSO.java index 128e13dd7ec0..e2ca3a76136c 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMDirectoriesPurgeRequestWithFSO.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMDirectoriesPurgeRequestWithFSO.java @@ -42,6 +42,7 @@ import org.apache.hadoop.ozone.audit.AuditLoggerType; import org.apache.hadoop.ozone.audit.OMSystemAction; import org.apache.hadoop.ozone.om.DeletingServiceMetrics; +import org.apache.hadoop.ozone.om.OMMetadataManager; import org.apache.hadoop.ozone.om.OMMetadataManager.VolumeBucketId; import org.apache.hadoop.ozone.om.OMMetrics; import org.apache.hadoop.ozone.om.OmMetadataManagerImpl; @@ -147,7 +148,7 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut omMetrics.decNumKeys(); omMetrics.incNumKeyDeletesInternal(); - OmBucketInfo omBucketInfo = getBucketInfo(omMetadataManager, + OmBucketInfo omBucketInfo = accumulatedBucketInfo(volBucketInfoMap, omMetadataManager, processed.volumeName, processed.bucketName); // bucketInfo can be null in case of delete volume or bucket // or key does not belong to bucket as bucket is recreated @@ -184,7 +185,7 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut omMetrics.decNumKeys(); omMetrics.incNumKeyDeletesInternal(); numSubFilesMoved++; - OmBucketInfo omBucketInfo = getBucketInfo(omMetadataManager, + OmBucketInfo omBucketInfo = accumulatedBucketInfo(volBucketInfoMap, omMetadataManager, processed.volumeName, processed.bucketName); // bucketInfo can be null in case of delete volume or bucket // or key does not belong to bucket as bucket is recreated @@ -205,7 +206,7 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut deletedDirNames.add(path.getDeletedDir()); BucketNameInfo bucketNameInfo = volumeBucketIdMap.get(new VolumeBucketId(path.getVolumeId(), path.getBucketId())); - OmBucketInfo omBucketInfo = getBucketInfo(omMetadataManager, + OmBucketInfo omBucketInfo = accumulatedBucketInfo(volBucketInfoMap, omMetadataManager, bucketNameInfo.getVolumeName(), bucketNameInfo.getBucketName()); if (omBucketInfo != null && omBucketInfo.getObjectID() == path.getBucketId()) { omBucketInfo.purgeSnapshotUsedNamespace(1); @@ -222,6 +223,14 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut deletingServiceMetrics.incrNumSubFilesMoved(numSubFilesMoved); deletingServiceMetrics.incrNumDirPurged(numDirsDeleted); + // All bucket locks are still held here, and every purge path has been processed, so the + // accumulated copies can be published as a whole. + for (Map.Entry, OmBucketInfo> entry : volBucketInfoMap.entrySet()) { + omMetadataManager.getBucketTable().addCacheEntry( + omMetadataManager.getBucketKey(entry.getKey().getLeft(), entry.getKey().getRight()), + entry.getValue(), context.getIndex()); + } + TransactionInfo transactionInfo = TransactionInfo.valueOf(context.getTermIndex()); if (fromSnapshotInfo != null) { fromSnapshotInfo.setLastTransactionInfo(transactionInfo.toByteString()); @@ -269,6 +278,18 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut getBucketLayout(), volBucketInfoMap, fromSnapshotInfo, openKeyInfoMap); } + /** + * Return the copy this request accumulates deltas on, fetching it on first use. All three purge + * paths can touch the same bucket, so they must share one copy or earlier deltas are lost. + */ + private OmBucketInfo accumulatedBucketInfo(Map, OmBucketInfo> volBucketInfoMap, + OMMetadataManager omMetadataManager, String volumeName, String bucketName) { + OmBucketInfo omBucketInfo = volBucketInfoMap.get(Pair.of(volumeName, bucketName)); + + return omBucketInfo != null ? omBucketInfo + : getBucketInfoForUpdate(omMetadataManager, volumeName, bucketName); + } + /** * Helper class to hold processed key information. */ diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyPurgeRequest.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyPurgeRequest.java index d4da86ef0909..01d009f3cc78 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyPurgeRequest.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyPurgeRequest.java @@ -142,7 +142,7 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut deletingServiceMetrics.setLastAOSTransactionInfo(transactionInfo); } List bucketInfoList = updateBucketSize(purgeKeysRequest.getBucketPurgeKeysSizeList(), - omMetadataManager); + omMetadataManager, context.getIndex()); if (LOG.isDebugEnabled()) { Map auditParams = new LinkedHashMap<>(); @@ -168,7 +168,7 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut } private List updateBucketSize(List bucketPurgeKeysSizeList, - OMMetadataManager omMetadataManager) throws OMException { + OMMetadataManager omMetadataManager, long trxnLogIndex) throws OMException { Map>> bucketPurgeKeysSizes = new HashMap<>(); List bucketKeyList = new ArrayList<>(); for (BucketPurgeKeysSize bucketPurgeKey : bucketPurgeKeysSizeList) { @@ -192,7 +192,7 @@ private List updateBucketSize(List bucketPurg String volumeName = volEntry.getKey(); for (Map.Entry> bucketEntry : volEntry.getValue().entrySet()) { String bucketName = bucketEntry.getKey(); - OmBucketInfo omBucketInfo = getBucketInfo(omMetadataManager, volumeName, bucketName); + OmBucketInfo omBucketInfo = getBucketInfoForUpdate(omMetadataManager, volumeName, bucketName); // Check null if bucket has been deleted. if (omBucketInfo != null) { boolean bucketUpdated = false; @@ -205,6 +205,8 @@ private List updateBucketSize(List bucketPurg } } if (bucketUpdated) { + omMetadataManager.getBucketTable().addCacheEntry( + omMetadataManager.getBucketKey(volumeName, bucketName), omBucketInfo, trxnLogIndex); bucketInfoList.add(omBucketInfo.copyObject()); } } diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyPurgeRequestAndResponse.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyPurgeRequestAndResponse.java index 1f05105d39f0..3d6a58dc1bdb 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyPurgeRequestAndResponse.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyPurgeRequestAndResponse.java @@ -40,12 +40,16 @@ import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.utils.TransactionInfo; import org.apache.hadoop.hdds.utils.db.BatchOperation; +import org.apache.hadoop.hdds.utils.db.cache.CacheKey; import org.apache.hadoop.ozone.om.OmSnapshot; import org.apache.hadoop.ozone.om.helpers.OmBucketInfo; import org.apache.hadoop.ozone.om.helpers.SnapshotInfo; import org.apache.hadoop.ozone.om.lock.IOzoneManagerLock; import org.apache.hadoop.ozone.om.request.OMRequestTestUtils; +import org.apache.hadoop.ozone.om.response.OMClientResponse; import org.apache.hadoop.ozone.om.response.key.OMKeyPurgeResponse; +import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.BucketNameInfo; +import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.BucketPurgeKeysSize; import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.DeletedKeys; import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.OMRequest; import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.OMResponse; @@ -137,6 +141,67 @@ private OMRequest preExecute(OMRequest originalOmRequest) throws IOException { return modifiedOmRequest; } + @Test + public void testPurgedSizesReachCacheAndDb() throws Exception { + Pair, List> deleteKeysAndRenamedEntry = + createAndDeleteKeysAndRenamedEntry(1, null); + + String bucketKey = omMetadataManager.getBucketKey(volumeName, bucketName); + OmBucketInfo bucketInfo = omMetadataManager.getBucketTable() + .getCacheValue(new CacheKey<>(bucketKey)).getCacheValue(); + bucketInfo.incrSnapshotUsedBytes(500L); + bucketInfo.incrSnapshotUsedNamespace(5L); + + BucketNameInfo bucketNameInfo = BucketNameInfo.newBuilder() + .setVolumeName(volumeName) + .setBucketName(bucketName) + .setBucketId(bucketInfo.getObjectID()) + .build(); + // A bucket id that does not match must leave the counters alone. + BucketNameInfo staleBucketNameInfo = BucketNameInfo.newBuilder() + .setVolumeName(volumeName) + .setBucketName(bucketName) + .setBucketId(bucketInfo.getObjectID() + 1) + .build(); + + // Two entries for the same bucket, so the deltas have to accumulate on one copy. + PurgeKeysRequest purgeKeysRequest = PurgeKeysRequest.newBuilder() + .addDeletedKeys(DeletedKeys.newBuilder() + .setVolumeName(volumeName) + .setBucketName(bucketName) + .addAllKeys(deleteKeysAndRenamedEntry.getKey())) + .addAllRenamedKeys(deleteKeysAndRenamedEntry.getValue()) + .addBucketPurgeKeysSize(BucketPurgeKeysSize.newBuilder() + .setBucketNameInfo(bucketNameInfo).setPurgedBytes(100L).setPurgedNamespace(1L)) + .addBucketPurgeKeysSize(BucketPurgeKeysSize.newBuilder() + .setBucketNameInfo(bucketNameInfo).setPurgedBytes(200L).setPurgedNamespace(2L)) + .addBucketPurgeKeysSize(BucketPurgeKeysSize.newBuilder() + .setBucketNameInfo(staleBucketNameInfo).setPurgedBytes(999L).setPurgedNamespace(9L)) + .build(); + + OMRequest omRequest = OMRequest.newBuilder() + .setPurgeKeysRequest(purgeKeysRequest) + .setCmdType(Type.PurgeKeys) + .setClientId(UUID.randomUUID().toString()) + .build(); + + OMClientResponse omClientResponse = + new OMKeyPurgeRequest(preExecute(omRequest)).validateAndUpdateCache(ozoneManager, 100L); + assertEquals(Status.OK, omClientResponse.getOMResponse().getStatus()); + + OmBucketInfo cached = omMetadataManager.getBucketTable().get(bucketKey); + assertEquals(500L - 300L, cached.getSnapshotUsedBytes()); + assertEquals(5L - 3L, cached.getSnapshotUsedNamespace()); + + try (BatchOperation batchOperation = omMetadataManager.getStore().initBatchOperation()) { + omClientResponse.checkAndUpdateDB(omMetadataManager, batchOperation); + omMetadataManager.getStore().commitBatchOperation(batchOperation); + } + OmBucketInfo durable = omMetadataManager.getBucketTable().getSkipCache(bucketKey); + assertEquals(cached.getSnapshotUsedBytes(), durable.getSnapshotUsedBytes()); + assertEquals(cached.getSnapshotUsedNamespace(), durable.getSnapshotUsedNamespace()); + } + @Test public void testValidateAndUpdateCache() throws Exception { // Create and Delete keys. The keys should be moved to DeletedKeys table From e105f720851c19298b16ac5e8ee15b8be2b1c1ca Mon Sep 17 00:00:00 2001 From: Chi-Hsuan Huang Date: Mon, 7 Sep 2026 22:20:09 +0800 Subject: [PATCH 10/14] HDDS-16161. Hold all bucket locks while aborting expired multipart uploads --- ...S3ExpiredMultipartUploadsAbortRequest.java | 329 ++++++++++-------- .../request/key/TestOMKeyCreateRequest.java | 4 +- ...S3ExpiredMultipartUploadsAbortRequest.java | 141 +++++++- 3 files changed, 302 insertions(+), 172 deletions(-) diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3ExpiredMultipartUploadsAbortRequest.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3ExpiredMultipartUploadsAbortRequest.java index cc28954ef385..6c9077ff2c0f 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3ExpiredMultipartUploadsAbortRequest.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3ExpiredMultipartUploadsAbortRequest.java @@ -25,6 +25,8 @@ import java.util.List; import java.util.Map; import java.util.SortedMap; +import java.util.stream.Collectors; +import org.apache.commons.lang3.tuple.Pair; import org.apache.hadoop.hdds.utils.db.cache.CacheKey; import org.apache.hadoop.hdds.utils.db.cache.CacheValue; import org.apache.hadoop.ozone.OzoneConsts; @@ -106,18 +108,44 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut Result result = null; Map> abortedMultipartUploads = new HashMap<>(); + // One accumulated copy per bucket, so a bucket listed more than once keeps a single entry. + Map, OmBucketInfo> bucketInfoMap = new HashMap<>(); + + OMMetadataManager omMetadataManager = ozoneManager.getMetadataManager(); + List bucketLockKeys = submittedExpiredMPUsPerBucket.stream() + .map(mpuByBucket -> Pair.of(mpuByBucket.getVolumeName(), mpuByBucket.getBucketName())) + .distinct() + .map(volBucketPair -> new String[]{volBucketPair.getLeft(), volBucketPair.getRight()}) + .collect(Collectors.toList()); + boolean acquiredLocks = false; try { + // Hold every bucket lock for the whole request, so the accumulated copies can be published + // once all buckets have been processed. A later bucket failing turns the whole request into + // an error response, which persists nothing. + mergeOmLockDetails(omMetadataManager.getLock().acquireWriteLocks(BUCKET_LOCK, bucketLockKeys)); + acquiredLocks = getOmLockDetails().isLockAcquired(); + for (ExpiredMultipartUploadsBucket mpuByBucket: submittedExpiredMPUsPerBucket) { - // For each bucket where the MPU will be aborted from, - // get its bucket lock and update the cache accordingly. updateTableCache(ozoneManager, trxnLogIndex, mpuByBucket, - abortedMultipartUploads); + abortedMultipartUploads, bucketInfoMap); } + for (Map.Entry, OmBucketInfo> entry : bucketInfoMap.entrySet()) { + omMetadataManager.getBucketTable().addCacheEntry( + omMetadataManager.getBucketKey(entry.getKey().getLeft(), entry.getKey().getRight()), + entry.getValue(), trxnLogIndex); + } + + // Hand the response its own copies, so the instances now in the cache are not shared with + // the DB batch. + Map> abortedMPUsToPersist = + abortedMultipartUploads.entrySet().stream() + .collect(Collectors.toMap(entry -> entry.getKey().copyObject(), Map.Entry::getValue)); + omClientResponse = new S3ExpiredMultipartUploadsAbortResponse( - omResponse.build(), abortedMultipartUploads + omResponse.build(), abortedMPUsToPersist ); result = Result.SUCCESS; @@ -128,6 +156,9 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut new S3ExpiredMultipartUploadsAbortResponse(createErrorOMResponse( omResponse, exception)); } finally { + if (acquiredLocks) { + mergeOmLockDetails(omMetadataManager.getLock().releaseWriteLocks(BUCKET_LOCK, bucketLockKeys)); + } if (omClientResponse != null) { omClientResponse.setOmLockDetails(getOmLockDetails()); } @@ -194,171 +225,161 @@ private void processResults(OMMetrics omMetrics, @SuppressWarnings("methodlength") private void updateTableCache(OzoneManager ozoneManager, long trxnLogIndex, ExpiredMultipartUploadsBucket mpusPerBucket, - Map> abortedMultipartUploads) + Map> abortedMultipartUploads, + Map, OmBucketInfo> bucketInfoMap) throws IOException { - boolean acquiredLock = false; String volumeName = mpusPerBucket.getVolumeName(); String bucketName = mpusPerBucket.getBucketName(); OMMetadataManager omMetadataManager = ozoneManager.getMetadataManager(); - OmBucketInfo omBucketInfo = null; - BucketLayout bucketLayout = null; - try { - mergeOmLockDetails(omMetadataManager.getLock() - .acquireWriteLock(BUCKET_LOCK, volumeName, bucketName)); - acquiredLock = getOmLockDetails().isLockAcquired(); - omBucketInfo = getBucketInfo(omMetadataManager, volumeName, bucketName); + OmBucketInfo omBucketInfo = bucketInfoMap.computeIfAbsent( + Pair.of(volumeName, bucketName), + volBucketPair -> getBucketInfoForUpdate(omMetadataManager, + volBucketPair.getLeft(), volBucketPair.getRight())); - if (omBucketInfo == null) { - LOG.warn("Volume: {}, Bucket: {} does not exist, skipping deletion.", - volumeName, bucketName); - return; - } + if (omBucketInfo == null) { + LOG.warn("Volume: {}, Bucket: {} does not exist, skipping deletion.", + volumeName, bucketName); + return; + } - // Do not use getBucketLayout since the expired MPUs request might - // contains MPUs from all kind of buckets - bucketLayout = omBucketInfo.getBucketLayout(); - - for (ExpiredMultipartUploadInfo expiredMPU: - mpusPerBucket.getMultipartUploadsList()) { - String expiredMPUKeyName = expiredMPU.getName(); - - // If the MPU key is no longer present in the table, MPU - // might have been completed / aborted, and should not be - // aborted. - OmMultipartKeyInfo omMultipartKeyInfo = - omMetadataManager.getMultipartInfoTable().get(expiredMPUKeyName); - - if (omMultipartKeyInfo != null) { - if (trxnLogIndex < omMultipartKeyInfo.getUpdateID()) { - LOG.warn("Transaction log index {} is smaller than " + - "the current updateID {} of MPU key {}, skipping deletion.", - trxnLogIndex, omMultipartKeyInfo.getUpdateID(), - expiredMPUKeyName); - continue; - } + // Do not use getBucketLayout since the expired MPUs request might + // contains MPUs from all kind of buckets + BucketLayout bucketLayout = omBucketInfo.getBucketLayout(); + + for (ExpiredMultipartUploadInfo expiredMPU: + mpusPerBucket.getMultipartUploadsList()) { + String expiredMPUKeyName = expiredMPU.getName(); + + // If the MPU key is no longer present in the table, MPU + // might have been completed / aborted, and should not be + // aborted. + OmMultipartKeyInfo omMultipartKeyInfo = + omMetadataManager.getMultipartInfoTable().get(expiredMPUKeyName); + + if (omMultipartKeyInfo != null) { + if (trxnLogIndex < omMultipartKeyInfo.getUpdateID()) { + LOG.warn("Transaction log index {} is smaller than " + + "the current updateID {} of MPU key {}, skipping deletion.", + trxnLogIndex, omMultipartKeyInfo.getUpdateID(), + expiredMPUKeyName); + continue; + } - // Set the UpdateID to current transactionLogIndex - omMultipartKeyInfo = omMultipartKeyInfo.toBuilder() - .setUpdateID(trxnLogIndex) - .build(); - - // Parse the multipart upload components (e.g. volume, bucket, key) - // from the multipartInfoTable db key - - OmMultipartUpload multipartUpload; - try { - multipartUpload = - OmMultipartUpload.from(expiredMPUKeyName); - } catch (IllegalArgumentException e) { - LOG.warn("Aborting expired MPU failed: MPU key: " + - expiredMPUKeyName + " has invalid structure, " + - "skipping this MPU."); - continue; - } + // Set the UpdateID to current transactionLogIndex + omMultipartKeyInfo = omMultipartKeyInfo.toBuilder() + .setUpdateID(trxnLogIndex) + .build(); - String multipartOpenKey; - try { - multipartOpenKey = - OMMultipartUploadUtils - .getMultipartOpenKey(multipartUpload.getVolumeName(), - multipartUpload.getBucketName(), - multipartUpload.getKeyName(), - multipartUpload.getUploadId(), omMetadataManager, - bucketLayout); - } catch (OMException ome) { - LOG.warn("Aborting expired MPU Failed: volume: " + - multipartUpload.getVolumeName() + ", bucket: " + - multipartUpload.getBucketName() + ", key: " + - multipartUpload.getKeyName() + ". Cannot parse the open key" + - "for this MPU, skipping this MPU."); - continue; - } + // Parse the multipart upload components (e.g. volume, bucket, key) + // from the multipartInfoTable db key + + OmMultipartUpload multipartUpload; + try { + multipartUpload = + OmMultipartUpload.from(expiredMPUKeyName); + } catch (IllegalArgumentException e) { + LOG.warn("Aborting expired MPU failed: MPU key: " + + expiredMPUKeyName + " has invalid structure, " + + "skipping this MPU."); + continue; + } - // When abort uploaded key, we need to subtract the PartKey length - // from the volume usedBytes. - long quotaReleased = 0; - long numParts; - List partsKeyInfoToDelete = new ArrayList<>(); - List partsTableKeysToDelete = new ArrayList<>(); - if (omMultipartKeyInfo.getSchemaVersion() - == OmMultipartKeyInfo.LEGACY_SCHEMA_VERSION) { - for (PartKeyInfo iterPartKeyInfo : omMultipartKeyInfo. - getPartKeyInfoMap()) { - quotaReleased += QuotaUtil.getReplicatedSize( - iterPartKeyInfo.getPartKeyInfo().getDataSize(), - omMultipartKeyInfo.getReplicationConfig()); - } - numParts = omMultipartKeyInfo.getPartKeyInfoMap().size(); - } else { - SortedMap tableParts = - OMMultipartUploadUtils.scanParts(omMetadataManager, - multipartUpload.getUploadId()); - quotaReleased += OMMultipartUploadUtils.getReplicatedSize( - tableParts, omMultipartKeyInfo.getReplicationConfig()); - partsKeyInfoToDelete.addAll(OMMultipartUploadUtils.toOmKeyInfoList( - tableParts, multipartUpload.getVolumeName(), - multipartUpload.getBucketName(), multipartUpload.getKeyName(), - omMultipartKeyInfo.getReplicationConfig())); - partsTableKeysToDelete.addAll(OMMultipartUploadUtils.getPartKeys( - multipartUpload.getUploadId(), tableParts)); - OMMultipartUploadUtils.addPartCleanupCacheEntries(omMetadataManager, - partsTableKeysToDelete, trxnLogIndex); - numParts = tableParts.size(); - } - omBucketInfo.incrUsedBytes(-quotaReleased); - - OmMultipartAbortInfo omMultipartAbortInfo = - new OmMultipartAbortInfo.Builder() - .setMultipartKey(expiredMPUKeyName) - .setMultipartOpenKey(multipartOpenKey) - .setMultipartKeyInfo(omMultipartKeyInfo) - .setBucketLayout(omBucketInfo.getBucketLayout()) - .setPartsKeyInfoToDelete(partsKeyInfoToDelete) - .setPartsTableKeysToDelete(partsTableKeysToDelete) - .build(); - - abortedMultipartUploads.computeIfAbsent(omBucketInfo, - k -> new ArrayList<>()).add(omMultipartAbortInfo); - - // Update cache of openKeyTable and multipartInfo table. - // No need to add the cache entries to delete table, as the entries - // in delete table are not used by any read/write operations. - - // Unlike normal MPU abort request where the MPU open keys needs - // to exist. For OpenKeyCleanupService run prior to - // HDDS-9017, these MPU open keys might already be deleted, - // causing "orphan" MPU keys (MPU entry exist in - // multipartInfoTable, but not in openKeyTable). - // We can skip this existence check and just delete the - // multipartInfoTable. The existence check can be re-added - // once there are no "orphan" keys - if (omMetadataManager.getOpenKeyTable(bucketLayout) - .isExist(multipartOpenKey)) { - omMetadataManager.getOpenKeyTable(bucketLayout) - .addCacheEntry(new CacheKey<>(multipartOpenKey), - CacheValue.get(trxnLogIndex)); - } - omMetadataManager.getMultipartInfoTable() - .addCacheEntry(new CacheKey<>(expiredMPUKeyName), - CacheValue.get(trxnLogIndex)); + String multipartOpenKey; + try { + multipartOpenKey = + OMMultipartUploadUtils + .getMultipartOpenKey(multipartUpload.getVolumeName(), + multipartUpload.getBucketName(), + multipartUpload.getKeyName(), + multipartUpload.getUploadId(), omMetadataManager, + bucketLayout); + } catch (OMException ome) { + LOG.warn("Aborting expired MPU Failed: volume: " + + multipartUpload.getVolumeName() + ", bucket: " + + multipartUpload.getBucketName() + ", key: " + + multipartUpload.getKeyName() + ". Cannot parse the open key" + + "for this MPU, skipping this MPU."); + continue; + } - ozoneManager.getMetrics().incNumExpiredMPUAborted(); - ozoneManager.getMetrics().incNumExpiredMPUPartsAborted(numParts); - LOG.debug("Expired MPU {} aborted containing {} parts.", - expiredMPUKeyName, numParts); + // When abort uploaded key, we need to subtract the PartKey length + // from the volume usedBytes. + long quotaReleased = 0; + long numParts; + List partsKeyInfoToDelete = new ArrayList<>(); + List partsTableKeysToDelete = new ArrayList<>(); + if (omMultipartKeyInfo.getSchemaVersion() + == OmMultipartKeyInfo.LEGACY_SCHEMA_VERSION) { + for (PartKeyInfo iterPartKeyInfo : omMultipartKeyInfo. + getPartKeyInfoMap()) { + quotaReleased += QuotaUtil.getReplicatedSize( + iterPartKeyInfo.getPartKeyInfo().getDataSize(), + omMultipartKeyInfo.getReplicationConfig()); + } + numParts = omMultipartKeyInfo.getPartKeyInfoMap().size(); } else { - LOG.debug("MPU key {} was not aborted, as it was not " + - "found in the multipart info table", expiredMPUKeyName); + SortedMap tableParts = + OMMultipartUploadUtils.scanParts(omMetadataManager, + multipartUpload.getUploadId()); + quotaReleased += OMMultipartUploadUtils.getReplicatedSize( + tableParts, omMultipartKeyInfo.getReplicationConfig()); + partsKeyInfoToDelete.addAll(OMMultipartUploadUtils.toOmKeyInfoList( + tableParts, multipartUpload.getVolumeName(), + multipartUpload.getBucketName(), multipartUpload.getKeyName(), + omMultipartKeyInfo.getReplicationConfig())); + partsTableKeysToDelete.addAll(OMMultipartUploadUtils.getPartKeys( + multipartUpload.getUploadId(), tableParts)); + OMMultipartUploadUtils.addPartCleanupCacheEntries(omMetadataManager, + partsTableKeysToDelete, trxnLogIndex); + numParts = tableParts.size(); } - } - } finally { - if (acquiredLock) { - mergeOmLockDetails(omMetadataManager.getLock() - .releaseWriteLock(BUCKET_LOCK, volumeName, bucketName)); + omBucketInfo.incrUsedBytes(-quotaReleased); + + OmMultipartAbortInfo omMultipartAbortInfo = + new OmMultipartAbortInfo.Builder() + .setMultipartKey(expiredMPUKeyName) + .setMultipartOpenKey(multipartOpenKey) + .setMultipartKeyInfo(omMultipartKeyInfo) + .setBucketLayout(omBucketInfo.getBucketLayout()) + .setPartsKeyInfoToDelete(partsKeyInfoToDelete) + .setPartsTableKeysToDelete(partsTableKeysToDelete) + .build(); + + abortedMultipartUploads.computeIfAbsent(omBucketInfo, + k -> new ArrayList<>()).add(omMultipartAbortInfo); + + // Update cache of openKeyTable and multipartInfo table. + // No need to add the cache entries to delete table, as the entries + // in delete table are not used by any read/write operations. + + // Unlike normal MPU abort request where the MPU open keys needs + // to exist. For OpenKeyCleanupService run prior to + // HDDS-9017, these MPU open keys might already be deleted, + // causing "orphan" MPU keys (MPU entry exist in + // multipartInfoTable, but not in openKeyTable). + // We can skip this existence check and just delete the + // multipartInfoTable. The existence check can be re-added + // once there are no "orphan" keys + if (omMetadataManager.getOpenKeyTable(bucketLayout) + .isExist(multipartOpenKey)) { + omMetadataManager.getOpenKeyTable(bucketLayout) + .addCacheEntry(new CacheKey<>(multipartOpenKey), + CacheValue.get(trxnLogIndex)); + } + omMetadataManager.getMultipartInfoTable() + .addCacheEntry(new CacheKey<>(expiredMPUKeyName), + CacheValue.get(trxnLogIndex)); + + ozoneManager.getMetrics().incNumExpiredMPUAborted(); + ozoneManager.getMetrics().incNumExpiredMPUPartsAborted(numParts); + LOG.debug("Expired MPU {} aborted containing {} parts.", + expiredMPUKeyName, numParts); + } else { + LOG.debug("MPU key {} was not aborted, as it was not " + + "found in the multipart info table", expiredMPUKeyName); } } - } } diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyCreateRequest.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyCreateRequest.java index 6fce4602a11d..353b04656e74 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyCreateRequest.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyCreateRequest.java @@ -69,10 +69,10 @@ import org.apache.hadoop.hdds.protocol.proto.HddsProtos.KeyValue; import org.apache.hadoop.hdds.scm.container.common.helpers.AllocatedBlock; import org.apache.hadoop.hdds.scm.container.common.helpers.ExcludeList; -import org.apache.hadoop.hdds.utils.db.BatchOperation; -import org.apache.hadoop.hdds.utils.db.cache.CacheKey; import org.apache.hadoop.hdds.scm.pipeline.Pipeline; import org.apache.hadoop.hdds.scm.pipeline.PipelineID; +import org.apache.hadoop.hdds.utils.db.BatchOperation; +import org.apache.hadoop.hdds.utils.db.cache.CacheKey; import org.apache.hadoop.ozone.OzoneAcl; import org.apache.hadoop.ozone.OzoneConsts; import org.apache.hadoop.ozone.om.OmConfig; diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3ExpiredMultipartUploadsAbortRequest.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3ExpiredMultipartUploadsAbortRequest.java index 031988ec0b05..871582e1c6b5 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3ExpiredMultipartUploadsAbortRequest.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3ExpiredMultipartUploadsAbortRequest.java @@ -25,6 +25,11 @@ import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertSame; import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.argThat; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.when; import java.util.ArrayList; import java.util.Arrays; @@ -38,10 +43,15 @@ import org.apache.hadoop.hdds.client.ReplicationFactor; import org.apache.hadoop.hdds.client.ReplicationType; import org.apache.hadoop.hdds.utils.UniqueId; +import org.apache.hadoop.hdds.utils.db.BatchOperation; +import org.apache.hadoop.hdds.utils.db.RocksDatabaseException; +import org.apache.hadoop.hdds.utils.db.Table; import org.apache.hadoop.hdds.utils.db.cache.CacheKey; import org.apache.hadoop.hdds.utils.db.cache.CacheValue; +import org.apache.hadoop.ozone.om.OMMetadataManager; import org.apache.hadoop.ozone.om.OMMetrics; import org.apache.hadoop.ozone.om.helpers.BucketLayout; +import org.apache.hadoop.ozone.om.helpers.OmBucketInfo; import org.apache.hadoop.ozone.om.helpers.OmKeyInfo; import org.apache.hadoop.ozone.om.helpers.OmKeyLocationInfoGroup; import org.apache.hadoop.ozone.om.helpers.OmMultipartKeyInfo; @@ -57,6 +67,7 @@ import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.OMRequest; import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.Status; import org.apache.hadoop.security.UserGroupInformation; +import org.junit.jupiter.api.Test; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.MethodSource; @@ -360,33 +371,121 @@ private void abortExpiredMPUsFromCache(String volumeName, String bucketName, assertEquals(Status.OK, omClientResponse.getOMResponse().getStatus()); } - private OMRequest createAbortExpiredMPURequest(String volumeName, - String bucketName, List mpuKeysToAbort) { + @Test + public void testLaterBucketFailureLeavesEarlierBucketUsageAlone() throws Exception { + this.bucketLayout = BucketLayout.DEFAULT; + final String volume = UUID.randomUUID().toString(); + final String bucket1 = UUID.randomUUID().toString(); + final String bucket2 = UUID.randomUUID().toString(); + OMRequestTestUtils.addVolumeAndBucketToDB(volume, bucket1, omMetadataManager, getBucketLayout()); + OMRequestTestUtils.addVolumeAndBucketToDB(volume, bucket2, omMetadataManager, getBucketLayout()); + + List bucket1MPUs = createMPUs(volume, bucket1, null, 1, 1, 100L); + List bucket2MPUs = createMPUs(volume, bucket2, null, 1, 1, 100L); + + String bucket1Key = omMetadataManager.getBucketKey(volume, bucket1); + OmBucketInfo bucket1CachedBefore = omMetadataManager.getBucketTable() + .getCacheValue(new CacheKey<>(bucket1Key)).getCacheValue(); + long bucket1UsedBytesBefore = bucket1CachedBefore.getUsedBytes(); + assertNotEquals(0L, bucket1UsedBytesBefore); + + // The second bucket's scan fails on a real table read, after the first bucket was processed. + Table multipartInfoTable = + spy(omMetadataManager.getMultipartInfoTable()); + doThrow(new RocksDatabaseException("injected")).when(multipartInfoTable) + .get(argThat(key -> key != null && key.contains(bucket2))); + OMMetadataManager spiedMetadataManager = spy(omMetadataManager); + doReturn(multipartInfoTable).when(spiedMetadataManager).getMultipartInfoTable(); + when(ozoneManager.getMetadataManager()).thenReturn(spiedMetadataManager); - List expiredMultipartUploads = mpuKeysToAbort - .stream().map(name -> - ExpiredMultipartUploadInfo.newBuilder().setName(name).build()) - .collect(Collectors.toList()); - ExpiredMultipartUploadsBucket expiredMultipartUploadsBucket = - ExpiredMultipartUploadsBucket.newBuilder() - .setVolumeName(volumeName) - .setBucketName(bucketName) - .addAllMultipartUploads(expiredMultipartUploads) - .build(); + OMRequest omRequest = doPreExecute(createAbortExpiredMPURequest( + expiredMPUsForBucket(volume, bucket1, bucket1MPUs), + expiredMPUsForBucket(volume, bucket2, bucket2MPUs))); + OMClientResponse omClientResponse = + new S3ExpiredMultipartUploadsAbortRequest(omRequest) + .validateAndUpdateCache(ozoneManager, 100L); + + assertNotEquals(Status.OK, omClientResponse.getOMResponse().getStatus()); + // Nothing is persisted for a failed request, so the first bucket's release must not reach the + // cache. The value catches an in-place update of the cached instance, the identity catches a + // published copy whose counters happen to match. + OmBucketInfo bucket1CachedAfter = omMetadataManager.getBucketTable() + .getCacheValue(new CacheKey<>(bucket1Key)).getCacheValue(); + assertEquals(bucket1UsedBytesBefore, bucket1CachedAfter.getUsedBytes()); + assertSame(bucket1CachedBefore, bucket1CachedAfter); + } + @Test + public void testRepeatedBucketAccumulatesIntoOneEntry() throws Exception { + this.bucketLayout = BucketLayout.DEFAULT; + final String volume = UUID.randomUUID().toString(); + final String bucket = UUID.randomUUID().toString(); + OMRequestTestUtils.addVolumeAndBucketToDB(volume, bucket, omMetadataManager, getBucketLayout()); + + List firstMPUs = createMPUs(volume, bucket, null, 1, 1, 100L); + List secondMPUs = createMPUs(volume, bucket, null, 1, 1, 100L); + + String bucketKey = omMetadataManager.getBucketKey(volume, bucket); + long usedBytesBefore = omMetadataManager.getBucketTable().get(bucketKey).getUsedBytes(); + + // The same bucket listed twice in one request must end up as a single accumulated entry. + OMRequest omRequest = doPreExecute(createAbortExpiredMPURequest( + expiredMPUsForBucket(volume, bucket, firstMPUs), + expiredMPUsForBucket(volume, bucket, secondMPUs))); + + OMClientResponse omClientResponse = + new S3ExpiredMultipartUploadsAbortRequest(omRequest) + .validateAndUpdateCache(ozoneManager, 100L); + assertEquals(Status.OK, omClientResponse.getOMResponse().getStatus()); + + assertNotInMultipartInfoTable(firstMPUs); + assertNotInMultipartInfoTable(secondMPUs); + + // Both listings share one accumulated copy, so the bucket row is written once and the cache + // cannot disagree with what is persisted. + long usedBytesCached = omMetadataManager.getBucketTable().get(bucketKey).getUsedBytes(); + // One 100-byte part per batch, charged at three-way replication. + assertEquals(600L, usedBytesBefore); + assertEquals(0L, usedBytesCached, "Both batches must release their part"); + + try (BatchOperation batchOperation = omMetadataManager.getStore().initBatchOperation()) { + omClientResponse.checkAndUpdateDB(omMetadataManager, batchOperation); + omMetadataManager.getStore().commitBatchOperation(batchOperation); + } + assertEquals(usedBytesCached, + omMetadataManager.getBucketTable().getSkipCache(bucketKey).getUsedBytes()); + } + + private ExpiredMultipartUploadsBucket expiredMPUsForBucket(String volumeName, + String bucketName, List mpuKeysToAbort) { + return ExpiredMultipartUploadsBucket.newBuilder() + .setVolumeName(volumeName) + .setBucketName(bucketName) + .addAllMultipartUploads(mpuKeysToAbort.stream() + .map(name -> ExpiredMultipartUploadInfo.newBuilder().setName(name).build()) + .collect(Collectors.toList())) + .build(); + } + + private OMRequest createAbortExpiredMPURequest(ExpiredMultipartUploadsBucket... buckets) { MultipartUploadsExpiredAbortRequest mpuExpiredAbortRequest = MultipartUploadsExpiredAbortRequest.newBuilder() - .addExpiredMultipartUploadsPerBucket(expiredMultipartUploadsBucket) + .addAllExpiredMultipartUploadsPerBucket(Arrays.asList(buckets)) .build(); return OMRequest.newBuilder() .setMultipartUploadsExpiredAbortRequest(mpuExpiredAbortRequest) - .setCmdType(OzoneManagerProtocolProtos - .Type.AbortExpiredMultiPartUploads) + .setCmdType(OzoneManagerProtocolProtos.Type.AbortExpiredMultiPartUploads) .setClientId(UUID.randomUUID().toString()) .build(); } + private OMRequest createAbortExpiredMPURequest(String volumeName, + String bucketName, List mpuKeysToAbort) { + return createAbortExpiredMPURequest( + expiredMPUsForBucket(volumeName, bucketName, mpuKeysToAbort)); + } + /** * Create MPus with randomized key name. */ @@ -515,6 +614,11 @@ private List createMPUsWithFSO(String volume, String bucket, */ private List createMPUs(String volume, String bucket, String key, int count, int numParts) throws Exception { + return createMPUs(volume, bucket, key, count, numParts, 0L); + } + + private List createMPUs(String volume, String bucket, + String key, int count, int numParts, long partSize) throws Exception { List mpuKeys = new ArrayList<>(); long trxnLogIndex = 1L; @@ -555,6 +659,10 @@ private List createMPUs(String volume, String bucket, long clientID = UniqueId.next(); OMRequest commitMultipartRequest = doPreExecuteCommitMPU( volume, bucket, keyName, clientID, multipartUploadID, j); + commitMultipartRequest = commitMultipartRequest.toBuilder().setCommitMultiPartUploadRequest( + commitMultipartRequest.getCommitMultiPartUploadRequest().toBuilder().setKeyArgs( + commitMultipartRequest.getCommitMultiPartUploadRequest().getKeyArgs().toBuilder() + .setDataSize(partSize))).build(); S3MultipartUploadCommitPartRequest s3MultipartUploadCommitPartRequest = getS3MultipartUploadCommitReq(commitMultipartRequest); @@ -562,7 +670,8 @@ private List createMPUs(String volume, String bucket, // Add key to open key table to be used in MPU commit processing OMRequestTestUtils.addKeyToTable( true, true, - volume, bucket, keyName, clientID, RatisReplicationConfig.getInstance(ONE), omMetadataManager); + volume, bucket, keyName, clientID, + omMetadataManager.getMultipartInfoTable().get(mpuKey).getReplicationConfig(), omMetadataManager); OMClientResponse commitResponse = s3MultipartUploadCommitPartRequest.validateAndUpdateCache( From f204b841622c43e4efb13943a589eec42a6f3979 Mon Sep 17 00:00:00 2001 From: Chi-Hsuan Huang Date: Mon, 7 Sep 2026 22:25:31 +0800 Subject: [PATCH 11/14] HDDS-16161. Use getBucketInfoForUpdate in S3MultipartUploadCompleteRequest --- .../multipart/S3MultipartUploadCompleteRequest.java | 11 +++-------- 1 file changed, 3 insertions(+), 8 deletions(-) diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadCompleteRequest.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadCompleteRequest.java index bc4c0c2c8739..33e1ee07f533 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadCompleteRequest.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadCompleteRequest.java @@ -175,15 +175,10 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut acquiredLock = getOmLockDetails().isLockAcquired(); validateBucketAndVolume(omMetadataManager, volumeName, bucketName); - // Work on a copy of the cached bucket so the namespace charge for - // recreating missing FSO parent directories (applied before parts are - // validated) is published only on success; a complete that fails with - // INVALID_PART must not leak it into the cache. See getBucketInfo. - OmBucketInfo omBucketInfo = getBucketInfo(omMetadataManager, + // The namespace charge for recreating missing FSO parent directories is applied before the + // parts are validated, so a complete that fails with INVALID_PART must not leak it. + OmBucketInfo omBucketInfo = getBucketInfoForUpdate(omMetadataManager, volumeName, bucketName); - if (omBucketInfo != null) { - omBucketInfo = omBucketInfo.copyObject(); - } List missingParentInfos; OMFileRequest.OMPathInfoWithFSO pathInfoFSO = OMFileRequest From 54ee3c71596d6dbf6354c6fa70e44079485ffb6e Mon Sep 17 00:00:00 2001 From: Chi-Hsuan Huang Date: Tue, 8 Sep 2026 00:16:58 +0800 Subject: [PATCH 12/14] HDDS-16161. Cover bucket usage publication in key delete and MPU tests Extend the existing success tests for OMKeyDeleteRequest, S3MultipartUploadAbortRequest and S3InitiateMultipartUploadRequestWithFSO so that dropping the addCacheEntry publication fails them. Verified by removing each publication and re-running the tests. --- .../ozone/om/request/key/TestOMKeyDeleteRequest.java | 8 ++++++++ .../TestS3InitiateMultipartUploadRequestWithFSO.java | 5 +++++ .../multipart/TestS3MultipartUploadAbortRequest.java | 10 ++++++++++ 3 files changed, 23 insertions(+) diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyDeleteRequest.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyDeleteRequest.java index 3878e3aca8b6..828d24bf2010 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyDeleteRequest.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyDeleteRequest.java @@ -24,6 +24,7 @@ import static org.junit.jupiter.api.Assertions.assertThrows; import java.util.UUID; +import org.apache.hadoop.hdds.utils.db.cache.CacheKey; import org.apache.hadoop.ozone.OzoneConsts; import org.apache.hadoop.ozone.om.exceptions.OMException; import org.apache.hadoop.ozone.om.helpers.BucketLayout; @@ -89,6 +90,12 @@ public void testValidateAndUpdateCache() throws Exception { OMKeyDeleteRequest omKeyDeleteRequest = getOmKeyDeleteRequest(modifiedOmRequest); + // The key was added straight to the table, so charge the bucket for it here; the delete + // releases it on a copy that only reaches the cache if the request publishes it. + String bucketKey = omMetadataManager.getBucketKey(volumeName, bucketName); + omMetadataManager.getBucketTable().getCacheValue(new CacheKey<>(bucketKey)) + .getCacheValue().incrUsedNamespace(1L); + OMClientResponse omClientResponse = omKeyDeleteRequest.validateAndUpdateCache(ozoneManager, 100L); @@ -99,6 +106,7 @@ public void testValidateAndUpdateCache() throws Exception { omKeyInfo = omMetadataManager.getKeyTable(getBucketLayout()).get(ozoneKey); assertNull(omKeyInfo); + assertEquals(0, omMetadataManager.getBucketTable().get(bucketKey).getUsedNamespace()); } @Test diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3InitiateMultipartUploadRequestWithFSO.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3InitiateMultipartUploadRequestWithFSO.java index d517be2bbba1..3c3d4c006f18 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3InitiateMultipartUploadRequestWithFSO.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/s3/multipart/TestS3InitiateMultipartUploadRequestWithFSO.java @@ -89,6 +89,11 @@ public void testValidateAndUpdateCache() throws Exception { long parentID = verifyDirectoriesInDB(dirs, volumeId, bucketId); + // The missing parents are charged on a copy, so they only show up here if the request + // published it. + assertEquals(dirs.size(), omMetadataManager.getBucketTable() + .get(omMetadataManager.getBucketKey(volumeName, bucketName)).getUsedNamespace()); + String multipartFileKey = omMetadataManager .getMultipartKey(volumeName, bucketName, keyName, modifiedRequest.getInitiateMultiPartUploadRequest().getKeyArgs() 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 3d529ad0acac..9c38af8c543d 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 @@ -19,12 +19,14 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNotSame; import static org.junit.jupiter.api.Assertions.assertNull; import java.io.IOException; 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.OmBucketInfo; import org.apache.hadoop.ozone.om.helpers.OmMultipartKeyInfo; import org.apache.hadoop.ozone.om.request.OMRequestTestUtils; import org.apache.hadoop.ozone.om.response.OMClientResponse; @@ -77,6 +79,10 @@ public void testValidateAndUpdateCache() throws Exception { S3MultipartUploadAbortRequest s3MultipartUploadAbortRequest = getS3MultipartUploadAbortReq(abortMPURequest); + String bucketKey = omMetadataManager.getBucketKey(volumeName, bucketName); + OmBucketInfo cachedBeforeAbort = omMetadataManager.getBucketTable() + .getCacheValue(new CacheKey<>(bucketKey)).getCacheValue(); + omClientResponse = s3MultipartUploadAbortRequest.validateAndUpdateCache(ozoneManager, 2L); @@ -94,6 +100,10 @@ public void testValidateAndUpdateCache() throws Exception { .getOpenKeyTable(s3MultipartUploadAbortRequest.getBucketLayout()) .get(multipartOpenKey)); + // No part was committed, so the released quota is zero and only the instance identity shows + // that the abort published the copy it released on. + assertNotSame(cachedBeforeAbort, omMetadataManager.getBucketTable() + .getCacheValue(new CacheKey<>(bucketKey)).getCacheValue()); } @Test From 09cd08e1d94457c846de5d05361596a24a915a39 Mon Sep 17 00:00:00 2001 From: Chi-Hsuan Huang Date: Tue, 8 Sep 2026 00:17:05 +0800 Subject: [PATCH 13/14] HDDS-16161. Tighten bucket usage assertions in existing failure tests The commit test asserted usedBytes, which is only charged after the injected failure point, so that half was vacuous. The partial delete test asserted a negative usedNamespace because the fixture never charged the bucket for the keys it added. --- .../ozone/om/request/key/TestOMKeyCommitRequest.java | 12 +++++------- .../om/request/key/TestOMKeysDeleteRequest.java | 5 +++++ 2 files changed, 10 insertions(+), 7 deletions(-) diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyCommitRequest.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyCommitRequest.java index fef410988328..b5152082c7d2 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyCommitRequest.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeyCommitRequest.java @@ -352,7 +352,8 @@ public void testAtomicCreateIfNotExistsCommitKeyAlreadyExists() throws Exception public void testFailedCommitDoesNotLeakBucketUsage() throws Exception { // Uncommitted blocks make the request build a pseudo key for deletion, which derives an // object id from the transaction index. An index above MAX_TRXN_ID makes that step throw - // after the bucket counters have already been changed. + // after incrUsedNamespace but before incrUsedBytes, so only the namespace charge is at + // risk of leaking here. List allocatedKeyLocationList = getKeyLocation(5); List allocatedBlockList = allocatedKeyLocationList .stream().map(OmKeyLocationInfo::getFromProtobuf) @@ -368,9 +369,7 @@ public void testFailedCommitDoesNotLeakBucketUsage() throws Exception { addKeyToOpenKeyTable(allocatedBlockList); String bucketKey = omMetadataManager.getBucketKey(volumeName, bucketName); - OmBucketInfo cachedBefore = omMetadataManager.getBucketTable().get(bucketKey); - long usedNamespaceBefore = cachedBefore.getUsedNamespace(); - long usedBytesBefore = cachedBefore.getUsedBytes(); + long usedNamespaceBefore = omMetadataManager.getBucketTable().get(bucketKey).getUsedNamespace(); // The mock returns 0 by default, so let it run the real check the OM performs. when(ozoneManager.getObjectIdFromTxId(anyLong())).thenAnswer(invocation -> OmUtils @@ -379,9 +378,8 @@ public void testFailedCommitDoesNotLeakBucketUsage() throws Exception { assertThrows(IllegalArgumentException.class, () -> omKeyCommitRequest .validateAndUpdateCache(ozoneManager, OmUtils.MAX_TRXN_ID + 1)); - OmBucketInfo cachedAfter = omMetadataManager.getBucketTable().get(bucketKey); - assertEquals(usedNamespaceBefore, cachedAfter.getUsedNamespace()); - assertEquals(usedBytesBefore, cachedAfter.getUsedBytes()); + assertEquals(usedNamespaceBefore, + omMetadataManager.getBucketTable().get(bucketKey).getUsedNamespace()); } @Test diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeysDeleteRequest.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeysDeleteRequest.java index 86a90a40e534..ac14439bae4f 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeysDeleteRequest.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMKeysDeleteRequest.java @@ -39,6 +39,7 @@ import java.util.UUID; import org.apache.hadoop.hdds.client.RatisReplicationConfig; import org.apache.hadoop.hdds.utils.db.BatchOperation; +import org.apache.hadoop.hdds.utils.db.cache.CacheKey; import org.apache.hadoop.ozone.om.OzoneManager; import org.apache.hadoop.ozone.om.exceptions.OMException; import org.apache.hadoop.ozone.om.request.OMRequestTestUtils; @@ -153,6 +154,10 @@ public void testPartialDeletePublishesBucketUsage() throws Exception { createPreRequisites(); String bucketKey = omMetadataManager.getBucketKey(volumeName, bucketName); + // createPreRequisites adds the keys straight to the table cache, so charge the bucket for + // them here; without it the released namespace would drive the counter negative. + omMetadataManager.getBucketTable().getCacheValue(new CacheKey<>(bucketKey)) + .getCacheValue().incrUsedNamespace(KEY_COUNT); long usedNamespaceBefore = omMetadataManager.getBucketTable().get(bucketKey).getUsedNamespace(); // One key does not exist, so the request succeeds partially. From 8f46ba08950aeb2854c032afde5e1f25e3bd0ad9 Mon Sep 17 00:00:00 2001 From: Chi-Hsuan Huang Date: Tue, 8 Sep 2026 00:17:05 +0800 Subject: [PATCH 14/14] HDDS-16161. Correct the read-modify-publish and cache javadoc getBucketInfoForUpdate stated the bucket write lock as an unconditional rule, but a caller under key path locking holds only the read lock and must not publish. The expired MPU abort request now also writes the multipart info and bucket table caches. --- .../apache/hadoop/ozone/om/request/key/OMKeyRequest.java | 7 ++++--- .../multipart/S3ExpiredMultipartUploadsAbortRequest.java | 4 ++-- 2 files changed, 6 insertions(+), 5 deletions(-) diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyRequest.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyRequest.java index 244c8c6638e3..d674976c1579 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyRequest.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyRequest.java @@ -967,9 +967,10 @@ public static OmBucketInfo getBucketInfo(OMMetadataManager omMetadataManager, *

* Mutations stay invisible until the caller publishes the copy with * {@code getBucketTable().addCacheEntry(...)}, after all fallible work and - * only on the path that persists the response. Hold the bucket write lock for - * the whole read-modify-publish, not just the publish, or a concurrent writer - * can be lost. + * only on the path that persists the response. A caller that publishes must + * hold the bucket write lock for the whole read-modify-publish, not just the + * publish, or a concurrent writer can be lost. A caller under key path + * locking holds only the bucket read lock, so it must not publish. */ @Nullable public static OmBucketInfo getBucketInfoForUpdate(OMMetadataManager omMetadataManager, diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3ExpiredMultipartUploadsAbortRequest.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3ExpiredMultipartUploadsAbortRequest.java index 6c9077ff2c0f..1d7b22aa3536 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3ExpiredMultipartUploadsAbortRequest.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3ExpiredMultipartUploadsAbortRequest.java @@ -61,8 +61,8 @@ /** * Handles requests to move both MPU open keys from the open key/file table and - * MPU part keys to delete table. Modifies the open key/file table cache only, - * and no underlying databases. + * MPU part keys to delete table. Modifies the open key/file, multipart info and + * bucket table caches, and no underlying databases. * The delete table cache does not need to be modified since it is not used * for client response validation. */