From 53de91ffa16399922ef0210f4e32905e4a5f4829 Mon Sep 17 00:00:00 2001 From: "Kou, Jilong" Date: Wed, 29 Jul 2026 15:06:30 +0800 Subject: [PATCH] Fix batch resend bug caused by prefetch limit Retained look-ahead data could exhaust the donor prefetch budget while rebuilding a resent batch, preventing a required earlier blob from being loaded and causing each retry to fail at the same cursor, sticking the resync progress at the same point. Allow blobs preceding the retained prefetch frontier to bypass the bounded look-ahead budget, and cover the scenario with a focused regression test. --- conanfile.py | 2 +- .../homestore_backend/pg_blob_iterator.cpp | 15 ++-- .../tests/homeobj_misc_tests.cpp | 81 ++++++++++++++++++- 3 files changed, 89 insertions(+), 9 deletions(-) diff --git a/conanfile.py b/conanfile.py index 4b6ddf57..9e46f9a6 100644 --- a/conanfile.py +++ b/conanfile.py @@ -10,7 +10,7 @@ class HomeObjectConan(ConanFile): name = "homeobject" - version = "4.2.7" + version = "4.2.8" homepage = "https://github.com/eBay/HomeObject" description = "Blob Store built on HomeStore" diff --git a/src/lib/homestore_backend/pg_blob_iterator.cpp b/src/lib/homestore_backend/pg_blob_iterator.cpp index 0087882d..5d9f3808 100644 --- a/src/lib/homestore_backend/pg_blob_iterator.cpp +++ b/src/lib/homestore_backend/pg_blob_iterator.cpp @@ -320,7 +320,11 @@ bool HSHomeObject::PGBlobIterator::prefetch_blobs_snapshot_data() { LOGD("prefetch_blobs_snapshot_data, inflight={}, idx={}, max_batch_size * 2 = {}", inflight_prefetch_bytes_, cur_start_blob_idx_, max_batch_size_ * 2); auto idx = cur_start_blob_idx_; - while (inflight_prefetch_bytes_ < max_batch_size_ * 2 && idx < cur_blob_list_.size()) { + // On batch resend, retained look-ahead may already consume part of the 2x budget. Allow missing blobs before the + // earliest retained blob to bypass the limit so the current batch can always be rebuilt. + const auto prefetch_frontier = prefetched_blobs_.empty() ? blob_id_t{0} : prefetched_blobs_.begin()->first; + while (idx < cur_blob_list_.size() && + (inflight_prefetch_bytes_ < max_batch_size_ * 2 || cur_blob_list_[idx].blob_id < prefetch_frontier)) { auto info = cur_blob_list_[idx++]; total_blobs++; // handle deleted object @@ -337,9 +341,8 @@ bool HSHomeObject::PGBlobIterator::prefetch_blobs_snapshot_data() { skipped_blobs++; continue; } - auto expect_blob_size = info.pbas.blk_count() * repl_dev_->get_blk_size(); - inflight_prefetch_bytes_ += expect_blob_size; - LOGD("will prefetch {}", info.blob_id); + inflight_prefetch_bytes_ += info.pbas.blk_count() * repl_dev_->get_blk_size(); + LOGD("Will prefetch {}: frontier {}, inflight {}", info.blob_id, prefetch_frontier, inflight_prefetch_bytes_); prefetch_list.emplace_back(info); } // POC: sort the prefetch_list by pbas, trying to let IO submitted to disk more sequential. @@ -426,8 +429,10 @@ bool HSHomeObject::PGBlobIterator::create_blobs_snapshot_data(sisl::io_blob_safe LOGE("blob {} not found in prefetched blob map", info.blob_id); break; } + auto const expect_blob_size = info.pbas.blk_count() * repl_dev_->get_blk_size(); auto res = std::move(it->second).get(); prefetched_blobs_.erase(it); + inflight_prefetch_bytes_ -= expect_blob_size; if (res.hasError()) { LOGE("blob {} hit error {}", info.blob_id, res.error()); @@ -436,8 +441,6 @@ bool HSHomeObject::PGBlobIterator::create_blobs_snapshot_data(sisl::io_blob_safe } std::vector< uint8_t > data(res->blob_.cbytes(), res->blob_.cbytes() + res->blob_.size()); blob_entries.push_back(CreateResyncBlobDataDirect(builder_, res->blob_id_, (uint8_t)res->state_, &data)); - auto const expect_blob_size = info.pbas.blk_count() * repl_dev_->get_blk_size(); - inflight_prefetch_bytes_ -= expect_blob_size; total_bytes += expect_blob_size; fetched_blobs++; } diff --git a/src/lib/homestore_backend/tests/homeobj_misc_tests.cpp b/src/lib/homestore_backend/tests/homeobj_misc_tests.cpp index 30a9dd65..6848063d 100644 --- a/src/lib/homestore_backend/tests/homeobj_misc_tests.cpp +++ b/src/lib/homestore_backend/tests/homeobj_misc_tests.cpp @@ -203,6 +203,80 @@ TEST_F(HomeObjectFixture, PGBlobIterator) { ASSERT_TRUE(pg_iter->update_cursor(objId(LAST_OBJ_ID))); } +// Regression test for a batch resend bug with retained look-ahead prefetch data, where the old bounded prefetch loop +// could stop after rebuilding only a prefix of the resent batch when retained bytes plus that prefix reached the 2x +// batch limit. The required latter blob(s) was then missing from prefetched_blobs_, and every retry failed at the +// same cursor. Derive a batch limit from actual stored blob sizes to reproduce that topology without depending on +// random sizes. +TEST_F(HomeObjectFixture, PGBlobIteratorBatchResendPrefetchFrontier) { + constexpr pg_id_t pg_id{1}; + constexpr uint64_t num_blobs{8}; + + create_pg(pg_id); + auto const shard = create_shard(pg_id, 64 * Mi, "shard meta"); + std::map< pg_id_t, std::vector< shard_id_t > > pg_shard_map{{pg_id, {shard.id}}}; + std::map< pg_id_t, blob_id_t > pg_blob_id{{pg_id, 0}}; + put_blobs(pg_shard_map, num_blobs, pg_blob_id); + + auto pg = _obj_inst->get_hs_pg(pg_id); + ASSERT_NE(pg, nullptr); + auto pg_iter = + std::make_shared< HSHomeObject::PGBlobIterator >(*_obj_inst, pg->pg_info_.replica_set_uuid, shard.create_lsn); + auto const shard_seq_num = HSHomeObject::get_sequence_num_from_shard_id(shard.id); + ASSERT_TRUE(pg_iter->update_cursor(objId(shard_seq_num, 0))); + ASSERT_TRUE(pg_iter->generate_shard_blob_list()); + + // Find a batch limit B from the actual stored sizes such that initial prefetch retains look-ahead, while on resend: + // retained look-ahead + a rebuilt batch prefix >= 2B before the final required batch blob is prefetched. + std::vector< uint64_t > prefix_bytes{0}; + for (auto const& info : pg_iter->cur_blob_list_) { + prefix_bytes.push_back(prefix_bytes.back() + info.pbas.blk_count() * pg_iter->repl_dev_->get_blk_size()); + } + bool topology_configured = false; + for (size_t batch_blobs = 2; batch_blobs < prefix_bytes.size() - 1 && !topology_configured; batch_blobs++) { + for (size_t prefetched_blobs = batch_blobs + 1; prefetched_blobs < prefix_bytes.size(); prefetched_blobs++) { + auto const min_batch_size = + std::max(prefix_bytes[batch_blobs - 1], prefix_bytes[prefetched_blobs - 1] / 2) + 1; + auto const max_batch_size = std::min( + {prefix_bytes[batch_blobs], prefix_bytes[prefetched_blobs] / 2, + (prefix_bytes[prefetched_blobs] - prefix_bytes[batch_blobs] + prefix_bytes[batch_blobs - 1]) / 2}); + if (min_batch_size <= max_batch_size) { + pg_iter->max_batch_size_ = min_batch_size; + topology_configured = true; + break; + } + } + } + ASSERT_TRUE(topology_configured); + + sisl::io_blob_safe shard_meta; + ASSERT_TRUE(pg_iter->create_shard_snapshot_data(shard_meta)); + objId batch_oid(shard_seq_num, 1); + ASSERT_TRUE(pg_iter->update_cursor(batch_oid)); + sisl::io_blob_safe blob_batch; + ASSERT_TRUE(pg_iter->create_blobs_snapshot_data(blob_batch)); + auto blob_msg = GetSizePrefixedResyncBlobDataBatch(blob_batch.cbytes() + sizeof(SyncMessageHeader)); + ASSERT_GT(blob_msg->blob_list()->size(), 1); + ASSERT_FALSE(pg_iter->prefetched_blobs_.empty()); + + uint64_t resend_prefix_bytes{0}; + for (uint32_t i = 0; i + 1 < blob_msg->blob_list()->size(); i++) { + auto const blob_id = blob_msg->blob_list()->Get(i)->blob_id(); + auto const info = std::ranges::find(pg_iter->cur_blob_list_, blob_id, &HSHomeObject::BlobInfo::blob_id); + ASSERT_NE(info, pg_iter->cur_blob_list_.end()); + resend_prefix_bytes += info->pbas.blk_count() * pg_iter->repl_dev_->get_blk_size(); + } + ASSERT_GE(pg_iter->inflight_prefetch_bytes_ + resend_prefix_bytes, pg_iter->max_batch_size_ * 2); + ASSERT_LT(blob_msg->blob_list()->Get(blob_msg->blob_list()->size() - 1)->blob_id(), + pg_iter->prefetched_blobs_.begin()->first); + + ASSERT_TRUE(pg_iter->update_cursor(batch_oid)); + sisl::io_blob_safe resent_blob_batch; + ASSERT_TRUE(pg_iter->create_blobs_snapshot_data(resent_blob_batch)); + ASSERT_EQ(resent_blob_batch.size(), blob_batch.size()); + ASSERT_EQ(std::memcmp(resent_blob_batch.cbytes(), blob_batch.cbytes(), blob_batch.size()), 0); +} + // Simulates a GC race where blob blkID changes between generate_shard_blob_list and load_blob_data. // Injects a cross-shard pbas into cur_blob_list_ so verify_blob fails (shard_id mismatch in blob header). // The retry inside load_blob_data_with_blkid re-reads the index, finds the updated pbas, and re-reads @@ -512,8 +586,11 @@ TEST_F(HomeObjectFixture, SnapshotReceiveHandler) { std::vector< flatbuffers::Offset< ResyncBlobData > > blob_entries; auto const batch_start_blob_id = cur_blob_id; for (uint64_t k = 0; k < num_blobs_per_batch; k++) { - auto blob_state = corrupt_dis(gen) <= corrupted_blob_percentage ? ResyncBlobState::CORRUPTED - : ResyncBlobState::NORMAL; + // An injected batch corruption must affect a normal blob so the receiver verifies and rejects it. + // CORRUPTED blobs intentionally bypass verification and are covered by clean batches below. + auto blob_state = is_corrupted_batch || corrupt_dis(gen) > corrupted_blob_percentage + ? ResyncBlobState::NORMAL + : ResyncBlobState::CORRUPTED; // Construct raw blob buffer auto blob = build_blob(cur_blob_id);