diff --git a/mooncake-store/include/client_service.h b/mooncake-store/include/client_service.h index 31b8cc48e7..631e4328bd 100644 --- a/mooncake-store/include/client_service.h +++ b/mooncake-store/include/client_service.h @@ -463,6 +463,9 @@ class Client { [[nodiscard]] tl::expected, ErrorCode> RemoveObjectHeartbeat(const UUID& client_id); + tl::expected AckRemoveObjectHeartbeat( + const UUID& client_id, const std::vector& tasks); + /** * @brief Stage a PROCESSING MEMORY replica for an existing key during * L2->L1 promotion. Returns the new replica's descriptor that the caller diff --git a/mooncake-store/include/file_storage.h b/mooncake-store/include/file_storage.h index 1805175f74..2439a2a082 100644 --- a/mooncake-store/include/file_storage.h +++ b/mooncake-store/include/file_storage.h @@ -54,8 +54,9 @@ class FileStorage { // Forward explicit-delete tombstone to the storage backend. // For BucketStorageBackend: marks tombstone + enables GC. // For other backends: no-op (default in StorageBackendInterface). - void MarkRemoved(const std::string& key); - void BatchMarkRemoved(const std::vector& keys); + tl::expected MarkRemoved(const std::string& key); + tl::expected BatchMarkRemoved( + const std::vector& keys); private: friend class FileStorageTest; diff --git a/mooncake-store/include/master_client.h b/mooncake-store/include/master_client.h index 8700dce44e..be6e1b9c24 100644 --- a/mooncake-store/include/master_client.h +++ b/mooncake-store/include/master_client.h @@ -493,14 +493,11 @@ class MasterClient { [[nodiscard]] tl::expected, ErrorCode> PromotionObjectHeartbeat(const UUID& client_id); - /** - * @brief Drain the removed_keys queue from master. Returns {tenant_id, - * key} pairs that were removed via Remove/BatchRemove and had LOCAL_DISK - * replicas on this client. The caller should MarkRemoved each key to - * trigger SSD tombstone + GC compaction. - */ + /** Fetch pending remove tasks without removing them from the queue. */ [[nodiscard]] tl::expected, ErrorCode> RemoveObjectHeartbeat(const UUID& client_id); + tl::expected AckRemoveObjectHeartbeat( + const UUID& client_id, const std::vector& tasks); /** * @brief Stage a PROCESSING MEMORY replica for an existing key during diff --git a/mooncake-store/include/master_service.h b/mooncake-store/include/master_service.h index f43217487b..f1f14e1990 100644 --- a/mooncake-store/include/master_service.h +++ b/mooncake-store/include/master_service.h @@ -727,16 +727,12 @@ class MasterService { auto PromotionObjectHeartbeat(const UUID& client_id) -> tl::expected, ErrorCode>; - /** - * @brief Drain the removed_keys queue for a client. - * - * Returns the list of {tenant_id, key} pairs that were removed via - * Remove/BatchRemove and had LOCAL_DISK replicas on this client. The - * client should call MarkRemoved on each key to trigger SSD tombstone - * marking + GC compaction. - */ + /** Fetch pending remove tasks without removing them from the queue. */ auto RemoveObjectHeartbeat(const UUID& client_id) -> tl::expected, ErrorCode>; + auto AckRemoveObjectHeartbeat( + const UUID& client_id, const std::vector& tasks) + -> tl::expected; /** * @brief Stage a PROCESSING MEMORY replica for an existing key. Allocates @@ -1544,6 +1540,8 @@ class MasterService { void FinalizeRemovedReplicasAfterDurable( const OpLogEntry& durable_entry, const std::vector& replica_ids, QuotaEraseMode quota_mode); + void EnqueueRemoveTasks(const std::vector& holder_ids, + const RemoveTaskItem& task); void FinalizeMetadataEraseAfterDurable(const OpLogEntry& durable_entry, QuotaEraseMode quota_mode); void FinalizeExpiredProcessingReplicasAfterDurable( diff --git a/mooncake-store/include/rpc_service.h b/mooncake-store/include/rpc_service.h index 2be2b87e9f..7abf3a34e1 100644 --- a/mooncake-store/include/rpc_service.h +++ b/mooncake-store/include/rpc_service.h @@ -219,13 +219,10 @@ class WrappedMasterService { tl::expected PollRemoveAll(const UUID& client_id); - tl::expected, ErrorCode> - RemoveObjectHeartbeat(const UUID& client_id); - - tl::expected, ErrorCode> - RemoveObjectHeartbeat(const UUID& client_id); - - tl::expected PollRemoveAll(const UUID& client_id); + tl::expected, ErrorCode> RemoveObjectHeartbeat( + const UUID& client_id); + tl::expected AckRemoveObjectHeartbeat( + const UUID& client_id, const std::vector& tasks); tl::expected ReportSsdCapacity( const UUID& client_id, int64_t ssd_total_capacity_bytes); diff --git a/mooncake-store/include/storage_backend.h b/mooncake-store/include/storage_backend.h index d9fe2584bb..b417c3c950 100644 --- a/mooncake-store/include/storage_backend.h +++ b/mooncake-store/include/storage_backend.h @@ -42,6 +42,10 @@ struct BucketMetadata { std::vector keys; std::vector metadatas; + // Persisted tombstones. Init filters these keys before rebuilding the + // in-memory object index. + std::vector tombstones; + // Runtime-only fields (not serialized) for safe deletion support // Tracks number of in-flight reads to enable safe bucket deletion mutable std::atomic inflight_reads_{0}; @@ -55,6 +59,10 @@ struct BucketMetadata { // Runtime-only (not serialized): true while a GC compaction is in flight // for this bucket, preventing re-entrant compaction. mutable std::atomic compacting_{false}; + // Runtime-only version of the bucket's deletion state. Compaction captures + // this value with its read snapshot and validates it before publishing a + // new bucket, so a concurrent deletion cannot publish stale data. + uint64_t generation_{0}; // Default constructor BucketMetadata() = default; @@ -65,10 +73,12 @@ struct BucketMetadata { data_size(other.data_size), keys(other.keys), metadatas(other.metadatas), + tombstones(other.tombstones), inflight_reads_(0), last_access_ns_(0), deleted_bytes_(0), - compacting_(false) {} + compacting_(false), + generation_(0) {} // Move constructor BucketMetadata(BucketMetadata&& other) noexcept @@ -76,10 +86,12 @@ struct BucketMetadata { data_size(other.data_size), keys(std::move(other.keys)), metadatas(std::move(other.metadatas)), + tombstones(std::move(other.tombstones)), inflight_reads_(0), last_access_ns_(0), deleted_bytes_(0), - compacting_(false) {} + compacting_(false), + generation_(0) {} // Copy assignment BucketMetadata& operator=(const BucketMetadata& other) { @@ -88,6 +100,7 @@ struct BucketMetadata { data_size = other.data_size; keys = other.keys; metadatas = other.metadatas; + tombstones = other.tombstones; // Don't copy runtime state } return *this; @@ -100,12 +113,13 @@ struct BucketMetadata { data_size = other.data_size; keys = std::move(other.keys); metadatas = std::move(other.metadatas); + tombstones = std::move(other.tombstones); // Don't move runtime state } return *this; } }; -YLT_REFL(BucketMetadata, data_size, keys, metadatas); +YLT_REFL(BucketMetadata, data_size, keys, metadatas, tombstones); /** * @brief RAII guard for tracking in-flight bucket reads. @@ -445,11 +459,16 @@ class StorageBackendInterface { // Default no-op: only BucketStorageBackend implements tombstone + GC. // File-per-key and other backends inherit the no-op (do not delete files). // Safe to call for keys not present in local storage (idempotent). - virtual void MarkRemoved(const std::string& /* key */) {} + virtual tl::expected MarkRemoved( + const std::string& /* key */) { + return {}; + } // Batch variant: mark multiple keys as removed in one lock acquisition. - virtual void BatchMarkRemoved( - const std::vector& /* keys */) {} + virtual tl::expected BatchMarkRemoved( + const std::vector& /* keys */) { + return {}; + } // Remove all persisted objects from disk. Called during RemoveAll to // clean up physical SSD files alongside master metadata deletion. virtual void RemoveAll() {} @@ -1047,8 +1066,10 @@ class BucketStorageBackend : public StorageBackendInterface { // Removes key from object_bucket_map_ (immediately invisible to // BatchLoad/IsExist) and bumps bucket deleted_bytes_. // Idempotent: no-op if key not in local storage. - void MarkRemoved(const std::string& key) override; - void BatchMarkRemoved(const std::vector& keys) override; + tl::expected MarkRemoved( + const std::string& key) override; + tl::expected BatchMarkRemoved( + const std::vector& keys) override; // Compact a single bucket: copy-on-write live keys to a new bucket, // atomically swap mappings, delete old bucket file after reads drain. @@ -1073,13 +1094,6 @@ class BucketStorageBackend : public StorageBackendInterface { private: // --- Background GC --- - // Select the best GC candidate bucket: deleted_bytes_>0, not compacting, - // highest deleted_ratio, then coldest last_access_ns_. - // Must be called with mutex_ held (exclusive). - // Returns buckets_.end() if no candidate. - std::map>::iterator - SelectGCCandidate(); - // Background GC thread entry point. void GCThreadFunc(); diff --git a/mooncake-store/src/client_service.cpp b/mooncake-store/src/client_service.cpp index 7e22f40354..9e9f611064 100644 --- a/mooncake-store/src/client_service.cpp +++ b/mooncake-store/src/client_service.cpp @@ -3510,6 +3510,11 @@ Client::RemoveObjectHeartbeat(const UUID& client_id) { return master_client_.RemoveObjectHeartbeat(client_id); } +tl::expected Client::AckRemoveObjectHeartbeat( + const UUID& client_id, const std::vector& tasks) { + return master_client_.AckRemoveObjectHeartbeat(client_id, tasks); +} + tl::expected Client::PromotionAllocStart( const std::string& key, uint64_t size, diff --git a/mooncake-store/src/file_storage.cpp b/mooncake-store/src/file_storage.cpp index 0520c213ba..c637b096c6 100644 --- a/mooncake-store/src/file_storage.cpp +++ b/mooncake-store/src/file_storage.cpp @@ -815,12 +815,14 @@ tl::expected FileStorage::IsEnableOffloading() { return enable_offloading; } -void FileStorage::MarkRemoved(const std::string& key) { - storage_backend_->MarkRemoved(key); +tl::expected FileStorage::MarkRemoved( + const std::string& key) { + return storage_backend_->MarkRemoved(key); } -void FileStorage::BatchMarkRemoved(const std::vector& keys) { - storage_backend_->BatchMarkRemoved(keys); +tl::expected FileStorage::BatchMarkRemoved( + const std::vector& keys) { + return storage_backend_->BatchMarkRemoved(keys); } tl::expected FileStorage::Heartbeat() { @@ -857,13 +859,26 @@ tl::expected FileStorage::Heartbeat() { auto remove_result = client_->RemoveObjectHeartbeat(client_->getClientId()); if (remove_result) { + bool all_marked = true; for (const auto& item : remove_result.value()) { auto storage_key = - MakeTenantScopedStorageKey(item.tenant_id, item.key); - storage_backend_->MarkRemoved(storage_key); + TenantId(item.tenant_id).MakeScopedKey(item.key); + auto mark_result = storage_backend_->MarkRemoved(storage_key); + if (!mark_result) { + all_marked = false; + LOG(ERROR) << "Failed to persist remove tombstone: " + << mark_result.error(); + break; + } } - if (!remove_result.value().empty()) { - VLOG(1) << "RemoveObjectHeartbeat drained " + if (all_marked && !remove_result.value().empty()) { + auto ack_result = client_->AckRemoveObjectHeartbeat( + client_->getClientId(), remove_result.value()); + if (!ack_result) { + LOG(ERROR) << "Failed to ACK remove tasks: " + << ack_result.error(); + } + VLOG(1) << "RemoveObjectHeartbeat processed " << remove_result.value().size() << " removed key(s) from master"; } diff --git a/mooncake-store/src/master_client.cpp b/mooncake-store/src/master_client.cpp index 0a06e73df9..5a266a696f 100644 --- a/mooncake-store/src/master_client.cpp +++ b/mooncake-store/src/master_client.cpp @@ -254,6 +254,11 @@ struct RpcNameTraits<&WrappedMasterService::RemoveObjectHeartbeat> { static constexpr const char* value = "RemoveObjectHeartbeat"; }; +template <> +struct RpcNameTraits<&WrappedMasterService::AckRemoveObjectHeartbeat> { + static constexpr const char* value = "AckRemoveObjectHeartbeat"; +}; + template <> struct RpcNameTraits<&WrappedMasterService::PromotionAllocStart> { static constexpr const char* value = "PromotionAllocStart"; @@ -1208,6 +1213,15 @@ MasterClient::RemoveObjectHeartbeat(const UUID& client_id) { std::vector>(client_id); } +tl::expected MasterClient::AckRemoveObjectHeartbeat( + const UUID& client_id, const std::vector& tasks) { + ScopedVLogTimer timer(1, "MasterClient::AckRemoveObjectHeartbeat"); + timer.LogRequest("client_id=", client_id.first, ":", client_id.second, + " tasks=", tasks.size()); + return invoke_rpc<&WrappedMasterService::AckRemoveObjectHeartbeat, void>( + client_id, tasks); +} + tl::expected MasterClient::PromotionAllocStart( const UUID& client_id, const std::string& key, uint64_t size, diff --git a/mooncake-store/src/master_service.cpp b/mooncake-store/src/master_service.cpp index dbce9168fa..2aef820aa5 100644 --- a/mooncake-store/src/master_service.cpp +++ b/mooncake-store/src/master_service.cpp @@ -1,6 +1,5 @@ #include "master_service.h" -#include #include #include #include @@ -1764,6 +1763,14 @@ void MasterService::FinalizeRemovedReplicasAfterDurable( const bool erased_local_disk = std::any_of( erased_replicas.begin(), erased_replicas.end(), [](const Replica& replica) { return replica.is_local_disk_replica(); }); + std::vector local_disk_holders; + for (const auto& replica : erased_replicas) { + if (!replica.is_local_disk_replica()) continue; + auto client_id = replica.get_local_disk_client_id(); + if (client_id.has_value()) { + local_disk_holders.push_back(client_id.value()); + } + } ReleaseLocalDiskUsage(erased_replicas); if (erased_local_disk) { shard.OnDiskReplicaRemoved(erased_local_disk, metadata); @@ -1774,6 +1781,11 @@ void MasterService::FinalizeRemovedReplicasAfterDurable( shard->tenants.erase(tenant_it); } } + if (erased_local_disk) { + EnqueueRemoveTasks( + local_disk_holders, + RemoveTaskItem{tenant_id.value(), durable_entry.object_key}); + } } void MasterService::FinalizeMetadataEraseAfterDurable( @@ -1810,6 +1822,7 @@ void MasterService::FinalizeExpiredProcessingReplicasAfterDurable( } auto& metadata = accessor.Get(); + auto replicas = PopReplicasWithCacheTotalAccounting( metadata, &Replica::fn_is_processing); if (!replicas.empty()) { @@ -5599,6 +5612,17 @@ auto MasterService::Remove(const std::string& key, const TenantId& tenant_id, } auto& metadata = accessor.Get(); + std::vector local_disk_holders; + metadata.VisitReplicas( + [](const Replica& replica) { + return replica.is_local_disk_replica(); + }, + [&local_disk_holders](Replica& replica) { + auto client_id = replica.get_local_disk_client_id(); + if (client_id.has_value()) { + local_disk_holders.push_back(client_id.value()); + } + }); if (!force && !metadata.IsLeaseExpired()) { VLOG(1) << "key=" << key << ", error=object_has_lease"; @@ -5636,10 +5660,15 @@ auto MasterService::Remove(const std::string& key, const TenantId& tenant_id, auto persist_result = AppendReservedOpLogWithDurableFinalize( std::move(reservation.value()), OpType::REMOVE, object_id.tenant_id.value(), key, {}, - [this, removed_ids = std::move(removed_ids)]( + [this, removed_ids = std::move(removed_ids), + local_disk_holders, + tenant_id_for_task = object_id.tenant_id.value(), key]( const OpLogEntry& durable_entry) { FinalizeRemovedReplicasAfterDurable( durable_entry, removed_ids, QuotaEraseMode::kFull); + EnqueueRemoveTasks( + local_disk_holders, + RemoveTaskItem{tenant_id_for_task, key}); }); if (!persist_result) { return tl::make_unexpected(persist_result.error()); @@ -5651,35 +5680,10 @@ auto MasterService::Remove(const std::string& key, const TenantId& tenant_id, // Before erasing metadata, collect LOCAL_DISK replica holders so we // can notify them to reclaim SSD space via RemoveObjectHeartbeat. - std::vector local_disk_holders; - metadata.VisitReplicas( - [](const Replica& replica) { - return replica.is_local_disk_replica(); - }, - [&local_disk_holders](Replica& replica) { - auto client_id = replica.get_local_disk_client_id(); - if (client_id.has_value()) { - local_disk_holders.push_back(client_id.value()); - } - }); - accessor.Erase(); // Push removed key to each LOCAL_DISK holder's removed_keys queue. - if (!local_disk_holders.empty()) { - ScopedLocalDiskSegmentAccess local_disk_segment_access = - segment_manager_.getLocalDiskSegmentAccess(); - auto& client_local_disk_segment = - local_disk_segment_access.getClientLocalDiskSegment(); - for (const auto& holder_id : local_disk_holders) { - auto it = client_local_disk_segment.find(holder_id); - if (it != client_local_disk_segment.end()) { - MutexLocker locker(&it->second->offloading_mutex_); - it->second->removed_keys.push_back( - RemoveTaskItem{tenant_id, key}); - } - } - } + EnqueueRemoveTasks(local_disk_holders, RemoveTaskItem{tenant_id.value(), key}); return {}; } @@ -6077,6 +6081,18 @@ auto MasterService::BatchRemove(const std::vector& keys, auto& metadata = it->second; + std::vector batch_local_disk_holders; + metadata.VisitReplicas( + [](const Replica& replica) { + return replica.is_local_disk_replica(); + }, + [&batch_local_disk_holders](Replica& replica) { + auto cid = replica.get_local_disk_client_id(); + if (cid.has_value()) { + batch_local_disk_holders.push_back(cid.value()); + } + }); + if (!force && !metadata.IsLeaseExpired(now)) { VLOG(1) << "key=" << key << ", error=object_has_lease"; results[original_idx] = @@ -6119,11 +6135,16 @@ auto MasterService::BatchRemove(const std::vector& keys, AppendReservedOpLogWithDurableFinalize( std::move(reservation.value()), OpType::REMOVE, normalized_tenant.value(), key, {}, - [this, removed_ids = std::move(removed_ids)]( + [this, removed_ids = std::move(removed_ids), + batch_local_disk_holders, + tenant_id = normalized_tenant.value(), key]( const OpLogEntry& durable_entry) { FinalizeRemovedReplicasAfterDurable( durable_entry, removed_ids, QuotaEraseMode::kFull); + EnqueueRemoveTasks( + batch_local_disk_holders, + RemoveTaskItem{tenant_id, key}); }); if (!persist_result) { results[original_idx] = @@ -6137,18 +6158,6 @@ auto MasterService::BatchRemove(const std::vector& keys, // Collect LOCAL_DISK replica holders before erasing, so we // can notify them to reclaim SSD space via RemoveObjectHeartbeat. - std::vector batch_local_disk_holders; - metadata.VisitReplicas( - [](const Replica& replica) { - return replica.is_local_disk_replica(); - }, - [&batch_local_disk_holders](Replica& replica) { - auto cid = replica.get_local_disk_client_id(); - if (cid.has_value()) { - batch_local_disk_holders.push_back(cid.value()); - } - }); - EraseMetadata(tenant_state, it, normalized_tenant, QuotaEraseMode::kFull, &shard); if (tenant_state.Empty()) { @@ -6156,21 +6165,9 @@ auto MasterService::BatchRemove(const std::vector& keys, } // Push removed key to each LOCAL_DISK holder's removed_keys queue. - if (!batch_local_disk_holders.empty()) { - ScopedLocalDiskSegmentAccess local_disk_segment_access = - segment_manager_.getLocalDiskSegmentAccess(); - auto& client_local_disk_segment = - local_disk_segment_access.getClientLocalDiskSegment(); - for (const auto& holder_id : batch_local_disk_holders) { - auto seg_it = client_local_disk_segment.find(holder_id); - if (seg_it != client_local_disk_segment.end()) { - MutexLocker locker( - &seg_it->second->offloading_mutex_); - seg_it->second->removed_keys.push_back( - RemoveTaskItem{normalized_tenant, key}); - } - } - } + EnqueueRemoveTasks( + batch_local_disk_holders, + RemoveTaskItem{normalized_tenant.value(), key}); results[original_idx] = {}; // Success } @@ -6407,12 +6404,51 @@ auto MasterService::RemoveObjectHeartbeat(const UUID& client_id) if (local_disk_segment_it == client_local_disk_segment.end()) { return tl::make_unexpected(ErrorCode::SEGMENT_NOT_FOUND); } - std::vector result; { MutexLocker locker(&local_disk_segment_it->second->offloading_mutex_); - result = std::move(local_disk_segment_it->second->removed_keys); + return local_disk_segment_it->second->removed_keys; } - return result; +} + +void MasterService::EnqueueRemoveTasks( + const std::vector& holder_ids, const RemoveTaskItem& task) { + if (holder_ids.empty()) return; + ScopedLocalDiskSegmentAccess access = + segment_manager_.getLocalDiskSegmentAccess(); + auto& segments = access.getClientLocalDiskSegment(); + for (const auto& holder_id : holder_ids) { + auto it = segments.find(holder_id); + if (it == segments.end()) continue; + MutexLocker locker(&it->second->offloading_mutex_); + if (std::find(it->second->removed_keys.begin(), + it->second->removed_keys.end(), task) == + it->second->removed_keys.end()) { + it->second->removed_keys.push_back(task); + } + } +} + +auto MasterService::AckRemoveObjectHeartbeat( + const UUID& client_id, const std::vector& tasks) + -> tl::expected { + std::shared_lock shared_lock(snapshot_mutex_); + ScopedLocalDiskSegmentAccess access = + segment_manager_.getLocalDiskSegmentAccess(); + auto& segments = access.getClientLocalDiskSegment(); + auto it = segments.find(client_id); + if (it == segments.end()) { + return tl::make_unexpected(ErrorCode::SEGMENT_NOT_FOUND); + } + MutexLocker locker(&it->second->offloading_mutex_); + auto& pending = it->second->removed_keys; + pending.erase(std::remove_if(pending.begin(), pending.end(), + [&tasks](const RemoveTaskItem& task) { + return std::find(tasks.begin(), + tasks.end(), task) != + tasks.end(); + }), + pending.end()); + return {}; } auto MasterService::ReportSsdCapacity(const UUID& client_id, diff --git a/mooncake-store/src/rpc_service.cpp b/mooncake-store/src/rpc_service.cpp index 206934855e..57abd06492 100644 --- a/mooncake-store/src/rpc_service.cpp +++ b/mooncake-store/src/rpc_service.cpp @@ -1735,6 +1735,14 @@ WrappedMasterService::RemoveObjectHeartbeat(const UUID& client_id) { return master_service_.RemoveObjectHeartbeat(client_id); } +tl::expected +WrappedMasterService::AckRemoveObjectHeartbeat( + const UUID& client_id, const std::vector& tasks) { + ScopedVLogTimer timer(1, "AckRemoveObjectHeartbeat"); + timer.LogRequest("action=ack_remove_heartbeat"); + return master_service_.AckRemoveObjectHeartbeat(client_id, tasks); +} + tl::expected WrappedMasterService::PromotionAllocStart( const UUID& client_id, const std::string& key, const std::string& tenant_id, @@ -1932,6 +1940,9 @@ void RegisterRpcService( server.register_handler< &mooncake::WrappedMasterService::RemoveObjectHeartbeat>( &wrapped_master_service); + server.register_handler< + &mooncake::WrappedMasterService::AckRemoveObjectHeartbeat>( + &wrapped_master_service); server .register_handler<&mooncake::WrappedMasterService::PromotionAllocStart>( &wrapped_master_service); diff --git a/mooncake-store/src/storage_backend.cpp b/mooncake-store/src/storage_backend.cpp index 7dc0b8b646..ab2723b897 100644 --- a/mooncake-store/src/storage_backend.cpp +++ b/mooncake-store/src/storage_backend.cpp @@ -2139,15 +2139,12 @@ tl::expected BucketStorageBackend::BatchLoad( #endif { // Fallback to per-key vector_read for non-UringFile (PosixFile). - for (const auto& plan : read_plans) { - int64_t actual_offset = plan.offset + plan.key_size; - iovec iov{plan.dest_slice.ptr, plan.dest_slice.size}; - SpDiag::PerfPoint pt_posix(PerfKey::GET_SSD_OWNER_LOAD_POSIX, - SpDiag::PerfLevel::MODULE); - pt_posix.Start(); - read_res = file->vector_read(&iov, 1, actual_offset); - pt_posix.End(read_res ? 0 : -1); - } + iovec iov{plan.dest_slice.ptr, plan.dest_slice.size}; + SpDiag::PerfPoint pt_posix(PerfKey::GET_SSD_OWNER_LOAD_POSIX, + SpDiag::PerfLevel::MODULE); + pt_posix.Start(); + read_res = file->vector_read(&iov, 1, actual_offset); + pt_posix.End(read_res ? 0 : -1); if (stats) { const auto read_us = std::chrono::duration_cast( @@ -2189,6 +2186,7 @@ tl::expected BucketStorageBackend::BatchLoad( return tl::make_unexpected(ErrorCode::FILE_READ_FAIL); } } + } } // bucket_guards go out of scope here, decrementing inflight_reads_ @@ -2336,13 +2334,23 @@ tl::expected BucketStorageBackend::Init() { total_size_ += metadata_it->second->data_size + metadata_it->second->meta_size; for (size_t i = 0; i < metadata_it->second->keys.size(); i++) { + const auto& key = metadata_it->second->keys[i]; + if (std::find(metadata_it->second->tombstones.begin(), + metadata_it->second->tombstones.end(), key) != + metadata_it->second->tombstones.end()) { + metadata_it->second->deleted_bytes_.fetch_add( + metadata_it->second->metadatas[i].key_size + + metadata_it->second->metadatas[i].data_size, + std::memory_order_relaxed); + continue; + } object_bucket_map_.emplace( - metadata_it->second->keys[i], - StorageObjectMetadata{ - metadata_it->first, - metadata_it->second->metadatas[i].offset, - metadata_it->second->metadatas[i].key_size, - metadata_it->second->metadatas[i].data_size, ""}); + key, StorageObjectMetadata{ + metadata_it->first, + metadata_it->second->metadatas[i].offset, + metadata_it->second->metadatas[i].key_size, + metadata_it->second->metadatas[i].data_size, + ""}); } } } @@ -3483,72 +3491,43 @@ void BucketStorageBackend::RemoveAll() { } // --- Explicit-delete-only GC --- -void BucketStorageBackend::MarkRemoved(const std::string& key) { +tl::expected BucketStorageBackend::MarkRemoved( + const std::string& key) { SharedMutexLocker lock(&mutex_); auto it = object_bucket_map_.find(key); if (it == object_bucket_map_.end()) { - return; + return {}; } int64_t bucket_id = it->second.bucket_id; int64_t freed = it->second.data_size + it->second.key_size; - object_bucket_map_.erase(it); - auto bucket_it = buckets_.find(bucket_id); - if (bucket_it != buckets_.end()) { - bucket_it->second->deleted_bytes_.fetch_add( - freed, std::memory_order_relaxed); + if (bucket_it == buckets_.end()) { + return tl::make_unexpected(ErrorCode::BUCKET_NOT_FOUND); } + auto bucket = bucket_it->second; + const auto object_metadata = it->second; + object_bucket_map_.erase(it); + bucket->tombstones.push_back(key); + auto persist_result = StoreBucketMetadata(bucket_id, bucket); + if (!persist_result) { + object_bucket_map_.emplace(key, object_metadata); + bucket->tombstones.pop_back(); + return tl::make_unexpected(persist_result.error()); + } + bucket->deleted_bytes_.fetch_add(freed, std::memory_order_relaxed); + ++bucket->generation_; + return {}; } -void BucketStorageBackend::BatchMarkRemoved( +tl::expected BucketStorageBackend::BatchMarkRemoved( const std::vector& keys) { - SharedMutexLocker lock(&mutex_); for (const auto& key : keys) { - auto it = object_bucket_map_.find(key); - if (it == object_bucket_map_.end()) continue; - int64_t bucket_id = it->second.bucket_id; - int64_t freed = it->second.data_size + it->second.key_size; - object_bucket_map_.erase(it); - - auto bucket_it = buckets_.find(bucket_id); - if (bucket_it != buckets_.end()) { - bucket_it->second->deleted_bytes_.fetch_add( - freed, std::memory_order_relaxed); + auto result = MarkRemoved(key); + if (!result) { + return result; } } -} - -std::map>::iterator -BucketStorageBackend::SelectGCCandidate() { - // Must be called with mutex_ held (exclusive). - // Find bucket with highest deleted_ratio; tie-break by coldest - // last_access_ns_ (smallest). Skip buckets with deleted_bytes_==0 - // or compacting_==true. - auto best_it = buckets_.end(); - double best_ratio = 0.0; - int64_t best_ts = std::numeric_limits::max(); - - for (auto it = buckets_.begin(); it != buckets_.end(); ++it) { - const auto& bucket = it->second; - int64_t deleted = - bucket->deleted_bytes_.load(std::memory_order_relaxed); - if (deleted <= 0) continue; - if (bucket->compacting_.load(std::memory_order_relaxed)) continue; - - int64_t data_size = bucket->data_size; - if (data_size <= 0) continue; - double ratio = - static_cast(deleted) / static_cast(data_size); - - int64_t ts = bucket->last_access_ns_.load(std::memory_order_relaxed); - - if (ratio > best_ratio || (ratio == best_ratio && ts < best_ts)) { - best_ratio = ratio; - best_ts = ts; - best_it = it; - } - } - return best_it; + return {}; } bool BucketStorageBackend::CompactBucket(int64_t bucket_id) { @@ -3570,6 +3549,7 @@ bool BucketStorageBackend::CompactBuckets( std::unordered_map, int64_t>> old_buckets; + std::unordered_map old_bucket_generations; { SharedMutexLocker lock(&mutex_); for (int64_t bid : bucket_ids) { @@ -3581,6 +3561,7 @@ bool BucketStorageBackend::CompactBuckets( int64_t ts = bucket->last_access_ns_.load( std::memory_order_relaxed); old_buckets[bid] = {bucket, ts}; + old_bucket_generations[bid] = bucket->generation_; for (size_t i = 0; i < bucket->keys.size(); ++i) { const auto& key = bucket->keys[i]; @@ -3799,6 +3780,23 @@ bool BucketStorageBackend::CompactBuckets( std::vector buckets_to_delete; { SharedMutexLocker lock(&mutex_); + for (const auto& [bid, pr] : old_buckets) { + auto current_it = buckets_.find(bid); + auto generation_it = old_bucket_generations.find(bid); + if (current_it == buckets_.end() || + generation_it == old_bucket_generations.end() || + current_it->second->generation_ != generation_it->second) { + // A delete changed the source bucket after the snapshot. The + // data file was built from a stale view, so do not publish it + // or replace any object mappings with it. + lock.unlock(); + CleanupOrphanedBucket(new_bucket_id); + reset_compacting(); + LOG(INFO) << "CompactBuckets discarded stale snapshot for " + << "bucket_id=" << bid; + return true; + } + } // Re-validate and remap each key in the new bucket. for (size_t i = 0; i < new_bucket->keys.size(); ++i) { const auto& key = new_bucket->keys[i]; diff --git a/mooncake-store/tests/e2e/gc_e2e_test.cpp b/mooncake-store/tests/e2e/gc_e2e_test.cpp index 2e4e686274..b39f03269b 100644 --- a/mooncake-store/tests/e2e/gc_e2e_test.cpp +++ b/mooncake-store/tests/e2e/gc_e2e_test.cpp @@ -26,7 +26,7 @@ #include #include -#include "client_buffer.hpp" +#include "client_buffer.h" #include "real_client.h" #include "test_server_helpers.h" #include "types.h" diff --git a/mooncake-store/tests/e2e/run_gc_e2e.sh b/mooncake-store/tests/e2e/run_gc_e2e.sh deleted file mode 100644 index 30a72ec98d..0000000000 --- a/mooncake-store/tests/e2e/run_gc_e2e.sh +++ /dev/null @@ -1,266 +0,0 @@ -#!/usr/bin/env bash -# run_gc_e2e.sh -# -# Multi-process integration test for the explicit-delete-only SSD GC. -# -# This is "Plan B": launches the REAL production binaries -# - mooncake_master (master service) -# - mooncake_client (real_client_main, the RPC server that owns the -# BucketStorageBackend + FileStorage offload path) -# and drives put/get/remove via the Python MooncakeDistributedStore client, -# which connects to the real_client via RPC and exercises -# RealClient::remove_internal -> FileStorage::MarkRemoved -> GC. -# -# Unlike gc_e2e_test.cpp (in-process), this uses separate OS processes, -# matching the production deployment topology. -# -# Prerequisites: -# - BUILD_DIR points to a build tree containing: -# mooncake-store/src/mooncake_master -# mooncake-store/src/mooncake_client -# mooncake-integration/store*.so (Python bindings) -# - Python environment with the mooncake wheel installed -# - A standalone HTTP metadata server (the master's embedded one or -# mooncake-wheel/mooncake/http_metadata_server.py) -# -# Usage: -# cd mooncake-store/tests/e2e -# BUILD_DIR=/path/to/build ./run_gc_e2e.sh -# -# Environment variables (all optional): -# BUILD_DIR Build tree root (default: ../../../build) -# WORK_DIR Temp working dir (default: /tmp/mooncake_gc_e2e) -# SSD_OFFLOAD_PATH SSD offload dir (default: $WORK_DIR/ssd_offload) -# MASTER_PORT Master RPC port (default: 50051) -# CLIENT_PORT RealClient RPC port (default: 50052) -# METADATA_PORT HTTP metadata port (default: 8080) -# GC_INTERVAL_MS GC scan interval (default: 200) -# GC_DELETED_RATIO Compaction threshold (default: 0.1) -# PAYLOAD_SIZE Key value size (default: 4194304 = 4MB) -# PASS_CODE Expected exit code for pass (default: 0) - -set -euo pipefail - -BUILD_DIR="${BUILD_DIR:-$(cd "$(dirname "$0")/../../.." && pwd)/build}" -WORK_DIR="${WORK_DIR:-/tmp/mooncake_gc_e2e}" -SSD_OFFLOAD_PATH="${SSD_OFFLOAD_PATH:-$WORK_DIR/ssd_offload}" -MASTER_PORT="${MASTER_PORT:-50051}" -CLIENT_PORT="${CLIENT_PORT:-50052}" -METADATA_PORT="${METADATA_PORT:-8080}" -GC_INTERVAL_MS="${GC_INTERVAL_MS:-200}" -GC_DELETED_RATIO="${GC_DELETED_RATIO:-0.1}" -PAYLOAD_SIZE="${PAYLOAD_SIZE:-4194304}" - -MASTER_BIN="$BUILD_DIR/mooncake-store/src/mooncake_master" -CLIENT_BIN="$BUILD_DIR/mooncake-store/src/mooncake_client" -METADATA_BIN="mooncake.http_metadata_server" # python module - -MASTER_ADDR="127.0.0.1:$MASTER_PORT" -METADATA_URL="http://127.0.0.1:$METADATA_PORT/metadata" -LOG_DIR="$WORK_DIR/logs" - -mkdir -p "$WORK_DIR" "$SSD_OFFLOAD_PATH" "$LOG_DIR" - -cleanup() { - local exit_code=$? - echo "[cleanup] stopping processes..." - [[ -n "${CLIENT_PID:-}" ]] && kill "$CLIENT_PID" 2>/dev/null || true - [[ -n "${MASTER_PID:-}" ]] && kill "$MASTER_PID" 2>/dev/null || true - [[ -n "${META_PID:-}" ]] && kill "$META_PID" 2>/dev/null || true - wait 2>/dev/null || true - if [[ $exit_code -eq 0 ]]; then - echo "[result] PASS" - else - echo "[result] FAIL (exit=$exit_code)" - fi - exit $exit_code -} -trap cleanup EXIT - -echo "============================================================" -echo " SSD GC Multi-Process E2E Test" -echo "============================================================" -echo " BUILD_DIR = $BUILD_DIR" -echo " WORK_DIR = $WORK_DIR" -echo " SSD_OFFLOAD_PATH = $SSD_OFFLOAD_PATH" -echo " GC_INTERVAL_MS = $GC_INTERVAL_MS" -echo " GC_DELETED_RATIO = $GC_DELETED_RATIO" -echo " PAYLOAD_SIZE = $PAYLOAD_SIZE" -echo "============================================================" - -# --- 1. Start HTTP metadata server (standalone, for transfer engine) --- -echo "[1/5] Starting HTTP metadata server on port $METADATA_PORT..." -python3 -m mooncake.http_metadata_server --port "$METADATA_PORT" \ - >"$LOG_DIR/metadata.log" 2>&1 & -META_PID=$! -sleep 1 -if ! kill -0 "$META_PID" 2>/dev/null; then - echo "[ERROR] metadata server failed to start (see $LOG_DIR/metadata.log)" - exit 1 -fi - -# --- 2. Start mooncake_master with offload enabled --- -echo "[2/5] Starting mooncake_master on port $MASTER_PORT..." -"$MASTER_BIN" \ - --port "$MASTER_PORT" \ - --metrics_port 9004 \ - --enable_offload \ - --default_kv_lease_ttl 0 \ - --log_dir "$LOG_DIR" \ - >"$LOG_DIR/master.log" 2>&1 & -MASTER_PID=$! -sleep 2 -if ! kill -0 "$MASTER_PID" 2>/dev/null; then - echo "[ERROR] master failed to start (see $LOG_DIR/master.log)" - exit 1 -fi - -# --- 3. Start mooncake_client (real_client_main) with SSD offload --- -# Export GC env vars BEFORE launching the client so BucketBackendConfig -# picks them up via FromEnvironment(). -echo "[3/5] Starting mooncake_client (real_client) on port $CLIENT_PORT..." -export MOONCAKE_OFFLOAD_STORAGE_BACKEND_DESCRIPTOR="bucket_storage_backend" -export MOONCAKE_OFFLOAD_BUCKET_EVICTION_POLICY="lru" -export MOONCAKE_OFFLOAD_DISABLE_SSD_EVICTION="true" -export MOONCAKE_OFFLOAD_BUCKET_GC_INTERVAL_MS="$GC_INTERVAL_MS" -export MOONCAKE_OFFLOAD_BUCKET_GC_DELETED_RATIO="$GC_DELETED_RATIO" -export MOONCAKE_OFFLOAD_FILE_STORAGE_PATH="$SSD_OFFLOAD_PATH" - -"$CLIENT_BIN" \ - --host "127.0.0.1" \ - --metadata_server "$METADATA_URL" \ - --master_server_address "$MASTER_ADDR" \ - --protocol tcp \ - --port "$CLIENT_PORT" \ - --global_segment_size "512MB" \ - --enable_offload \ - --start_offload_rpc_server \ - >"$LOG_DIR/client.log" 2>&1 & -CLIENT_PID=$! -sleep 3 -if ! kill -0 "$CLIENT_PID" 2>/dev/null; then - echo "[ERROR] real_client failed to start (see $LOG_DIR/client.log)" - exit 1 -fi - -CLIENT_RPC_ADDR="127.0.0.1:$CLIENT_PORT" -echo "[info] real_client RPC at $CLIENT_RPC_ADDR" - -# --- 4. Drive workload via Python client --- -# The Python MooncakeDistributedStore connects to the real_client RPC -# server and issues put/get/remove through RealClient handlers, which is -# the only path that triggers MarkRemoved -> GC. -echo "[4/5] Running GC workload via Python client..." - -export PYTHONPATH="$BUILD_DIR/mooncake-integration:${PYTHONPATH:-}" - -python3 - "$CLIENT_RPC_ADDR" "$MASTER_ADDR" "$METADATA_URL" \ - "$SSD_OFFLOAD_PATH" "$PAYLOAD_SIZE" "$LOG_DIR" <<'PYEOF' -import os, sys, time, glob - -client_rpc = sys.argv[1] -master_addr = sys.argv[2] -metadata_url = sys.argv[3] -ssd_path = sys.argv[4] -payload_size = int(sys.argv[5]) -log_dir = sys.argv[6] - -import mooncake.store as mc - -# Connect a Python client to the real_client RPC server. -# MooncakeDistributedStore routes put/get/remove to the real_client -# process, exercising RealClient::remove_internal -> MarkRemoved. -client = mc.MooncakeDistributedStore( - local_hostname="127.0.0.1:50070", - metadata_server=metadata_url, - master_server_addr=master_addr, - global_segment_size=256 * 1024 * 1024, - local_buffer_size=128 * 1024 * 1024, - protocol="tcp", -) - -def count_bucket_files(): - return len(glob.glob(os.path.join(ssd_path, "*.bucket"))) - -v1 = b"A" * payload_size -v2 = b"B" * payload_size -v3 = b"C" * payload_size - -# Put 3 keys -> offloaded to bucket(s) on SSD. -print("[py] putting 3 keys...") -assert client.put("gc_b_k1", v1) == 0, "put k1 failed" -assert client.put("gc_b_k2", v2) == 0, "put k2 failed" -assert client.put("gc_b_k3", v3) == 0, "put k3 failed" - -# Wait for offload to complete (master queues on PutEnd, heartbeat drains). -print("[py] waiting for offload...") -for _ in range(50): - try: - got = client.get("gc_b_k1") - if got == v1: - break - except Exception: - pass - time.sleep(0.2) -else: - print("[py][ERROR] offload timed out", file=sys.stderr) - sys.exit(1) - -buckets_before = count_bucket_files() -print(f"[py] buckets_before={buckets_before}") - -# Remove the middle key -> tombstone, triggers GC compaction. -print("[py] removing gc_b_k2...") -assert client.remove("gc_b_k2") == 0, "remove k2 failed" - -# Wait for GC compaction: survivors must stay readable, removed key gone. -print("[py] waiting for GC compaction...") -deadline = time.time() + 30 -reclaimed = False -while time.time() < deadline: - # Survivors must remain readable with correct data. - try: - assert client.get("gc_b_k1") == v1, "k1 corrupted during GC" - assert client.get("gc_b_k3") == v3, "k3 corrupted during GC" - except AssertionError as e: - print(f"[py][ERROR] {e}", file=sys.stderr) - sys.exit(1) - # Removed key must stay gone. - try: - client.get("gc_b_k2") - print("[py][ERROR] removed key k2 reappeared", file=sys.stderr) - sys.exit(1) - except Exception: - pass # expected: get fails - # Detect compaction: bucket file count changed. - buckets_now = count_bucket_files() - if buckets_now > 0 and buckets_now != buckets_before: - reclaimed = True - print(f"[py] compaction detected: buckets {buckets_before}->{buckets_now}") - break - time.sleep(0.3) - -# Final integrity check. -assert client.get("gc_b_k1") == v1, "k1 final check failed" -assert client.get("gc_b_k3") == v3, "k3 final check failed" - -if not reclaimed: - print("[py][WARN] GC reclamation not detected within timeout " - "(may still pass if data integrity holds)") -print("[py] PASS: survivors intact, removed key gone") -PYEOF - -PY_EXIT=$? -if [[ $PY_EXIT -ne 0 ]]; then - echo "[ERROR] Python workload failed (exit=$PY_EXIT, see $LOG_DIR/client.log)" - exit $PY_EXIT -fi - -echo "[5/5] Verifying SSD file reclamation..." -FINAL_BUCKETS=$(find "$SSD_OFFLOAD_PATH" -name "*.bucket" 2>/dev/null | wc -l) -echo "[info] final bucket files: $FINAL_BUCKETS" - -echo "============================================================" -echo " GC E2E: PASS" -echo "============================================================" -exit 0