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/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 1fb1c22ebeba..5b91f66cf3fc 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(); 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/OMKeyCommitRequest.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCommitRequest.java index 6e4386aae5f8..24c2b6b3e253 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 @@ -195,7 +195,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()) { @@ -407,6 +407,9 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut omBucketInfo.incrUsedBytes(correctedSpace); + omMetadataManager.getBucketTable().addCacheEntry( + omMetadataManager.getBucketKey(volumeName, bucketName), omBucketInfo, trxnLogIndex); + omResponse.setCommitKeyResponse(CommitKeyResponse.newBuilder() .setModificationTime(commitKeyArgs.getModificationTime()) .build()); 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 8b3d32bdc1f7..8e21cc8a4156 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 @@ -130,7 +130,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"; @@ -351,6 +351,9 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut omBucketInfo.incrUsedBytes(correctedSpace); + omMetadataManager.getBucketTable().addCacheEntry( + omMetadataManager.getBucketKey(volumeName, bucketName), omBucketInfo, trxnLogIndex); + omResponse.setCommitKeyResponse(CommitKeyResponse.newBuilder() .setModificationTime(commitKeyArgs.getModificationTime()) .build()); 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 71d49df7bfcd..7c1056b51f56 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/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(); 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 1aa65f362544..6cab024cd964 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/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/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..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 @@ -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,24 @@ 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. 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, + 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 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 9a60ebbd391c..0d7ff10251df 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 @@ -237,7 +237,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. @@ -260,6 +260,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/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..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 @@ -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; @@ -59,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. */ @@ -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/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 bd7c738812ba..99111ff5e097 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 89d82cc1605a..e292b0284d42 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 @@ -314,6 +314,10 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut } commitResponseBuilder.setModificationTime(keyArgs.getModificationTime()); omResponse.setCommitMultiPartUploadResponse(commitResponseBuilder); + + omMetadataManager.getBucketTable().addCacheEntry( + omMetadataManager.getBucketKey(volumeName, bucketName), omBucketInfo, trxnLogIndex); + omClientResponse = getOmClientResponse(ozoneManager, keyVersionsToDeleteMap, openKey, omKeyInfo, multipartKey, multipartKeyInfo, multipartPartKey, 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 1eccefccb9a0..7133ff405fc4 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 @@ -177,15 +177,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 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"; 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..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 @@ -348,6 +348,40 @@ 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 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) + .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); + 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 + .getObjectIdFromTxId(omMetadataManager.getOmEpoch(), invocation.getArgument(0))); + + assertThrows(IllegalArgumentException.class, () -> omKeyCommitRequest + .validateAndUpdateCache(ozoneManager, OmUtils.MAX_TRXN_ID + 1)); + + assertEquals(usedNamespaceBefore, + omMetadataManager.getBucketTable().get(bucketKey).getUsedNamespace()); + } + @Test public void testValidateAndUpdateCacheWithUncommittedBlocks() throws Exception { 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..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 @@ -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; @@ -70,6 +71,8 @@ import org.apache.hadoop.hdds.scm.container.common.helpers.ExcludeList; 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; @@ -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( 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/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 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..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 @@ -38,6 +38,8 @@ 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.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; @@ -147,6 +149,48 @@ 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); + // 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. + 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 +326,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 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( 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