From 95603969ec58e794808cee1d735935ce8c11dd23 Mon Sep 17 00:00:00 2001 From: Jie Yao Date: Wed, 30 Sep 2026 15:35:33 +0800 Subject: [PATCH 1/2] Sync scrub superblk timestamps between leader and followers via RPC Followers previously had no visibility into a PG's last_deep_scrub_timestamp/ last_shallow_scrub_timestamp, which the leader persists locally after each scrub. Add a PUSH_SCRUB_SUPERBLK data RPC channel and a periodic leader-side timer that pushes these timestamps to followers; each side keeps the newer value per field and the response is merged back into the leader. --- conanfile.py | 2 +- src/lib/homestore_backend/gc_manager.cpp | 4 +- src/lib/homestore_backend/gc_manager.hpp | 44 ++++------ src/lib/homestore_backend/hs_blob_manager.cpp | 4 +- src/lib/homestore_backend/hs_homeobject.hpp | 4 + src/lib/homestore_backend/hs_pg_manager.cpp | 39 +++++++++ .../replication_state_machine.hpp | 2 +- src/lib/homestore_backend/scrub_manager.cpp | 87 +++++++++++++++++++ src/lib/homestore_backend/scrub_manager.hpp | 22 +++++ .../snapshot_receive_handler.cpp | 18 ++-- .../tests/homeobj_misc_tests.cpp | 20 ++--- .../tests/hs_scrubber_tests.cpp | 84 ++++++++++++++++++ 12 files changed, 277 insertions(+), 53 deletions(-) diff --git a/conanfile.py b/conanfile.py index f032c15c4..86a2cbcf7 100644 --- a/conanfile.py +++ b/conanfile.py @@ -10,7 +10,7 @@ class HomeObjectConan(ConanFile): name = "homeobject" - version = "4.3.5" + version = "4.3.6" homepage = "https://github.com/eBay/HomeObject" description = "Blob Store built on HomeStore" diff --git a/src/lib/homestore_backend/gc_manager.cpp b/src/lib/homestore_backend/gc_manager.cpp index 82f5e4791..cf14075d2 100644 --- a/src/lib/homestore_backend/gc_manager.cpp +++ b/src/lib/homestore_backend/gc_manager.cpp @@ -13,8 +13,8 @@ SISL_LOGGING_DECL(gcmgr) #define RECOVERD_GC_TASK_ID 0 #define GCLOG(level, gc_task_id, pg_id, shard_id, msg, ...) \ - LOG##level##MOD(gcmgr, "[gc_task_id={}, pg_id={}, shard_id=0x{:x}] " msg, gc_task_id, pg_id, shard_id, \ - ##__VA_ARGS__) + LOG##level##MOD(gcmgr, "[gc_task_id={}, pg_id={}, shard_id=0x{:x}] " msg, gc_task_id, pg_id, \ + shard_id, ##__VA_ARGS__) #define GCLOGT(gc_task_id, pg_id, shard_id, msg, ...) GCLOG(TRACE, gc_task_id, pg_id, shard_id, msg, ##__VA_ARGS__) #define GCLOGD(gc_task_id, pg_id, shard_id, msg, ...) GCLOG(DEBUG, gc_task_id, pg_id, shard_id, msg, ##__VA_ARGS__) diff --git a/src/lib/homestore_backend/gc_manager.hpp b/src/lib/homestore_backend/gc_manager.hpp index 7e03153bd..39b7608ac 100644 --- a/src/lib/homestore_backend/gc_manager.hpp +++ b/src/lib/homestore_backend/gc_manager.hpp @@ -155,14 +155,10 @@ class GCManager { // Backlog / pressure snapshot gauges. Values refreshed once per scan cycle by // GCManager::scan_chunks_for_gc; worst-case staleness = gc_scan_interval_sec. - REGISTER_GAUGE(pending_gc_bytes, - "Total reclaimable garbage bytes in PG-owned chunks on this pdev"); - REGISTER_GAUGE(eligible_gc_bytes, - "Reclaimable bytes currently eligible for normal GC on this pdev"); - REGISTER_GAUGE(eligible_gc_chunk_count, - "Chunks currently eligible for normal GC on this pdev"); - REGISTER_GAUGE(pending_normal_gc_task_count, - "Normal-priority GC tasks queued or running on this pdev"); + REGISTER_GAUGE(pending_gc_bytes, "Total reclaimable garbage bytes in PG-owned chunks on this pdev"); + REGISTER_GAUGE(eligible_gc_bytes, "Reclaimable bytes currently eligible for normal GC on this pdev"); + REGISTER_GAUGE(eligible_gc_chunk_count, "Chunks currently eligible for normal GC on this pdev"); + REGISTER_GAUGE(pending_normal_gc_task_count, "Normal-priority GC tasks queued or running on this pdev"); // Distribution of the pending backlog by garbage-ratio bucket. We register 10 // gauges under a single Prometheus metric name (`pending_gc_chunks_ratio`), @@ -179,18 +175,16 @@ class GCManager { // shape. Prometheus text output is unaffected — it uses HELP (unchanged across // registrations) and disambiguates series by labels. static constexpr std::array< const char*, 10 > kRatioBucketLabels = { - "00-10", "10-20", "20-30", "30-40", "40-50", - "50-60", "60-70", "70-80", "80-90", "90-100"}; + "00-10", "10-20", "20-30", "30-40", "40-50", "50-60", "60-70", "70-80", "80-90", "90-100"}; for (size_t i = 0; i < kRatioBucketLabels.size(); ++i) { const auto lo = i * 10; const auto hi = (i + 1) * 10; - const auto desc = fmt::format( - "Snapshot count of pending chunks with garbage ratio in ({}, {}]% " - "(bucket={})", - lo, hi, kRatioBucketLabels[i]); - ratio_bucket_indices_[i] = m_impl_ptr->register_gauge( - "pending_gc_chunks_ratio", desc, "" /* report_name */, - sisl::metric_label{"bucket", kRatioBucketLabels[i]}); + const auto desc = fmt::format("Snapshot count of pending chunks with garbage ratio in ({}, {}]% " + "(bucket={})", + lo, hi, kRatioBucketLabels[i]); + ratio_bucket_indices_[i] = + m_impl_ptr->register_gauge("pending_gc_chunks_ratio", desc, "" /* report_name */, + sisl::metric_label{"bucket", kRatioBucketLabels[i]}); } register_me_to_farm(); @@ -230,9 +224,8 @@ class GCManager { // Bypass GAUGE_UPDATE for the same reason as bucket registration: we need to // address 10 distinct gauge indices that share one metric name. for (size_t i = 0; i < ratio_bucket_indices_.size(); ++i) { - m_impl_ptr->gauge_update( - ratio_bucket_indices_[i], - static_cast< int64_t >(gc_actor_.get_pending_ratio_bucket(i))); + m_impl_ptr->gauge_update(ratio_bucket_indices_[i], + static_cast< int64_t >(gc_actor_.get_pending_ratio_bucket(i))); } } @@ -315,12 +308,8 @@ class GCManager { // Snapshot readers used by pdev_gc_metrics::on_gather. Return the last value published // by GCManager::scan_chunks_for_gc for this pdev; 0 before the first scan completes. - uint64_t get_pending_gc_bytes() const { - return m_pending_gc_bytes.load(std::memory_order_relaxed); - } - uint64_t get_eligible_gc_bytes() const { - return m_eligible_gc_bytes.load(std::memory_order_relaxed); - } + uint64_t get_pending_gc_bytes() const { return m_pending_gc_bytes.load(std::memory_order_relaxed); } + uint64_t get_eligible_gc_bytes() const { return m_eligible_gc_bytes.load(std::memory_order_relaxed); } uint32_t get_eligible_gc_chunk_count() const { return m_eligible_gc_chunk_count.load(std::memory_order_relaxed); } @@ -332,8 +321,7 @@ class GCManager { // the totals locally over all chunks on this pdev, then hand them in via one call so the // metrics stay internally consistent within a scan cycle. Between-gauge drift is bounded // by one scan interval; individual scalars are aligned and therefore torn-read safe. - void publish_scan_snapshot(uint64_t pending_bytes, uint64_t eligible_bytes, - uint32_t eligible_chunks, + void publish_scan_snapshot(uint64_t pending_bytes, uint64_t eligible_bytes, uint32_t eligible_chunks, const std::array< uint32_t, 10 >& ratio_buckets) { m_pending_gc_bytes.store(pending_bytes, std::memory_order_relaxed); m_eligible_gc_bytes.store(eligible_bytes, std::memory_order_relaxed); diff --git a/src/lib/homestore_backend/hs_blob_manager.cpp b/src/lib/homestore_backend/hs_blob_manager.cpp index 6dc82be30..e1ec6846a 100644 --- a/src/lib/homestore_backend/hs_blob_manager.cpp +++ b/src/lib/homestore_backend/hs_blob_manager.cpp @@ -14,8 +14,8 @@ SISL_LOGGING_DECL(blobmgr) #define BLOG(level, trace_id, shard_id, blob_id, msg, ...) \ LOG##level##MOD(blobmgr, "[traceID={},shardID=0x{:x},pg={},shard=0x{:x},blob={}] " msg, trace_id, shard_id, \ - (shard_id >> homeobject::shard_width), (shard_id & homeobject::shard_mask), blob_id, \ - ##__VA_ARGS__) + (shard_id >> homeobject::shard_width), (shard_id & homeobject::shard_mask), \ + blob_id, ##__VA_ARGS__) #define BLOGT(trace_id, shard_id, blob_id, msg, ...) BLOG(TRACE, trace_id, shard_id, blob_id, msg, ##__VA_ARGS__) #define BLOGD(trace_id, shard_id, blob_id, msg, ...) BLOG(DEBUG, trace_id, shard_id, blob_id, msg, ##__VA_ARGS__) diff --git a/src/lib/homestore_backend/hs_homeobject.hpp b/src/lib/homestore_backend/hs_homeobject.hpp index 9a34448bb..1a7fbbf5e 100644 --- a/src/lib/homestore_backend/hs_homeobject.hpp +++ b/src/lib/homestore_backend/hs_homeobject.hpp @@ -439,9 +439,11 @@ class HSHomeObject : public HomeObjectImpl { * RPC handlers for scrub: * 1. on_scrub_req_received: receive the scrub req from leader * 2. on_scrub_result_received: receive the scrub result from followers + * 3. on_scrub_timestamp_push_received: receive the pushed scrub superblk timestamps from the leader */ void on_scrub_req_received(boost::intrusive_ptr< sisl::GenericRpcData >& rpc_data); void on_scrub_result_received(boost::intrusive_ptr< sisl::GenericRpcData >& rpc_data); + void on_scrub_timestamp_push_received(boost::intrusive_ptr< sisl::GenericRpcData >& rpc_data); /** * Register data RPC handlers for this PG @@ -572,6 +574,8 @@ class HSHomeObject : public HomeObjectImpl { inline const static std::string PUSH_SCRUB_REQ{"PUSH_SCRUB_REQ"}; // return scrub result to leader inline const static std::string PUSH_SCRUB_RESULT{"PUSH_SCRUB_RESULT"}; + // sync last_deep_scrub_timestamp/last_shallow_scrub_timestamp between leader and followers + inline const static std::string PUSH_SCRUB_TIMESTAMP{"PUSH_SCRUB_TIMESTAMP"}; class PGBlobIterator { public: diff --git a/src/lib/homestore_backend/hs_pg_manager.cpp b/src/lib/homestore_backend/hs_pg_manager.cpp index 7ad853e2c..7d140c027 100644 --- a/src/lib/homestore_backend/hs_pg_manager.cpp +++ b/src/lib/homestore_backend/hs_pg_manager.cpp @@ -1182,6 +1182,14 @@ void HSHomeObject::HS_PG::register_data_rpc_handlers() { } else { LOGW("PUSH_SCRUB_RESULT RPC handler already registered for pg={}", pg_id); } + + success = + repl_dev_->add_data_rpc_service(PUSH_SCRUB_TIMESTAMP, bind_this(HS_PG::on_scrub_timestamp_push_received, 1)); + if (success) { + LOGI("Successfully registered PUSH_SCRUB_TIMESTAMP RPC handler for pg={}", pg_id); + } else { + LOGW("PUSH_SCRUB_TIMESTAMP RPC handler already registered for pg={}", pg_id); + } } void HSHomeObject::HS_PG::on_scrub_req_received(boost::intrusive_ptr< sisl::GenericRpcData >& rpc_data) { @@ -1329,6 +1337,37 @@ void HSHomeObject::HS_PG::on_scrub_result_received(boost::intrusive_ptr< sisl::G scrub_mgr->handle_scrub_req_resp(pg_id, scrub_result); } +void HSHomeObject::HS_PG::on_scrub_timestamp_push_received(boost::intrusive_ptr< sisl::GenericRpcData >& rpc_data) { + const auto pg_id = pg_info_.id; + LOGD("Received scrub superblk push for pg={}", pg_id); + + auto const& incoming_buf = rpc_data->request_blob(); + if (incoming_buf.size() != sizeof(ScrubManager::scrub_timestamp_info)) { + LOGW("scrub superblk push received with invalid buffer size for pg={}, size={}", pg_id, incoming_buf.size()); + rpc_data->send_response(); + return; + } + + ScrubManager::scrub_timestamp_info incoming; + std::memcpy(&incoming, incoming_buf.cbytes(), sizeof(incoming)); + + auto scrub_mgr = home_obj_.scrub_manager(); + if (!scrub_mgr) { + LOGW("ScrubManager is not initialized in HS_PG::on_scrub_timestamp_push_received for pg={}", pg_id); + rpc_data->send_response(); + return; + } + + auto merged = + std::make_shared< ScrubManager::scrub_timestamp_info >(scrub_mgr->sync_scrub_timestamp(pg_id, incoming)); + + sisl::io_blob_list_t blob_list; + blob_list.emplace_back(reinterpret_cast< uint8_t* >(merged.get()), + static_cast< uint32_t >(sizeof(ScrubManager::scrub_timestamp_info)), false); + rpc_data->set_comp_cb([merged](boost::intrusive_ptr< sisl::GenericRpcData >&) {}); + rpc_data->send_response(blob_list); +} + // NOTE: caller should hold the _pg_lock const HSHomeObject::HS_PG* HSHomeObject::_get_hs_pg_unlocked(pg_id_t pg_id) const { auto iter = _pg_map.find(pg_id); diff --git a/src/lib/homestore_backend/replication_state_machine.hpp b/src/lib/homestore_backend/replication_state_machine.hpp index 75d2d4187..6d01ea590 100644 --- a/src/lib/homestore_backend/replication_state_machine.hpp +++ b/src/lib/homestore_backend/replication_state_machine.hpp @@ -277,7 +277,7 @@ class ReplicationStateMachine : public homestore::ReplDevListener { std::pair< homestore::repl_lsn_t, homestore::chunk_num_t > get_no_space_left_error_info() const; - void handle_no_space_left(homestore ::repl_lsn_t lsn, homestore ::chunk_num_t chunk_id); + void handle_no_space_left(homestore::repl_lsn_t lsn, homestore::chunk_num_t chunk_id); }; } // namespace homeobject diff --git a/src/lib/homestore_backend/scrub_manager.cpp b/src/lib/homestore_backend/scrub_manager.cpp index 58086b70f..27ec63a4a 100644 --- a/src/lib/homestore_backend/scrub_manager.cpp +++ b/src/lib/homestore_backend/scrub_manager.cpp @@ -230,7 +230,23 @@ bool ScrubManager::is_eligible_for_shallow_scrub(const pg_id_t& pg_id) { return false; } +void ScrubManager::add_missing_pg_scrub_superblks() { + // pgs created before scrub manager tracked scrub superblocks (or that otherwise missed getting one, e.g. a crash + // between create_pg and add_pg) won't have an entry in m_pg_scrub_sb_map. Backfill those here. + std::vector< pg_id_t > pg_ids; + m_hs_home_object->get_pg_ids(pg_ids); + for (auto const& pg_id : pg_ids) { + if (m_pg_scrub_sb_map.find(pg_id) == m_pg_scrub_sb_map.end()) { + LOGINFOMOD(scrubmgr, "pg={} has no scrub superblock, backfilling one", pg_id); + add_pg(pg_id); + } + } +} + void ScrubManager::start() { + // backfill scrub superblocks for any pre-existing pgs that don't have one yet. + add_missing_pg_scrub_superblks(); + // 1 set scrub task handling threads. // TODO :: make thread count configurable, thread number is the most concurrent scrub tasks that can be handled // concurrently. Too many concurrent scrub tasks may bring too much pressure to the node @@ -269,6 +285,10 @@ void ScrubManager::start() { // TODO: make the interval configurable, for now set it to 5 seconds m_retry_timer_hdl = iomanager.schedule_thread_timer(5ull * 1000 * 1000 * 1000, true, nullptr /*cookie*/, [this](void*) { check_scrub_timeouts(); }); + // TODO: make the interval configurable, for now set it to 30 seconds + m_push_scrub_timestamp_timer_hdl = + iomanager.schedule_thread_timer(30ull * 1000 * 1000 * 1000, true, nullptr /*cookie*/, + [this](void*) { push_scrub_timestamp_to_followers(); }); }); LOGINFOMOD(scrubmgr, "scrub manager started!"); } @@ -286,6 +306,11 @@ void ScrubManager::stop() { iomanager.cancel_timer(m_retry_timer_hdl, true); m_retry_timer_hdl = iomgr::null_timer_handle; } + + if (m_push_scrub_timestamp_timer_hdl != iomgr::null_timer_handle) { + iomanager.cancel_timer(m_push_scrub_timestamp_timer_hdl, true); + m_push_scrub_timestamp_timer_hdl = iomgr::null_timer_handle; + } }); m_scrub_timer_fiber = nullptr; } @@ -1072,6 +1097,68 @@ std::optional< ScrubManager::pg_scrub_superblk > ScrubManager::get_scrub_superbl return *(*(it->second)); } +ScrubManager::scrub_timestamp_info ScrubManager::sync_scrub_timestamp(const pg_id_t pg_id, + const scrub_timestamp_info& incoming) { + auto it = m_pg_scrub_sb_map.find(pg_id); + if (it == m_pg_scrub_sb_map.end()) { + LOGWARNMOD(scrubmgr, "scrub superblk not found for pg={}, cannot sync scrub superblk", pg_id); + return incoming; + } + + auto& sb = *(it->second); + bool updated = false; + if (incoming.last_deep_scrub_timestamp > sb->last_deep_scrub_timestamp) { + sb->last_deep_scrub_timestamp = incoming.last_deep_scrub_timestamp; + updated = true; + } + if (incoming.last_shallow_scrub_timestamp > sb->last_shallow_scrub_timestamp) { + sb->last_shallow_scrub_timestamp = incoming.last_shallow_scrub_timestamp; + updated = true; + } + if (updated) { sb.write(); } + + return scrub_timestamp_info{sb->last_deep_scrub_timestamp, sb->last_shallow_scrub_timestamp}; +} + +void ScrubManager::push_scrub_timestamp_to_followers() { + for (auto const& [pg_id, sb] : m_pg_scrub_sb_map) { + auto hs_pg = m_hs_home_object->get_hs_pg(pg_id); + if (!hs_pg || !hs_pg->repl_dev_ || !hs_pg->repl_dev_->is_leader()) { continue; } + + auto outgoing = std::make_shared< scrub_timestamp_info >( + scrub_timestamp_info{(*sb)->last_deep_scrub_timestamp, (*sb)->last_shallow_scrub_timestamp}); + + const auto& self_id = m_hs_home_object->our_uuid(); + for (const auto& member : hs_pg->pg_info_.members) { + if (member.id == self_id) continue; + + sisl::io_blob_list_t blob_list; + blob_list.emplace_back(reinterpret_cast< uint8_t* >(outgoing.get()), sizeof(scrub_timestamp_info), false); + + hs_pg->repl_dev_->data_request_bidirectional(member.id, HSHomeObject::PUSH_SCRUB_TIMESTAMP, blob_list) + .via(folly::getGlobalIOExecutor()) + .thenValue([this, pg_id, peer_id = member.id, outgoing](auto&& response) { + if (response.hasError()) { + LOGWARNMOD(scrubmgr, "failed to push scrub superblk to peer {} for pg={}, error={}", peer_id, + pg_id, response.error()); + return; + } + + auto const& resp_blob = response.value().response_blob(); + if (resp_blob.size() != sizeof(scrub_timestamp_info)) { + LOGWARNMOD(scrubmgr, "invalid scrub superblk push response from peer {} for pg={}, size={}", + peer_id, pg_id, resp_blob.size()); + return; + } + + scrub_timestamp_info incoming; + std::memcpy(&incoming, resp_blob.cbytes(), sizeof(incoming)); + sync_scrub_timestamp(pg_id, incoming); + }); + } + } +} + ScrubManager::PGScrubContext::PGScrubContext(uint64_t task_id, const HSHomeObject::HS_PG* hs_pg) : task_id(task_id), hs_pg(hs_pg) { // Build m_active_batch once; PG membership is stable for the lifetime of this scrub task. diff --git a/src/lib/homestore_backend/scrub_manager.hpp b/src/lib/homestore_backend/scrub_manager.hpp index 8c4fb591c..f9e58f0f6 100644 --- a/src/lib/homestore_backend/scrub_manager.hpp +++ b/src/lib/homestore_backend/scrub_manager.hpp @@ -69,6 +69,15 @@ class ScrubManager { }; #pragma pack() + // wire payload for PUSH_SCRUB_TIMESTAMP: sent by the leader to a follower, and echoed back by the follower with + // its own (possibly merged) values. +#pragma pack(1) + struct scrub_timestamp_info { + uint64_t last_deep_scrub_timestamp{0}; + uint64_t last_shallow_scrub_timestamp{0}; + }; +#pragma pack() + // scrub req struct scrub_req { scrub_req() = default; @@ -282,6 +291,17 @@ class ScrubManager { void save_scrub_superblk(const pg_id_t pg_id, const bool is_deep_scrub, bool force_update = true); void add_scrub_req(std::shared_ptr< scrub_req > req); + // Merges incoming (peer's) scrub superblk timestamps into the local one for pg_id, keeping the newer value for + // each field independently, persisting if anything changed, and returning the resulting local values. + // Called both by the follower's PUSH_SCRUB_TIMESTAMP RPC handler (incoming = leader's values) and by the leader + // itself once a follower's response comes back (incoming = follower's values). + scrub_timestamp_info sync_scrub_timestamp(const pg_id_t pg_id, const scrub_timestamp_info& incoming); + + // leader-only: periodically push last_deep_scrub_timestamp/last_shallow_scrub_timestamp to followers of every pg + // this node leads, and merge back whatever the followers report as newer. Exposed publicly (in addition to being + // called from the periodic timer) so it can be triggered on-demand, e.g. from tests. + void push_scrub_timestamp_to_followers(); + // local scrub std::shared_ptr< scrub_result > local_scrub_blob(std::shared_ptr< scrub_req > req); std::shared_ptr< scrub_result > local_scrub_meta(std::shared_ptr< scrub_req > req); @@ -351,6 +371,7 @@ class ScrubManager { } }; + void add_missing_pg_scrub_superblks(); void scan_pg_for_scrub(); void handle_pg_scrub_task(scrub_task task); bool is_eligible_for_deep_scrub(const pg_id_t& pg_id); @@ -366,6 +387,7 @@ class ScrubManager { iomgr::timer_handle_t m_scrub_timer_hdl{iomgr::null_timer_handle}; iomgr::timer_handle_t m_retry_timer_hdl{iomgr::null_timer_handle}; + iomgr::timer_handle_t m_push_scrub_timestamp_timer_hdl{iomgr::null_timer_handle}; iomgr::io_fiber_t m_scrub_timer_fiber{nullptr}; HSHomeObject* m_hs_home_object{nullptr}; MPMCPriorityQueue< scrub_task > m_scrub_task_queue; diff --git a/src/lib/homestore_backend/snapshot_receive_handler.cpp b/src/lib/homestore_backend/snapshot_receive_handler.cpp index 4ea5cf0da..624a60fe4 100644 --- a/src/lib/homestore_backend/snapshot_receive_handler.cpp +++ b/src/lib/homestore_backend/snapshot_receive_handler.cpp @@ -35,8 +35,8 @@ int HSHomeObject::SnapshotReceiveHandler::process_pg_snapshot_data(ResyncPGMetaD pg_member.priority = member->priority(); pg_info.members.insert(pg_member); } - LOGD("Resync PG membership: pg={}, expected_members={}, members={}", pg_meta.pg_id(), - pg_info.expected_member_num, pg_meta.members()->size()) + LOGD("Resync PG membership: pg={}, expected_members={}, members={}", pg_meta.pg_id(), pg_info.expected_member_num, + pg_meta.members()->size()) #ifdef _PRERELEASE if (iomgr_flip::instance()->test_flip("snapshot_receiver_pg_error")) { @@ -105,8 +105,8 @@ int HSHomeObject::SnapshotReceiveHandler::process_shard_snapshot_data(ResyncShar } ctx_->shard_cursor = shard_meta.shard_id(); ctx_->cur_batch_num = 0; - LOGI("Processed resync shard metadata: pg={}, shard_id=0x{:x}, state={}", shard_meta.pg_id(), - shard_meta.shard_id(), shard_meta.state()); + LOGI("Processed resync shard metadata: pg={}, shard_id=0x{:x}, state={}", shard_meta.pg_id(), shard_meta.shard_id(), + shard_meta.state()); return 0; } @@ -267,8 +267,8 @@ int HSHomeObject::SnapshotReceiveHandler::process_blobs_snapshot_data(ResyncBlob } #endif auto blob_id = blob->blob_id(); - LOGT("Writing resync blob: pg={}, shard_id=0x{:x}, blob={}, blkid={}", ctx_->pg_id, ctx_->shard_cursor, - blob_id, blk_id.to_string()); + LOGT("Writing resync blob: pg={}, shard_id=0x{:x}, blob={}, blkid={}", ctx_->pg_id, ctx_->shard_cursor, blob_id, + blk_id.to_string()); // ToDo: limit the max concurrent? futs.emplace_back( @@ -315,9 +315,9 @@ int HSHomeObject::SnapshotReceiveHandler::process_blobs_snapshot_data(ResyncBlob if (!all_io_submitted || ec != std::error_code{}) { if (!all_io_submitted) { - LOGE( - "Failed to submit complete resync shard batch: pg={}, shard_id=0x{:x}, batch={}, expected_blobs={}, submitted_blobs={}", - ctx_->pg_id, ctx_->shard_cursor, batch_num, data_blobs.blob_list()->size(), futs.size()); + LOGE("Failed to submit complete resync shard batch: pg={}, shard_id=0x{:x}, batch={}, expected_blobs={}, " + "submitted_blobs={}", + ctx_->pg_id, ctx_->shard_cursor, batch_num, data_blobs.blob_list()->size(), futs.size()); } else { LOGE("Failed to write resync shard batch: pg={}, shard_id=0x{:x}, batch={}, error_code={}, error={}", ctx_->pg_id, ctx_->shard_cursor, batch_num, ec.value(), ec.message()); diff --git a/src/lib/homestore_backend/tests/homeobj_misc_tests.cpp b/src/lib/homestore_backend/tests/homeobj_misc_tests.cpp index 6848063dc..2c0de9c21 100644 --- a/src/lib/homestore_backend/tests/homeobj_misc_tests.cpp +++ b/src/lib/homestore_backend/tests/homeobj_misc_tests.cpp @@ -508,10 +508,9 @@ TEST_F(HomeObjectFixture, SnapshotReceiveHandler) { for (uint64_t i = 1; i <= num_shards_per_pg; i++) { shard_ids.push_back(i); } - auto pg_entry = - CreateResyncPGMetaDataDirect(builder, pg_id, &uuid, pg->pg_info_.size, pg->pg_info_.expected_member_num, - pg->pg_info_.chunk_size, blob_seq_num, num_shards_per_pg, &members, &shard_ids, - blob_seq_num /* total_blobs_to_transfer */); + auto pg_entry = CreateResyncPGMetaDataDirect( + builder, pg_id, &uuid, pg->pg_info_.size, pg->pg_info_.expected_member_num, pg->pg_info_.chunk_size, + blob_seq_num, num_shards_per_pg, &members, &shard_ids, blob_seq_num /* total_blobs_to_transfer */); builder.Finish(pg_entry); auto pg_meta = GetResyncPGMetaData(builder.GetBufferPointer()); auto ret = handler->process_pg_snapshot_data(*pg_meta); @@ -672,7 +671,8 @@ TEST_F(HomeObjectFixture, SnapshotReceiveHandler) { builder, retry_blob_id, static_cast< uint8_t >(ResyncBlobState::NORMAL), &retry_data)); blob_map[retry_blob_id] = std::make_tuple< Blob, bool >(std::move(retry_blob), false); } - builder.Finish(CreateResyncBlobDataBatchDirect(builder, &retry_blob_entries, j == num_batches_per_shard)); + builder.Finish( + CreateResyncBlobDataBatchDirect(builder, &retry_blob_entries, j == num_batches_per_shard)); auto retry_blob_batch = GetResyncBlobDataBatch(builder.GetBufferPointer()); ASSERT_EQ(handler->process_blobs_snapshot_data(*retry_blob_batch, j, j == num_batches_per_shard), 0); } else { @@ -759,10 +759,9 @@ TEST_F(HomeObjectFixture, SnapshotReceiveHandlerAllocatorResyncAfterCrash) { members.push_back(CreateMemberDirect(builder, &id, member.name.c_str(), priority)); } std::vector< uint64_t > shard_ids = {1}; - auto pg_entry = CreateResyncPGMetaDataDirect(builder, pg_id, &uuid, pg->pg_info_.size, - pg->pg_info_.expected_member_num, pg->pg_info_.chunk_size, - num_blobs_total /*blob_seq_num*/, 1 /*shard_seq_num*/, &members, - &shard_ids); + auto pg_entry = CreateResyncPGMetaDataDirect( + builder, pg_id, &uuid, pg->pg_info_.size, pg->pg_info_.expected_member_num, pg->pg_info_.chunk_size, + num_blobs_total /*blob_seq_num*/, 1 /*shard_seq_num*/, &members, &shard_ids); builder.Finish(pg_entry); ASSERT_EQ(handler->process_pg_snapshot_data(*GetResyncPGMetaData(builder.GetBufferPointer())), 0); builder.Reset(); @@ -816,7 +815,8 @@ TEST_F(HomeObjectFixture, SnapshotReceiveHandlerAllocatorResyncAfterCrash) { uint64_t total_bytes{0}; for (blob_id_t blob_id = start_blob_id; blob_id < end_blob_id; blob_id++) { auto blob = build_blob(blob_id); - total_bytes += sisl::round_up(sizeof(HSHomeObject::BlobHeader), _obj_inst->_data_block_size) + blob.body.size(); + total_bytes += + sisl::round_up(sizeof(HSHomeObject::BlobHeader), _obj_inst->_data_block_size) + blob.body.size(); } return total_bytes; }; diff --git a/src/lib/homestore_backend/tests/hs_scrubber_tests.cpp b/src/lib/homestore_backend/tests/hs_scrubber_tests.cpp index eb06b1d44..1345f0f7b 100644 --- a/src/lib/homestore_backend/tests/hs_scrubber_tests.cpp +++ b/src/lib/homestore_backend/tests/hs_scrubber_tests.cpp @@ -769,6 +769,90 @@ TEST_F(HomeObjectFixture, ScrubSuperblockPersistenceTest) { g_helper->sync(); } +// Test PUSH_SCRUB_TIMESTAMP sync: the leader periodically pushes its last_deep_scrub_timestamp/ +// last_shallow_scrub_timestamp to followers; whichever side has the newer value for each field wins +// and propagates to the other side. Exercises both directions: leader newer (follower catches up) +// and follower newer (leader catches up). +TEST_F(HomeObjectFixture, ScrubSuperblkPushSyncTest) { + const pg_id_t pg_id = 1; + create_pg(pg_id); + auto scrub_mgr = _obj_inst->scrub_manager(); + + g_helper->sync(); + + // ===== Leader newer: follower should catch up to the leader's values ===== + std::this_thread::sleep_for(std::chrono::seconds(2)); + run_on_pg_leader(pg_id, [&]() { + scrub_mgr->save_scrub_superblk(pg_id, true /* is_deep_scrub */, true /* force_update */); + scrub_mgr->save_scrub_superblk(pg_id, false /* is_deep_scrub */, true /* force_update */); + }); + + g_helper->sync(); + + uint64_t follower_deep_before = 0; + uint64_t follower_shallow_before = 0; + run_on_pg_follower(pg_id, [&]() { + auto sb = scrub_mgr->get_scrub_superblk(pg_id); + ASSERT_TRUE(sb.has_value()); + follower_deep_before = sb->last_deep_scrub_timestamp; + follower_shallow_before = sb->last_shallow_scrub_timestamp; + }); + + run_on_pg_leader(pg_id, [&]() { scrub_mgr->push_scrub_timestamp_to_followers(); }); + + run_on_pg_follower(pg_id, [&]() { + auto deadline = std::chrono::steady_clock::now() + std::chrono::seconds(15); + bool caught_up = false; + while (std::chrono::steady_clock::now() < deadline) { + auto sb = scrub_mgr->get_scrub_superblk(pg_id); + ASSERT_TRUE(sb.has_value()); + if (sb->last_deep_scrub_timestamp > follower_deep_before && + sb->last_shallow_scrub_timestamp > follower_shallow_before) { + caught_up = true; + break; + } + std::this_thread::sleep_for(std::chrono::milliseconds(200)); + } + EXPECT_TRUE(caught_up) << "Follower should catch up to leader's newer scrub superblk timestamps"; + }); + + g_helper->sync(); + + // ===== Follower newer: leader should catch up to the follower's values ===== + std::this_thread::sleep_for(std::chrono::seconds(2)); + run_on_pg_follower(pg_id, [&]() { + scrub_mgr->save_scrub_superblk(pg_id, true /* is_deep_scrub */, true /* force_update */); + scrub_mgr->save_scrub_superblk(pg_id, false /* is_deep_scrub */, true /* force_update */); + }); + + g_helper->sync(); + + run_on_pg_leader(pg_id, [&]() { + auto before_sb = scrub_mgr->get_scrub_superblk(pg_id); + ASSERT_TRUE(before_sb.has_value()); + const auto leader_deep_before = before_sb->last_deep_scrub_timestamp; + const auto leader_shallow_before = before_sb->last_shallow_scrub_timestamp; + + scrub_mgr->push_scrub_timestamp_to_followers(); + + auto deadline = std::chrono::steady_clock::now() + std::chrono::seconds(15); + bool caught_up = false; + while (std::chrono::steady_clock::now() < deadline) { + auto after_sb = scrub_mgr->get_scrub_superblk(pg_id); + ASSERT_TRUE(after_sb.has_value()); + if (after_sb->last_deep_scrub_timestamp > leader_deep_before && + after_sb->last_shallow_scrub_timestamp > leader_shallow_before) { + caught_up = true; + break; + } + std::this_thread::sleep_for(std::chrono::milliseconds(200)); + } + EXPECT_TRUE(caught_up) << "Leader should catch up to follower's newer scrub superblk timestamps"; + }); + + g_helper->sync(); +} + // Test cancel scrub task TEST_F(HomeObjectFixture, CancelScrubTaskTest) { const pg_id_t pg_id = 1; From 748f19b559e3b3c71015056d60673a6f1e4db7d7 Mon Sep 17 00:00:00 2001 From: Jie Yao Date: Thu, 1 Oct 2026 16:31:50 +0800 Subject: [PATCH 2/2] add lock for scrub superblk --- src/lib/homestore_backend/scrub_manager.cpp | 72 +++++++++++++-------- src/lib/homestore_backend/scrub_manager.hpp | 11 +++- 2 files changed, 55 insertions(+), 28 deletions(-) diff --git a/src/lib/homestore_backend/scrub_manager.cpp b/src/lib/homestore_backend/scrub_manager.cpp index 27ec63a4a..1af06b5d2 100644 --- a/src/lib/homestore_backend/scrub_manager.cpp +++ b/src/lib/homestore_backend/scrub_manager.cpp @@ -822,9 +822,13 @@ ScrubManager::submit_scrub_task(const pg_id_t& pg_id, const bool is_deep, SCRUB_ return folly::makeSemiFuture(std::shared_ptr< ScrubManager::ShallowScrubReport >(nullptr)); } - const auto& pg_scrub_sb = *(ps_scrub_super_blk_it->second); - const auto last_scrub_time = - is_deep ? pg_scrub_sb->last_deep_scrub_timestamp : pg_scrub_sb->last_shallow_scrub_timestamp; + auto& pg_scrub_sb_entry = *(ps_scrub_super_blk_it->second); + uint64_t last_scrub_time; + { + std::lock_guard lock(pg_scrub_sb_entry.mutex); + last_scrub_time = is_deep ? pg_scrub_sb_entry.sb->last_deep_scrub_timestamp + : pg_scrub_sb_entry.sb->last_shallow_scrub_timestamp; + } auto [promise, future] = folly::makePromiseContract< std::shared_ptr< ShallowScrubReport > >(); ScrubManager::scrub_task task(last_scrub_time, pg_id, is_deep, trigger_type, std::move(promise)); @@ -1030,7 +1034,7 @@ void ScrubManager::remove_pg(const pg_id_t pg_id) { } LOGINFOMOD(scrubmgr, "removed pg={} in scrub manager!", pg_id); - it->second->destroy(); + it->second->sb.destroy(); m_pg_scrub_sb_map.erase(it); } @@ -1038,21 +1042,21 @@ void ScrubManager::remove_pg(const pg_id_t pg_id) { void ScrubManager::on_pg_scrub_meta_blk_found( sisl::byte_view const& buf, void* meta_cookie, std::vector< homestore::superblk< pg_scrub_superblk > >& stale_pg_scrub_sbs) { - auto sb = std::make_shared< homestore::superblk< pg_scrub_superblk > >(); - (*sb).load(buf, meta_cookie); - const auto pg_id = (*sb)->pg_id; + auto entry = std::make_shared< PgScrubSbEntry >(); + entry->sb.load(buf, meta_cookie); + const auto pg_id = entry->sb->pg_id; auto hs_pg = m_hs_home_object->get_hs_pg(pg_id); if (!hs_pg) { // this is a stale pg scrub superblock, we just log and destroy it. LOGINFOMOD(scrubmgr, "cannot find pg={}, destroy stale scrub superblock", pg_id); - stale_pg_scrub_sbs.emplace_back(std::move(*sb)); + stale_pg_scrub_sbs.emplace_back(std::move(entry->sb)); return; } - const auto last_deep_scrub_time = (*sb)->last_deep_scrub_timestamp; - const auto last_shallow_scrub_time = (*sb)->last_shallow_scrub_timestamp; + const auto last_deep_scrub_time = entry->sb->last_deep_scrub_timestamp; + const auto last_shallow_scrub_time = entry->sb->last_shallow_scrub_timestamp; - m_pg_scrub_sb_map.emplace(pg_id, std::move(sb)); + m_pg_scrub_sb_map.emplace(pg_id, std::move(entry)); LOGINFOMOD(scrubmgr, "loaded scrub superblock for pg={}, last_deep_scrub_time={}, last_shallow_scrub_time={}", pg_id, last_deep_scrub_time, last_shallow_scrub_time); } @@ -1064,24 +1068,27 @@ void ScrubManager::save_scrub_superblk(const pg_id_t pg_id, const bool is_deep_s auto it = m_pg_scrub_sb_map.find(pg_id); if (it == m_pg_scrub_sb_map.end()) { // Create new superblock for this PG - auto sb = std::make_shared< homestore::superblk< pg_scrub_superblk > >(pg_scrub_meta_name); - (*sb).create(sizeof(pg_scrub_superblk)); - (*sb)->pg_id = pg_id; - (*sb)->last_deep_scrub_timestamp = current_time; - (*sb)->last_shallow_scrub_timestamp = current_time; - (*sb).write(); - m_pg_scrub_sb_map.emplace(pg_id, std::move(sb)); + auto entry = std::make_shared< PgScrubSbEntry >(pg_scrub_meta_name); + entry->sb.create(sizeof(pg_scrub_superblk)); + entry->sb->pg_id = pg_id; + entry->sb->last_deep_scrub_timestamp = current_time; + entry->sb->last_shallow_scrub_timestamp = current_time; + entry->sb.write(); + m_pg_scrub_sb_map.emplace(pg_id, std::move(entry)); return; } if (force_update) { - // Update existing superblock + // Update existing superblock. Guard against concurrent mutation/persist from sync_scrub_timestamp (RPC + // handler / leader's follower-response callbacks) racing on the same pg's superblk. + auto& entry = *(it->second); + std::lock_guard lock(entry.mutex); if (is_deep_scrub) { - (*(it->second))->last_deep_scrub_timestamp = current_time; + entry.sb->last_deep_scrub_timestamp = current_time; } else { - (*(it->second))->last_shallow_scrub_timestamp = current_time; + entry.sb->last_shallow_scrub_timestamp = current_time; } - (*(it->second)).write(); + entry.sb.write(); } else { LOGINFOMOD(scrubmgr, "skip updating scrub superblock for pg={} since there is no scrub progress update", pg_id); } @@ -1094,7 +1101,8 @@ std::optional< ScrubManager::pg_scrub_superblk > ScrubManager::get_scrub_superbl return std::nullopt; } - return *(*(it->second)); + std::lock_guard lock(it->second->mutex); + return *(it->second->sb); } ScrubManager::scrub_timestamp_info ScrubManager::sync_scrub_timestamp(const pg_id_t pg_id, @@ -1105,7 +1113,12 @@ ScrubManager::scrub_timestamp_info ScrubManager::sync_scrub_timestamp(const pg_i return incoming; } - auto& sb = *(it->second); + // Can be invoked concurrently for the same pg_id, e.g. once per follower from the leader's + // push_scrub_timestamp_to_followers() response callbacks, or by overlapping PUSH_SCRUB_TIMESTAMP RPCs on a + // follower. Serialize the compare-update-persist sequence so updates aren't lost or torn. + auto& entry = *(it->second); + std::lock_guard lock(entry.mutex); + auto& sb = entry.sb; bool updated = false; if (incoming.last_deep_scrub_timestamp > sb->last_deep_scrub_timestamp) { sb->last_deep_scrub_timestamp = incoming.last_deep_scrub_timestamp; @@ -1121,12 +1134,17 @@ ScrubManager::scrub_timestamp_info ScrubManager::sync_scrub_timestamp(const pg_i } void ScrubManager::push_scrub_timestamp_to_followers() { - for (auto const& [pg_id, sb] : m_pg_scrub_sb_map) { + for (auto const& [pg_id, entry] : m_pg_scrub_sb_map) { auto hs_pg = m_hs_home_object->get_hs_pg(pg_id); if (!hs_pg || !hs_pg->repl_dev_ || !hs_pg->repl_dev_->is_leader()) { continue; } - auto outgoing = std::make_shared< scrub_timestamp_info >( - scrub_timestamp_info{(*sb)->last_deep_scrub_timestamp, (*sb)->last_shallow_scrub_timestamp}); + scrub_timestamp_info outgoing_info; + { + std::lock_guard lock(entry->mutex); + outgoing_info = + scrub_timestamp_info{entry->sb->last_deep_scrub_timestamp, entry->sb->last_shallow_scrub_timestamp}; + } + auto outgoing = std::make_shared< scrub_timestamp_info >(outgoing_info); const auto& self_id = m_hs_home_object->our_uuid(); for (const auto& member : hs_pg->pg_info_.members) { diff --git a/src/lib/homestore_backend/scrub_manager.hpp b/src/lib/homestore_backend/scrub_manager.hpp index f9e58f0f6..f7275bd3e 100644 --- a/src/lib/homestore_backend/scrub_manager.hpp +++ b/src/lib/homestore_backend/scrub_manager.hpp @@ -78,6 +78,15 @@ class ScrubManager { }; #pragma pack() + // pg_scrub_superblk is mutated and persisted from multiple contexts (local scrub completion, the + // PUSH_SCRUB_TIMESTAMP RPC handler, and the leader's per-follower response callbacks, possibly concurrently for + // the same pg). The mutex guards the read-modify-write-persist sequence on sb. + struct PgScrubSbEntry { + explicit PgScrubSbEntry(const std::string& sub_name = "") : sb(sub_name) {} + homestore::superblk< pg_scrub_superblk > sb; + std::mutex mutex; + }; + // scrub req struct scrub_req { scrub_req() = default; @@ -393,7 +402,7 @@ class ScrubManager { MPMCPriorityQueue< scrub_task > m_scrub_task_queue; std::shared_ptr< folly::IOThreadPoolExecutor > m_scrub_executor; folly::ConcurrentHashMap< pg_id_t, std::shared_ptr< PGScrubContext > > m_pg_scrub_ctx_map; - folly::ConcurrentHashMap< pg_id_t, std::shared_ptr< homestore::superblk< pg_scrub_superblk > > > m_pg_scrub_sb_map; + folly::ConcurrentHashMap< pg_id_t, std::shared_ptr< PgScrubSbEntry > > m_pg_scrub_sb_map; std::shared_ptr< folly::IOThreadPoolExecutor > m_scrub_req_executor; };