From 0be20764df1d0d9a26f76960f003bbe8be46841c Mon Sep 17 00:00:00 2001 From: Niko Nastonen Date: Thu, 10 Sep 2026 10:08:04 -0700 Subject: [PATCH] SDSTOR-20656: clear wbc chunk selector --- conanfile.py | 2 +- .../tests/test_volume_chunk_selector.cpp | 207 +++++++++++ src/lib/volume/volume_chunk_selector.cpp | 343 ++++++++++++------ src/lib/volume/volume_chunk_selector.hpp | 32 +- 4 files changed, 472 insertions(+), 112 deletions(-) diff --git a/conanfile.py b/conanfile.py index ab53d813..1251fe3c 100644 --- a/conanfile.py +++ b/conanfile.py @@ -10,7 +10,7 @@ class HomeBlocksConan(ConanFile): name = "homeblocks" - version = "6.0.10" + version = "6.0.11" homepage = "https://github.com/eBay/HomeBlocks" description = "Block Store built on HomeStore" diff --git a/src/lib/volume/tests/test_volume_chunk_selector.cpp b/src/lib/volume/tests/test_volume_chunk_selector.cpp index 9312be97..e2a8e234 100644 --- a/src/lib/volume/tests/test_volume_chunk_selector.cpp +++ b/src/lib/volume/tests/test_volume_chunk_selector.cpp @@ -9,6 +9,12 @@ #include #include "test_common.hpp" +#include +#include +#include +#include +#include + SISL_LOGGING_INIT(HOMEBLOCKS_LOG_MODS) SISL_OPTION_GROUP(test_volume_chunk_selector, (num_vols, "", "num_vols", "number of volumes", ::cxxopts::value< uint32_t >()->default_value("2"), @@ -78,11 +84,25 @@ uint16_t VChunk::get_chunk_id() const { return m_internal_chunk->get_chunk_id(); blk_num_t VChunk::get_total_blks() const { return m_internal_chunk->get_total_blks(); } uint64_t VChunk::size() const { return m_internal_chunk->size(); } + void VChunk::reset() {} + cshared< Chunk > VChunk::get_internal_chunk() const { return m_internal_chunk; } +// void VChunk::reset_block_allocator() {} + } // namespace homestore +template < typename Pred > +void wait_until(Pred&& pred, std::chrono::milliseconds timeout = std::chrono::milliseconds{5000}, + std::chrono::milliseconds poll = std::chrono::milliseconds{50}) { + auto deadline = std::chrono::steady_clock::now() + timeout; + while (!pred() && std::chrono::steady_clock::now() < deadline) { + std::this_thread::sleep_for(poll); + } + ASSERT_TRUE(pred()); +} + class ChunkSelectorTest : public ::testing::Test { public: ChunkSelectorTest() {} @@ -148,6 +168,9 @@ TEST_F(ChunkSelectorTest, AllocateReleaseChunksTest) { chunk_sel->release_chunks(i /* ordinal */); } + wait_until([&] { return chunk_sel->num_free_chunks() == init_num_free_chunks; }, std::chrono::seconds{10}, + std::chrono::milliseconds{100}); + // All the chunks will be free. RELEASE_ASSERT_EQ(init_num_free_chunks, chunk_sel->num_free_chunks(), "num free chunks mismatch"); } @@ -232,6 +255,190 @@ TEST_F(ChunkSelectorTest, RecoverChunksTest) { RELEASE_ASSERT(chunk, "Chunk not available"); } +TEST_F(ChunkSelectorTest, ConcurrentAllocateSameVolumeTest) { + auto chunk_sel = + std::make_shared< VolumeChunkSelector >("test", [this](uint64_t, const std::vector< chunk_num_t >&) {}); + + add_chunks_per_pdev(chunk_sel, 2 /*pdevs*/, 10 /*chunks per pdev*/); + + std::vector< chunk_num_t > r1, r2; + uint32_t pdev1{UINT32_MAX}, pdev2{UINT32_MAX}; + std::atomic< bool > go{false}; + + auto worker = [&](std::vector< chunk_num_t >& out, uint32_t& pdev) { + while (!go.load(std::memory_order_acquire)) { + std::this_thread::yield(); + } + out = chunk_sel->allocate_init_chunks(0 /* same volume */, 64 * Ki, pdev); + }; + + std::thread t1(worker, std::ref(r1), std::ref(pdev1)); + std::thread t2(worker, std::ref(r2), std::ref(pdev2)); + + go.store(true, std::memory_order_release); + t1.join(); + t2.join(); + + ASSERT_FALSE(r1.empty()); + ASSERT_FALSE(r2.empty()); + + std::sort(r1.begin(), r1.end()); + std::sort(r2.begin(), r2.end()); + + EXPECT_EQ(r1, r2); + EXPECT_EQ(pdev1, pdev2); + + auto chunks = chunk_sel->get_chunks(0); + ASSERT_EQ(chunks.size(), r1.size()); + + std::unordered_set< chunk_num_t > uniq{r1.begin(), r1.end()}; + EXPECT_EQ(uniq.size(), r1.size()); +} + +TEST_F(ChunkSelectorTest, SelectChunkReturnsNullWhenOutOfSpaceTest) { + auto chunk_sel = + std::make_shared< VolumeChunkSelector >("test", [this](uint64_t, const std::vector< chunk_num_t >&) {}); + + // Single chunk volume, but no space in the chunk + auto chunk = std::make_shared< homestore::Chunk >(0 /*pdev*/, 0 /*chunk_id*/); + chunk->set_available_blks(0); + chunk_sel->add_chunk(chunk); + + uint32_t pdev_id{UINT32_MAX}; + auto ids = chunk_sel->allocate_init_chunks(0 /* ordinal */, 8 * Ki /* <= one chunk */, pdev_id, false); + ASSERT_FALSE(ids.empty()); + + homestore::blk_alloc_hints hints; + hints.application_hint = 0; + + auto start = std::chrono::steady_clock::now(); + auto out = chunk_sel->select_chunk(1 /* nblks */, hints); + auto dur = std::chrono::steady_clock::now() - start; + + EXPECT_EQ(out, nullptr); + EXPECT_LT(dur, std::chrono::seconds(2)); +} + +TEST_F(ChunkSelectorTest, ConcurrentSelectAndReleaseTest) { + auto chunk_sel = + std::make_shared< VolumeChunkSelector >("test", [this](uint64_t, const std::vector< chunk_num_t >&) {}); + + add_chunks_per_pdev(chunk_sel, 2 /*pdevs*/, 10 /*chunks per pdev*/); + const auto total_free_before_alloc = chunk_sel->num_free_chunks(); + + uint32_t pdev_id{UINT32_MAX}; + auto ids = chunk_sel->allocate_init_chunks(0 /* ordinal */, 64 * Ki, pdev_id); + ASSERT_FALSE(ids.empty()); + + homestore::blk_alloc_hints hints; + hints.application_hint = 0; + + std::atomic< bool > done{false}; + std::atomic< uint32_t > select_success{0}; + std::atomic< uint32_t > select_null{0}; + + std::thread selector([&] { + while (!done.load(std::memory_order_acquire)) { + auto chunk = chunk_sel->select_chunk(1 /* nblks */, hints); + if (chunk) { + ++select_success; + } else { + ++select_null; + break; + } + } + }); + + std::this_thread::sleep_for(std::chrono::milliseconds(100)); + std::thread releaser([&] { + chunk_sel->release_chunks(0 /* ordinal */); + done.store(true, std::memory_order_release); + }); + + releaser.join(); + selector.join(); + + wait_until([&] { return chunk_sel->num_free_chunks() == total_free_before_alloc; }, + std::chrono::milliseconds{10000}); + + EXPECT_EQ(chunk_sel->num_free_chunks(), total_free_before_alloc); + EXPECT_TRUE(chunk_sel->get_chunks(0).empty()); + EXPECT_GT(select_success.load() + select_null.load(), 0u); +} + +TEST_F(ChunkSelectorTest, ResizeCallbackBlockedThenReleaseTest) { + std::promise< void > cb_entered_promise; + auto cb_entered = cb_entered_promise.get_future(); + + std::promise< void > allow_cb_exit_promise; + auto allow_cb_exit = allow_cb_exit_promise.get_future(); + + auto chunk_sel = std::make_shared< VolumeChunkSelector >( + "test", [&, this](uint64_t, const std::vector< chunk_num_t >& chunk_ids) { + EXPECT_FALSE(chunk_ids.empty()); + cb_entered_promise.set_value(); + allow_cb_exit.wait(); + }); + + add_chunks_per_pdev(chunk_sel, 1 /*pdevs*/, 10 /*chunks per pdev*/); + const auto total_free_before_alloc = chunk_sel->num_free_chunks(); + + uint32_t pdev_id{UINT32_MAX}; + auto ids = chunk_sel->allocate_init_chunks(0 /* ordinal */, 64 * Ki /* max_num_chunks > 1 */, pdev_id, true); + ASSERT_FALSE(ids.empty()); + + // Exhaust currently active chunks so select_chunk() must try resize + auto active_chunks = chunk_sel->get_chunks(0); + ASSERT_FALSE(active_chunks.empty()); + for (auto& c : active_chunks) { + c->get_internal_chunk()->set_available_blks(0); + } + + homestore::blk_alloc_hints hints; + hints.application_hint = 0; + + std::thread selector([&] { (void)chunk_sel->select_chunk(1 /* nblks */, hints); }); + + ASSERT_EQ(cb_entered.wait_for(std::chrono::seconds(2)), std::future_status::ready); + + std::thread releaser([&] { chunk_sel->release_chunks(0 /* ordinal */); }); + + std::this_thread::sleep_for(std::chrono::milliseconds(100)); + allow_cb_exit_promise.set_value(); + + selector.join(); + releaser.join(); + + wait_until([&] { return chunk_sel->num_free_chunks() == total_free_before_alloc; }, + std::chrono::milliseconds{10000}); + + EXPECT_EQ(chunk_sel->num_free_chunks(), total_free_before_alloc); + EXPECT_TRUE(chunk_sel->get_chunks(0).empty()); +} + +TEST_F(ChunkSelectorTest, RecoverReleaseReallocateTest) { + auto chunk_sel = + std::make_shared< VolumeChunkSelector >("test", [this](uint64_t, const std::vector< chunk_num_t >&) {}); + + add_chunks_per_pdev(chunk_sel, 1 /*pdevs*/, 10 /*chunks per pdev*/); + const auto total_free = chunk_sel->num_free_chunks(); + + std::vector< homestore::chunk_num_t > recovered_ids{0, 1, 2}; + ASSERT_TRUE(chunk_sel->recover_chunks(0 /* ordinal */, 0 /* pdev */, 48 * Ki, recovered_ids)); + + auto recovered_chunks = chunk_sel->get_chunks(0); + ASSERT_EQ(recovered_chunks.size(), recovered_ids.size()); + + chunk_sel->release_chunks(0 /* ordinal */); + + wait_until([&] { return chunk_sel->num_free_chunks() == total_free; }, std::chrono::milliseconds{10000}); + + uint32_t pdev_id{UINT32_MAX}; + auto new_ids = chunk_sel->allocate_init_chunks(1 /* new ordinal */, 48 * Ki, pdev_id, true); + ASSERT_FALSE(new_ids.empty()); + EXPECT_EQ(pdev_id, 0u); +} + int main(int argc, char* argv[]) { int parsed_argc = argc; ::testing::InitGoogleTest(&parsed_argc, argv); diff --git a/src/lib/volume/volume_chunk_selector.cpp b/src/lib/volume/volume_chunk_selector.cpp index 8846b994..21b465f5 100644 --- a/src/lib/volume/volume_chunk_selector.cpp +++ b/src/lib/volume/volume_chunk_selector.cpp @@ -15,6 +15,7 @@ #include "volume_chunk_selector.hpp" #include "hb_internal.hpp" #include +#include namespace homeblocks { @@ -29,17 +30,25 @@ void VolumeChunkSelector::add_chunk(homestore::cshared< Chunk >& chunk) { auto vol_chunk = std::make_shared< HBChunk >(chunk); auto chunk_id = homestore::VChunk(chunk).get_chunk_id(); auto pdev_id = homestore::VChunk(chunk).get_pdev_id(); + + LOGDEBUG("Adding chunk id {} to selector {}", chunk_id, m_module_name); + std::lock_guard lock(m_chunk_sel_mutex); m_all_chunks.emplace(chunk_id, vol_chunk); m_per_dev_chunks[pdev_id].emplace(chunk_id, vol_chunk); - LOGDEBUG("Adding chunk id {} to selector {}", chunk_id, m_module_name); } std::vector< chunk_num_t > VolumeChunkSelector::allocate_init_chunks(uint64_t volume_ordinal, uint64_t volume_size, uint32_t& pdev_id, bool lazy_alloc) { RELEASE_ASSERT(volume_ordinal < m_volume_chunks.size(), "Invalid ordinal for volume {}", volume_ordinal); + + std::unique_lock lock{m_chunk_sel_mutex}; + + // Fast path: volume already exists if (m_volume_chunks[volume_ordinal] != nullptr) { LOGW("Already allocated chunks for volume={}", volume_ordinal); auto volc = m_volume_chunks[volume_ordinal]; + pdev_id = volc->pdev; + std::vector< chunk_num_t > chunk_ids; for (auto& chunk : volc->m_chunks) { if (chunk) { chunk_ids.emplace_back(chunk->get_chunk_id()); } @@ -47,17 +56,13 @@ std::vector< chunk_num_t > VolumeChunkSelector::allocate_init_chunks(uint64_t vo return chunk_ids; } - uint64_t chunk_size{0}; - { - std::lock_guard lock(m_chunk_sel_mutex); - if (m_all_chunks.empty()) { - LOGE("No chunks available in system for volume={}", volume_ordinal); - return {}; - } - - chunk_size = m_all_chunks.begin()->second->size(); + if (m_all_chunks.empty()) { + LOGE("No chunks available in system for volume={}", volume_ordinal); + return {}; } + const uint64_t chunk_size = m_all_chunks.begin()->second->size(); + auto volc = std::make_shared< VolumeChunksInfo >(); volc->ordinal = volume_ordinal; volc->max_num_chunks = std::max(1UL, (volume_size + chunk_size - 1) / chunk_size); @@ -72,20 +77,20 @@ std::vector< chunk_num_t > VolumeChunkSelector::allocate_init_chunks(uint64_t vo // Initially we create num_chunks_per_vol_init active chunks. auto chunks = allocate_init_chunks_from_pdev(volc->num_active_chunks, volc->max_num_chunks); if (chunks.empty()) { - LOGE("Couldnt allocate chunks for volume={}", volume_ordinal); + LOGE("Couldn't allocate chunks for volume={}", volume_ordinal); return {}; } volc->pdev = (*chunks.begin())->get_pdev_id(); pdev_id = volc->pdev; volc->m_chunks.resize(volc->max_num_chunks); - m_volume_chunks[volume_ordinal] = volc; std::string str; uint64_t idx = 0; std::vector< chunk_num_t > chunk_ids; + chunk_ids.reserve(chunks.size()); + for (auto& chunk : chunks) { - // Add the chunks to the volume chunk list. RELEASE_ASSERT(chunk->m_vol_ordinal == INVALID_VOL_ORDINAL, "Chunk assigned to volume {}", chunk->m_vol_ordinal); chunk->m_vol_ordinal = volume_ordinal; @@ -94,6 +99,8 @@ std::vector< chunk_num_t > VolumeChunkSelector::allocate_init_chunks(uint64_t vo fmt::format_to(std::back_inserter(str), "{} ", chunk->get_chunk_id()); } + m_volume_chunks[volume_ordinal] = volc; + LOGI("Allocating initial module={} num_chunks={} for volume={} chunks={}", m_module_name, chunk_ids.size(), volume_ordinal, str); return chunk_ids; @@ -105,69 +112,97 @@ homestore::cshared< Chunk > VolumeChunkSelector::select_chunk(homestore::blk_cou if (!hints.application_hint) { return nullptr; } uint64_t volume_ordinal = hints.application_hint.value(); - // We dont take lock on volumes vector and volume chunks vector - // as they precreated and never changed - if (volume_ordinal >= m_volume_chunks.size() || !m_volume_chunks[volume_ordinal]) { return nullptr; } - auto volc = m_volume_chunks[volume_ordinal]; + shared< VolumeChunksInfo > volc; - // TODO to remove , keep trak of number of freed and alloc blks. - // Add chunk_selector interface to have additional functions on_alloc_blk, on_free_blk in homestore. - uint64_t total_blks = 0, available_blks = 0; - for (auto& chunk : volc->m_chunks) { - if (!chunk) continue; - total_blks += chunk->get_total_blks(); - available_blks += chunk->available_blks(); + { + std::shared_lock lock{m_chunk_sel_mutex}; + + if (volume_ordinal >= m_volume_chunks.size()) { return nullptr; } + + volc = m_volume_chunks[volume_ordinal]; + if (!volc || volc->releasing.load(std::memory_order_acquire)) { return nullptr; } + + volc->inflight_selects.fetch_add(1, std::memory_order_acq_rel); } - do { -#ifdef _PRERELEASE - if (iomgr_flip::instance()->test_flip("vol_num_chunks_force_resize_op")) { - // this is to simulate no blks available. - LOGI("volume resize op flip is set."); - resize_volume_num_chunks(nblks, volc); + // Hand-made RAII: decrement the counter on every exit path + struct PinGuard { + shared< VolumeChunksInfo > volc; + ~PinGuard() { + if (volc) { volc->inflight_selects.fetch_sub(1, std::memory_order_acq_rel); } } -#endif + } pin{volc}; - // If the ratio of available_blks to total_blks is less than half or there is a request of nblks - // more than the available blks and there is room for more chunks then resize. - auto usage_ratio = (float)available_blks / total_blks; - if ((nblks > available_blks || usage_ratio < 0.5) && (volc->num_active_chunks.load() < volc->max_num_chunks)) { - // Check if number chunks needs to be increased. - resize_volume_num_chunks(nblks, volc); - } + // Recheck after pinning in case release started immediately after unlock. + if (volc->releasing.load(std::memory_order_acquire)) { return nullptr; } + + constexpr auto k_sleep = std::chrono::milliseconds{100}; + constexpr auto k_max_wait = std::chrono::seconds{5}; + auto wait_deadline = std::chrono::steady_clock::time_point{}; + + do { + const auto num_active = volc->num_active_chunks.load(std::memory_order_acquire); + if (num_active == 0) { return nullptr; } - // This is the fastpath where we try to allocate the blks from the active chunks. - // Traverse through active chunks in the vector and find the first chunk - // which has some available blks. It may not satisfy all the nblks, in that case - // virtual_dev will call select_chunk again. - uint64_t num_active_chunks = volc->num_active_chunks; - for (uint64_t i = 0; i < num_active_chunks; i++) { - // Lock-free round-robin over the active chunks (shared atomic cursor). - auto const idx = volc->m_next_chunk_index.fetch_add(1, std::memory_order_relaxed) % num_active_chunks; + for (uint64_t i = 0; i < num_active; ++i) { + auto idx = volc->m_next_chunk_index.fetch_add(1, std::memory_order_relaxed) % num_active; auto chunk = volc->m_chunks[idx]; if (chunk && chunk->available_blks() > 0) { return chunk->get_internal_chunk(); } } + if (volc->releasing.load(std::memory_order_acquire)) { return nullptr; } + + if (num_active >= volc->max_num_chunks) { + LOGW("Volume {} is out of space: active={} max={}", volc->ordinal, num_active, volc->max_num_chunks); + return nullptr; + } + + auto rr = resize_volume_num_chunks(nblks, volc); + if (rr == ResizeResult::NoCapacity || rr == ResizeResult::Releasing) { return nullptr; } + + if (rr == ResizeResult::Busy || rr == ResizeResult::Started) { + if (wait_deadline == std::chrono::steady_clock::time_point{}) { + wait_deadline = std::chrono::steady_clock::now() + k_max_wait; + } + + if (std::chrono::steady_clock::now() >= wait_deadline) { + LOGW("Timed out waiting for volume {} resize progress, resize_op={}", volc->ordinal, + static_cast< int >(volc->resize_op.load(std::memory_order_acquire))); + return nullptr; + } + } else { + // reset bounded-wait tracking on any non-waiting outcome + wait_deadline = std::chrono::steady_clock::time_point{}; + } + LOGT("Waiting to allocate more chunks active={} total={}", volc->num_active_chunks.load(), volc->max_num_chunks); LOGT("{}", dump_chunks()); - std::this_thread::sleep_for(std::chrono::milliseconds(100)); + + std::this_thread::sleep_for(k_sleep); + } while (true); return {}; } -void VolumeChunkSelector::resize_volume_num_chunks(homestore::blk_count_t nblks, shared< VolumeChunksInfo > volc) { - auto idle = ResizeOp::Idle, inprogress = ResizeOp::InProgress; - auto status = resize_op.compare_exchange_strong(idle, inprogress); - if (!status) { - // Some other thread is in process of adding the chunks. - return; - } +VolumeChunkSelector::ResizeResult VolumeChunkSelector::resize_volume_num_chunks(homestore::blk_count_t nblks, + shared< VolumeChunksInfo > volc) { + // Don't resize while releasing + if (!volc || volc->releasing.load(std::memory_order_acquire)) { return ResizeResult::Releasing; } + + // Some other thread is in process of adding the chunks + auto idle = ResizeOp::Idle; + auto in_progress = ResizeOp::InProgress; + if (!volc->resize_op.compare_exchange_strong(idle, in_progress)) { return ResizeResult::Busy; } // TODO chunk select will have on_alloc_blk, on_free_blk + // Only scan the published active prefix. Readers should treat + // [0, num_active_chunks) as the only visible chunk range uint64_t total_blks = 0, available_blks = 0; - for (auto& chunk : volc->m_chunks) { + const auto active = volc->num_active_chunks.load(std::memory_order_acquire); + for (uint64_t i = 0; i < active; ++i) { + auto const& chunk = volc->m_chunks[i]; if (!chunk) continue; total_blks += chunk->get_total_blks(); available_blks += chunk->available_blks(); @@ -182,52 +217,124 @@ void VolumeChunkSelector::resize_volume_num_chunks(homestore::blk_count_t nblks, } #endif if (!force_resize) { - auto usage_ratio = (float)available_blks / total_blks; - if ((nblks < available_blks && usage_ratio > 0.5)) { - // Check again if another thread already did the resize. + auto usage_ratio = total_blks ? (float)available_blks / total_blks : 0.0f; + if ((nblks < available_blks && usage_ratio > 0.5f)) { + // Check again if another thread already did the resize LOGI("Another thread already completed the resize op."); - return; + volc->resize_op.store(ResizeOp::Idle, std::memory_order_release); + return ResizeResult::NotNeeded; + } + } + + { + std::shared_lock lock{m_chunk_sel_mutex}; + const auto active_now = volc->num_active_chunks.load(std::memory_order_acquire); + const auto remaining_slots = volc->max_num_chunks - active_now; + const auto free_on_pdev = m_per_dev_chunks.contains(volc->pdev) ? m_per_dev_chunks.at(volc->pdev).size() : 0; + + if (remaining_slots == 0 || free_on_pdev == 0) { + volc->resize_op.store(ResizeOp::Idle, std::memory_order_release); + return ResizeResult::NoCapacity; } } - // Spawn background task to create new chunks. + // Spawn background task to create new chunks LOGD("Initiating op to resize num chunks for module={} volume={} available={} total={}", m_module_name, volc->ordinal, available_blks, total_blks); + iomanager.run_on_forget(iomgr::reactor_regex::random_worker, [volc, this]() mutable { - std::string str; - auto num_chunks_to_alloc = std::min(static_cast< uint64_t >(num_chunks_per_resize), - (volc->max_num_chunks - volc->num_active_chunks.load())); - auto chunks = allocate_resize_chunks_from_pdev(volc->pdev, num_chunks_to_alloc); - RELEASE_ASSERT(!chunks.empty(), "No chunks available for resize volume={}", volc->ordinal) - - // Add the new chunks to the active chunks - auto indx = volc->num_active_chunks.load(); - for (auto chunk : chunks) { - volc->m_chunks[indx] = chunk; - indx++; - fmt::format_to(std::back_inserter(str), "{}({}) ", chunk->get_chunk_id(), chunk->get_pdev_id()); + const auto num_chunks_to_alloc = + std::min(static_cast< uint64_t >(num_chunks_per_resize), + (volc->max_num_chunks - volc->num_active_chunks.load(std::memory_order_acquire))); + + if (num_chunks_to_alloc == 0) { + volc->resize_op.store(ResizeOp::Idle); + return; } + auto new_chunks = allocate_resize_chunks_from_pdev(volc->pdev, num_chunks_to_alloc); + if (new_chunks.empty()) { + LOGW("No chunks available for resize volume={}", volc->ordinal); + volc->resize_op.store(ResizeOp::Idle, std::memory_order_release); + return; + } + + // Build the metadata payload first, but do NOT publish new chunks yet. + // That way select_chunk() still only sees the old active prefix. std::vector< chunk_num_t > chunk_ids; - for (auto& chunk : volc->m_chunks) { - if (chunk) { chunk_ids.emplace_back(chunk->get_chunk_id()); } + bool release_started = false; + { + std::shared_lock lock{m_chunk_sel_mutex}; + + if (volc->releasing.load(std::memory_order_acquire)) { + release_started = true; + } else { + const auto active = volc->num_active_chunks.load(std::memory_order_acquire); + chunk_ids.reserve(active + new_chunks.size()); + + for (uint64_t i = 0; i < active; ++i) { + auto const& chunk = volc->m_chunks[i]; + if (chunk) { chunk_ids.emplace_back(chunk->get_chunk_id()); } + } + } + } + + // If release started after we grabbed chunks from the free pool, + // put them back instead of attaching them to the dying volume + if (release_started) { + std::unique_lock lock{m_chunk_sel_mutex}; + for (auto& chunk : new_chunks) { + m_per_dev_chunks[chunk->get_pdev_id()].emplace(chunk->get_chunk_id(), chunk); + } + volc->resize_op.store(ResizeOp::Idle); + return; } - // Persist the new chunk ids to the metablk of volume - // before making them active chunks. Invoke the registered callback. + for (auto& chunk : new_chunks) { + chunk_ids.emplace_back(chunk->get_chunk_id()); + } + + // Persist first, publish second m_update_vol_sb_cb(volc->ordinal, chunk_ids); - // Update the number of active chunks and compelete the resize operation. - volc->num_active_chunks = indx; - resize_op.store(ResizeOp::Idle); - LOGI("Resize op done. Allocated more chunks for volume={} total={} new={} new_chunks={}", m_module_name, - volc->ordinal, volc->num_active_chunks.load(), chunks.size(), str); + std::string str; + { + std::unique_lock lock{m_chunk_sel_mutex}; + + // Recheck after metadata update. Release may have started while + // callback was running; if so, do not publish these chunks. + if (volc->releasing.load(std::memory_order_acquire)) { + for (auto& chunk : new_chunks) { + m_per_dev_chunks[chunk->get_pdev_id()].emplace(chunk->get_chunk_id(), chunk); + } + volc->resize_op.store(ResizeOp::Idle); + return; + } + + // Publish by writing chunk pointers first, then bumping + // num_active_chunks. Readers only trust the active prefix. + auto idx = volc->num_active_chunks.load(std::memory_order_relaxed); + for (auto& chunk : new_chunks) { + RELEASE_ASSERT(chunk->m_vol_ordinal == INVALID_VOL_ORDINAL, "Chunk assigned to volume {}", + chunk->m_vol_ordinal); + chunk->m_vol_ordinal = volc->ordinal; + volc->m_chunks[idx++] = chunk; + fmt::format_to(std::back_inserter(str), "{}({}) ", chunk->get_chunk_id(), chunk->get_pdev_id()); + } + + volc->num_active_chunks.store(idx, std::memory_order_release); + } + + volc->resize_op.store(ResizeOp::Idle); + LOGI("Resize op done. Allocated more chunks for volume={} total={} new={} new_chunks={}", volc->ordinal, + volc->num_active_chunks.load(), new_chunks.size(), str); }); + + return ResizeResult::Started; } std::vector< shared< VolumeChunkSelector::HBChunk > > VolumeChunkSelector::allocate_init_chunks_from_pdev(uint64_t init_chunks, uint64_t total_chunks) { - std::lock_guard lock(m_chunk_sel_mutex); std::vector< shared< HBChunk > > result; RELEASE_ASSERT(init_chunks <= total_chunks, "Invalid chunks requested"); for (auto& [pdev, pdev_chunks] : m_per_dev_chunks) { @@ -266,10 +373,11 @@ VolumeChunkSelector::allocate_resize_chunks_from_pdev(uint32_t pdev_id, uint64_t bool VolumeChunkSelector::recover_chunks(uint64_t volume_ordinal, uint32_t pdev, uint64_t volume_size, const std::vector< chunk_num_t >& chunk_ids) { - std::lock_guard lock(m_chunk_sel_mutex); + std::unique_lock lock(m_chunk_sel_mutex); auto volc = m_volume_chunks[volume_ordinal]; RELEASE_ASSERT(!volc, "volume already exists"); + if (m_all_chunks.empty()) { return false; } auto chunk_size = m_all_chunks.begin()->second->size(); volc = std::make_shared< VolumeChunksInfo >(); volc->ordinal = volume_ordinal; @@ -306,26 +414,57 @@ bool VolumeChunkSelector::recover_chunks(uint64_t volume_ordinal, uint32_t pdev, return true; } +// Release the active chunks back to the per device chunk pool void VolumeChunkSelector::release_chunks(uint64_t volume_ordinal) { - // Release the active chunks back to the per device chunk pool. - std::lock_guard lock(m_chunk_sel_mutex); - std::string str; - uint64_t count = 0; - auto volc = m_volume_chunks[volume_ordinal]; - RELEASE_ASSERT(volc, "volume doesnt exists"); + shared< VolumeChunksInfo > volc; + + { + std::unique_lock lock(m_chunk_sel_mutex); + volc = std::exchange(m_volume_chunks[volume_ordinal], nullptr); + RELEASE_ASSERT(volc, "volume doesnt exists"); + volc->releasing.store(true, std::memory_order_release); + } + + // Wait until no selector is still walking this volume and no resize is running + while (volc->inflight_selects.load(std::memory_order_acquire) != 0 || + volc->resize_op.load(std::memory_order_acquire) != ResizeOp::Idle) { + std::this_thread::yield(); + } + + auto release_fn = [this, volume_ordinal, volc]() mutable { + std::string str; + uint64_t cnt{}; + + for (auto& chunk : volc->m_chunks) { + if (!chunk) { continue; } - for (auto chunk : volc->m_chunks) { - if (chunk) { - chunk->m_vol_ordinal = INVALID_VOL_ORDINAL; - m_per_dev_chunks[chunk->get_pdev_id()].emplace(chunk->get_chunk_id(), chunk); fmt::format_to(std::back_inserter(str), "{} ", chunk->get_chunk_id()); - count++; + ++cnt; + + if (homestore::hs() && homestore::hs()->has_index_service()) [[likely]] { + homestore::hs()->index_service().wb_cache().evict_chunk_blkids(*chunk->get_internal_chunk()); + } + + { + std::unique_lock lock{m_chunk_sel_mutex}; + chunk->reset(); + m_per_dev_chunks[chunk->get_pdev_id()].emplace(chunk->get_chunk_id(), chunk); + } } - } - m_volume_chunks[volume_ordinal] = nullptr; - LOGI("Released chunks for volume={} num_chunks={}", volume_ordinal, count); - LOGDEBUG("Released chunks={}", str); + volc->m_chunks.clear(); + volc->num_active_chunks.store(0, std::memory_order_release); + + LOGI("Released chunks for volume={} num_chunks={}", volume_ordinal, cnt); + LOGDEBUG("Released chunks={}", str); + }; + + if (homestore::hs() && homestore::hs()->has_index_service()) [[likely]] { + // Clear wbc entries in a separate thread + iomanager.run_on_forget(iomgr::reactor_regex::random_worker, std::move(release_fn)); + } else { + release_fn(); + } } void VolumeChunkSelector::foreach_chunks(std::function< void(homestore::cshared< Chunk >&) >&& cb) { @@ -335,7 +474,7 @@ void VolumeChunkSelector::foreach_chunks(std::function< void(homestore::cshared< } std::vector< shared< VolumeChunkSelector::HBChunk > > VolumeChunkSelector::get_chunks(uint64_t volume_ordinal) { - std::lock_guard lock(m_chunk_sel_mutex); + std::shared_lock lock(m_chunk_sel_mutex); std::vector< shared< VolumeChunkSelector::HBChunk > > chunks; RELEASE_ASSERT(volume_ordinal < m_volume_chunks.size(), "Invalid ordinal for volume {}", volume_ordinal); @@ -348,7 +487,7 @@ std::vector< shared< VolumeChunkSelector::HBChunk > > VolumeChunkSelector::get_c } uint64_t VolumeChunkSelector::num_free_chunks() const { - std::lock_guard lock(m_chunk_sel_mutex); + std::shared_lock lock(m_chunk_sel_mutex); uint64_t count = 0; for (const auto& [pdev, chunks] : m_per_dev_chunks) { count += chunks.size(); @@ -357,7 +496,7 @@ uint64_t VolumeChunkSelector::num_free_chunks() const { } void VolumeChunkSelector::dump_per_pdev_chunks() const { - std::lock_guard lock(m_chunk_sel_mutex); + std::shared_lock lock(m_chunk_sel_mutex); for (const auto& [pdev, chunks] : m_per_dev_chunks) { std::string str; for (const auto& [chunk_id, _] : chunks) { @@ -368,7 +507,7 @@ void VolumeChunkSelector::dump_per_pdev_chunks() const { } std::string VolumeChunkSelector::dump_chunks() const { - std::lock_guard lock(m_chunk_sel_mutex); + std::shared_lock lock(m_chunk_sel_mutex); std::string str; for (uint32_t i = 0; i < m_volume_chunks.size(); i++) { if (!m_volume_chunks[i]) { continue; } diff --git a/src/lib/volume/volume_chunk_selector.hpp b/src/lib/volume/volume_chunk_selector.hpp index 3ad2c3c9..0f377155 100644 --- a/src/lib/volume/volume_chunk_selector.hpp +++ b/src/lib/volume/volume_chunk_selector.hpp @@ -30,12 +30,28 @@ using Chunk = homestore::Chunk; class VolumeChunkSelector : public homestore::ChunkSelector { static constexpr homestore::chunk_num_t num_chunks_per_vol_init = 1; static constexpr homestore::chunk_num_t num_chunks_per_resize = 3; - static constexpr uint64_t INVALID_VOL_ORDINAL = UINT64_MAX; + static constexpr uint64_t INVALID_VOL_ORDINAL = UINT64_MAX; // not owned by any volume struct HBChunk : public homestore::VChunk { HBChunk(homestore::cshared< Chunk >& chunk) : homestore::VChunk(chunk) {} ~HBChunk() = default; uint64_t m_vol_ordinal{INVALID_VOL_ORDINAL}; + + void reset() { + m_vol_ordinal = INVALID_VOL_ORDINAL; + homestore::VChunk::reset(); + } + }; + +private: + enum class ResizeOp { Idle, InProgress }; + + enum class ResizeResult { + Busy, // another thread is resizing + Started, // resize worker launched + NotNeeded, // enough free blocks already + NoCapacity, // cannot add more chunks + Releasing }; struct VolumeChunksInfo { @@ -54,6 +70,10 @@ class VolumeChunkSelector : public homestore::ChunkSelector { std::atomic< uint32_t > m_next_chunk_index{0}; uint64_t ordinal; uint32_t pdev; + + std::atomic< ResizeOp > resize_op{ResizeOp::Idle}; + std::atomic< uint32_t > inflight_selects{}; + std::atomic< bool > releasing{false}; }; public: @@ -88,16 +108,11 @@ class VolumeChunkSelector : public homestore::ChunkSelector { private: std::vector< shared< HBChunk > > allocate_init_chunks_from_pdev(uint64_t init_chunks, uint64_t total_chunks); std::vector< shared< HBChunk > > allocate_resize_chunks_from_pdev(uint32_t pdev, uint64_t num_chunks); - void resize_volume_num_chunks(homestore::blk_count_t nblks, shared< VolumeChunksInfo > volc); + ResizeResult resize_volume_num_chunks(homestore::blk_count_t nblks, shared< VolumeChunksInfo > volc); void dump_per_pdev_chunks() const; std::string dump_chunks() const; private: - enum class ResizeOp { - Idle, - InProgress, - }; - // Store volume chunks details with index as volume ordinal. std::vector< shared< VolumeChunksInfo > > m_volume_chunks; @@ -111,9 +126,8 @@ class VolumeChunkSelector : public homestore::ChunkSelector { // for allocation. This pool is used for allocation of chunks to volume. // Chunks once allocated to volume are removed from this pool. std::unordered_map< uint64_t, ChunkMap > m_per_dev_chunks; - mutable std::mutex m_chunk_sel_mutex; + mutable std::shared_mutex m_chunk_sel_mutex; UpdateVolSbCb m_update_vol_sb_cb; - std::atomic< ResizeOp > resize_op{ResizeOp::Idle}; std::string m_module_name; };