Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion conanfile.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@

class HomestoreConan(ConanFile):
name = "homestore"
version = "7.5.12"
version = "7.5.13"

homepage = "https://github.com/eBay/Homestore"
description = "HomeStore Storage Engine"
Expand Down
15 changes: 12 additions & 3 deletions src/lib/blkdata_svc/blk_read_tracker.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -46,14 +46,23 @@ void BlkReadTracker::merge(const BlkId& blkid, int64_t new_ref_count,
return false;
});
} else if (new_ref_count < 0) {
// This is a remove operation
// This is a remove operation. Extract waiters inside the lambda so they destruct after lambda returns
// and the write lock is released. Firing m_cb() under the write lock can have unexpected issue
// (e.g., destroy cp_guard and trigger cp flush, which deadlocks any iomgr thread that is concurrently
// trying to acquire the same write lock in its IO completion path.
folly::small_vector< blk_track_waiter_ptr, 8 > fired_waiters;
m_pending_reads_map.upsert_or_delete(
base_blkid, [new_ref_count, &base_blkid](BlkTrackRecord& rec, bool existing) {
base_blkid, [new_ref_count, &base_blkid, &fired_waiters](BlkTrackRecord& rec, bool existing) {
HS_DBG_ASSERT_EQ(existing, true, "Decrement a ref count (blk: {}) which does not exist in map",
base_blkid.to_string());
rec.m_ref_cnt += new_ref_count;
return (rec.m_ref_cnt == 0);
if (rec.m_ref_cnt == 0) {
fired_waiters = std::move(rec.m_waiters);
return true;
}
return false;
});
// fired_waiters destructs here, outside the write lock -> m_cb() fires safely
} else {
// this is wait_on operation
m_pending_reads_map.update(base_blkid, [&waiter_rescheduled, &waiter](BlkTrackRecord& rec) {
Expand Down
84 changes: 84 additions & 0 deletions src/tests/test_blk_read_tracker.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -13,10 +13,14 @@
* specific language governing permissions and limitations under the License.
*
*********************************************************************************/
#include <atomic>
#include <chrono>
#include <future>
#include <memory>
#include <mutex>
#include <random>
#include <string>
#include <thread>
#include <vector>
#include <list>
#include <gtest/gtest.h>
Expand Down Expand Up @@ -550,6 +554,86 @@ TEST_F(BlkReadTrackerTest, TestThreadedInsertRemoveAndWait2) {
LOGINFO("Step 4: Test Passed.");
}

// Regression: waiter callback must not run while upsert_or_delete holds the bucket write lock.
//
// When ref_cnt drops to 0, delete n runs under the write lock and ~blk_track_waiter fires m_cb().
// If m_cb() blocks on an external resource (e.g. sync IO) while the external thread needs the same
// write lock, the two threads can deadlock:
//
// remove_thread: remove() and hold write lock in hashmap -> m_cb() -> waits for IO
// io_thread: holds IO -> insert()/remove() -> waits for hashmap write lock
//
// disk_io_mtx stands in for "IO not finished yet". Without the fix, remove_thread holds the
// write lock inside m_cb(); with the fix, m_cb() runs after the lock is released.
//
// Execution timeline (logs use [phase] tags matching this order):
// [setup] insert pending read; register waiter whose m_cb() waits on disk IO
// [io_thread] acquire disk_io_mtx (IO busy); signal ready
// [main] wait until IO is held; start remove_thread
// [remove_thread] remove() -> ref_cnt 0 -> m_cb() fires
// [m_cb] signal started; block on disk_io_mtx
// [io_thread] insert()/remove() on same bucket (needs hashmap write lock)
// [io_thread] release disk_io_mtx
// [m_cb] finish; [remove_thread] return; [main] assert no deadlock
TEST_F(BlkReadTrackerTest, WaiterCallbackRunsOutsideWriteLock) {
// --- [setup] ---
LOGINFO("[setup] insert pending read on a single bucket.");
get_inst()->set_entries_per_record(8);
const BlkId bid{0, 8, 0}; // all ops use the same bucket write lock
get_inst()->insert(bid);

std::mutex disk_io_mtx;
std::promise< void > io_thread_ready;
std::promise< void > callback_started;
std::atomic< bool > callback_done{false};

LOGINFO("[setup] register waiter; m_cb() will block until disk IO completes.");
get_inst()->wait_on(bid, [&]() {
LOGINFO("[m_cb] entered (ref_cnt reached 0).");
callback_started.set_value();
std::lock_guard lock(disk_io_mtx); // remove_thread path: m_cb() waits for IO
callback_done.store(true);
LOGINFO("[m_cb] finished after acquiring disk IO.");
});

// --- [io_thread] + [main] handoff ---
LOGINFO("[main] start io_thread.");
std::thread io_thread([&]() {
std::lock_guard lock(disk_io_mtx); // io_thread holds IO; m_cb() cannot finish yet
io_thread_ready.set_value();
LOGINFO("[io_thread] disk IO held; waiting for m_cb() to start.");

callback_started.get_future().wait();
LOGINFO("[io_thread] m_cb() started; trying insert()/remove() on same bucket.");

// Without fix: remove_thread still holds hashmap write lock -> blocks here -> deadlock.
// With fix: write lock already released -> proceeds while m_cb() waits for IO.
get_inst()->insert(bid);
get_inst()->remove(bid);
LOGINFO("[io_thread] hashmap ops completed.");
});

io_thread_ready.get_future().wait();
LOGINFO("[main] disk IO is held; start remove_thread to trigger m_cb().");

// --- [remove_thread] ---
std::promise< void > remove_done;
std::thread remove_thread([&]() {
get_inst()->remove(bid); // ref_cnt -> 0, fires m_cb()
remove_done.set_value();
LOGINFO("[remove_thread] remove() returned (m_cb() completed).");
});

constexpr auto k_timeout = std::chrono::seconds(10);
ASSERT_EQ(remove_done.get_future().wait_for(k_timeout), std::future_status::ready)
<< "deadlock: m_cb() appears to run under the bucket write lock";

remove_thread.join();
io_thread.join();
ASSERT_TRUE(callback_done.load());
LOGINFO("[main] passed - no deadlock, m_cb() was invoked.");
}

SISL_OPTION_GROUP(test_blk_read_tracker,
(num_threads, "", "num_threads", "number of threads",
::cxxopts::value< uint32_t >()->default_value("2"), "number"));
Expand Down
Loading