Skip to content
Merged
Show file tree
Hide file tree
Changes from 4 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 @@ -10,7 +10,7 @@

class HomeBlocksConan(ConanFile):
name = "homeblocks"
version = "6.0.9"
version = "6.0.10"

homepage = "https://github.com/eBay/HomeBlocks"
description = "Block Store built on HomeStore"
Expand Down
56 changes: 26 additions & 30 deletions src/lib/craft/craft_repl_dev.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -365,25 +365,17 @@ unique< CraftJournalBackend > make_homestore_journal_backend(shared< homestore::
return std::make_unique< HomeStoreCraftJournalBackend >(std::move(logstore), vol_ordinal, lba_size);
}

// ─── HomeStoreCraftCheckpointTrigger (SDSTOR-22888) ──────────────────────────
//
// Thin wrapper over homestore::cp_mgr(). One instance is shared by every volume's CraftReplDev.

class HomeStoreCraftCheckpointTrigger : public CraftCheckpointTrigger {
public:
async_status trigger_cp_flush(bool force) override {
checkpoint_trigger_fn_t make_homestore_checkpoint_trigger_fn() {
return [](bool force) -> async_status {
Comment thread
raakella1 marked this conversation as resolved.
if (co_await homestore::cp_mgr().trigger_cp_flush(force)) co_return ok();
// cp_mgr().trigger_cp_flush() returns false, synchronously, when a flush is already in
// progress (cp_mgr.cpp: m_in_flush_phase). That's expected and harmless when force=false
// (this trigger's only current caller, fire_checkpoint_trigger, coalesces with any in-flight
// flush on purpose). Only a force=true false return is a real failure worth surfacing.
if (!force) co_return ok();
co_return std::unexpected(make_error_condition(volume_error::INTERNAL_ERROR));
}
};

unique< CraftCheckpointTrigger > make_homestore_checkpoint_trigger() {
return std::make_unique< HomeStoreCraftCheckpointTrigger >();
};
}

// ─── constructor ──────────────────────────────────────────────────────────────
Expand Down Expand Up @@ -787,18 +779,18 @@ bool CraftReplDev::checkpoint_interval_crossed_locked(int64_t commit_lsn_snapsho
}

void CraftReplDev::fire_checkpoint_trigger(int64_t commit_lsn_snapshot) {
if (checkpoint_trigger_ == nullptr) {
if (!checkpoint_trigger_) {
LOGW("commit_lsn={} crossed checkpoint interval but no checkpoint_trigger_ wired -- skipping",
commit_lsn_snapshot);
return;
}
// force=false: let this coalesce with any checkpoint already in flight rather than forcing
// back-to-back flushes under high commit throughput (see CraftCheckpointTrigger's doc comment).
// back-to-back flushes under high commit throughput (see checkpoint_trigger_fn_t's doc comment).
// Detached (fire-and-forget): nothing here depends on the flush completing. A failure is logged,
// not propagated, same posture as catch-up/fetch failures elsewhere in this class.
auto self = shared_from_this();
detail::detach([self, commit_lsn_snapshot]() -> async_status {
if (auto cp = co_await self->checkpoint_trigger_->trigger_cp_flush(false); !cp)
if (auto cp = co_await self->checkpoint_trigger_(false); !cp)
LOGE("checkpoint trigger failed at commit_lsn={}: {}", commit_lsn_snapshot, cp.error().message());
co_return ok();
}());
Expand All @@ -812,8 +804,8 @@ void CraftReplDev::fire_checkpoint_trigger(int64_t commit_lsn_snapshot) {
// every subsequent write()/keep_alive() retries the advance. Never holds a lock across the co_await
// read_slot() suspension point below (same rule fetch_data's doc comment already establishes).

async_result< int64_t > CraftReplDev::commit_impl(int64_t upto_lsn, write_index_fn_t const& write_fn,
delete_index_fn_t const& delete_fn) {
async_result< int64_t > CraftReplDev::commit_impl(int64_t upto_lsn, write_index_fn_t write_fn,
delete_index_fn_t delete_fn) {
int64_t commit_lsn, last_append_lsn;
{
std::lock_guard lk{state_mu_};
Expand All @@ -823,12 +815,23 @@ async_result< int64_t > CraftReplDev::commit_impl(int64_t upto_lsn, write_index_
last_append_lsn = state_.last_append_lsn;
}
// Guaranteed reset on every exit path (stall, success, or error): a local RAII object's destructor
// runs when the coroutine frame unwinds, exactly like a plain function's locals on return.
// runs when the coroutine frame unwinds, exactly like a plain function's locals on return. Also
// runs the checkpoint-interval check here rather than only on the happy-path tail: folding it into
// this same critical section means every exit -- including an early error return mid-loop -- still
// gets a chance to fire the checkpoint (a slot that keeps failing on retry no longer permanently
// starves it), and merges what would otherwise be two separate state_mu_ acquisitions into one.
struct RunningGuard {
CraftReplDev* self;
~RunningGuard() {
std::lock_guard lk{self->state_mu_};
self->commit_running_ = false;
int64_t final_commit_lsn;
bool should_checkpoint;
{
std::lock_guard lk{self->state_mu_};
self->commit_running_ = false;
final_commit_lsn = self->state_.commit_lsn;
should_checkpoint = self->checkpoint_interval_crossed_locked(final_commit_lsn);
}
if (should_checkpoint) self->fire_checkpoint_trigger(final_commit_lsn);
}
} guard{this};

Expand Down Expand Up @@ -946,15 +949,8 @@ async_result< int64_t > CraftReplDev::commit_impl(int64_t upto_lsn, write_index_
state_.commit_lsn = lsn;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

you can maintain a local variable final_commmit_lsn which you can define above the struct RunningGuard and have RunningGuard hold a ptr to it.

int64_t final_commmit_lsn;
struct RunningGuard {
    CraftReplDev* self;
    int64_t* final_commmit_lsn_ptr;
    ~RunningGuard() {
            bool should_checkpoint;
            {
                std::lock_guard lk{self->state_mu_};
                self->state_.commit_lsn = *final_commmit_lsn_ptr;
                self->commit_running_ = false;
                should_checkpoint = self->checkpoint_interval_crossed_locked(final_commit_lsn);
            }
            if (should_checkpoint) self->fire_checkpoint_trigger(final_commit_lsn);
            }
    } guard(this, &final_commit_lsn);

Update that variable without holding any lock at the end of each loop (line 749)
You can also get rid of the lock before return and use return final_commmit_lsn

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

ACK

}

int64_t final_commit_lsn;
bool should_checkpoint;
{
std::lock_guard lk{state_mu_};
final_commit_lsn = state_.commit_lsn;
should_checkpoint = checkpoint_interval_crossed_locked(final_commit_lsn);
}
if (should_checkpoint) fire_checkpoint_trigger(final_commit_lsn);
co_return final_commit_lsn;
std::lock_guard lk{state_mu_};
co_return state_.commit_lsn;
}

async_result< int64_t > CraftReplDev::commit(int64_t upto_lsn) {
Expand All @@ -969,14 +965,14 @@ async_result< int64_t > CraftReplDev::commit(int64_t upto_lsn) {
delete_index_fn_t delete_fn = [this](lba_t s, lba_t e, std::vector< homestore::blk_id >& freed) {
return indx_tbl_->delete_lba_range(s, e, freed);
};
auto r = co_await commit_impl(upto_lsn, write_fn, delete_fn);
auto r = co_await commit_impl(upto_lsn, std::move(write_fn), std::move(delete_fn));
co_return r;
}

#ifdef _PRERELEASE
async_result< int64_t > CraftReplDev::commit_with(int64_t upto_lsn, write_index_fn_t write_fn,
delete_index_fn_t delete_fn) {
auto r = co_await commit_impl(upto_lsn, write_fn, delete_fn);
auto r = co_await commit_impl(upto_lsn, std::move(write_fn), std::move(delete_fn));
co_return r;
}
#endif
Expand Down
48 changes: 20 additions & 28 deletions src/lib/craft/craft_repl_dev.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -153,26 +153,19 @@ class CraftPeerFetcher {
virtual ~CraftPeerFetcher() = default;
};

// ─── CraftCheckpointTrigger ───────────────────────────────────────────────────
// ─── checkpoint_trigger_fn_t ──────────────────────────────────────────────────
//
// Abstraction over HomeStore's checkpoint manager (homestore::cp_mgr().trigger_cp_flush()).
// Injected into CraftReplDev so unit tests (which compile craft_repl_dev.cpp directly against a
// mock journal backend, with no running HomeStore instance -- see test_craft_raft_entries.cpp) can
// exercise the trigger without touching HomeStore. Production code passes
// HomeStoreCraftCheckpointTrigger (defined in craft_repl_dev.cpp). Default (null) leaves the
// trigger stubbed -- same posture as CraftPeerFetcher.

class CraftCheckpointTrigger {
public:
virtual async_status trigger_cp_flush(bool force) = 0;
virtual ~CraftCheckpointTrigger() = default;
};

// Factory that wraps homestore::cp_mgr(). One instance is shared by every volume's CraftReplDev
// (there is exactly one CPManager per HomeStore instance), unlike make_homestore_journal_backend
// which is per-volume -- so CraftReplDev takes this via a non-owning pointer (set_checkpoint_trigger),
// not ownership at construction. Tests inject MockCraftCheckpointTrigger directly.
unique< CraftCheckpointTrigger > make_homestore_checkpoint_trigger();
// A callable trigger_cp_flush(bool force) -- abstraction over HomeStore's checkpoint manager
// (homestore::cp_mgr().trigger_cp_flush()). Injected into CraftReplDev so unit tests can exercise
// the trigger without touching HomeStore. Production code passes make_homestore_checkpoint_trigger_fn()'s result.
// Default (empty) leaves the trigger stubbed -- same posture as CraftPeerFetcher.
using checkpoint_trigger_fn_t = std::function< async_status(bool force) >;

// Factory that wraps homestore::cp_mgr() (there is exactly one CPManager per HomeStore instance,
// unlike make_homestore_journal_backend which is per-volume). Each CraftReplDev stores its own copy
// of the returned callable by value (set_checkpoint_trigger) -- callers don't need to keep the
// factory's result alive themselves. Tests inject a plain callable directly.
checkpoint_trigger_fn_t make_homestore_checkpoint_trigger_fn();

// ─── CraftReplDev ─────────────────────────────────────────────────────────────
//
Expand Down Expand Up @@ -289,7 +282,7 @@ class CraftReplDev : public std::enable_shared_from_this< CraftReplDev > {
// would let a later read serve stale data from a write that no longer exists in the journal). Called
// only during login (quiesced -- no concurrent writes). commit_lsn is NOT changed.
// FIXME(S4/S7): before dropping entries here, force a completed checkpoint --
// co_await checkpoint_trigger_->trigger_cp_flush(true). Once entries above/below lsn are gone, the journal is
// co_await checkpoint_trigger_(true). Once entries above/below lsn are gone, the journal is
// no longer a durable record of them; if HomeStore's checkpoint has only been requested and not yet
// completed, a crash in between loses that data. HomeStore's own IndexTable::destroy() hits the
// identical problem and force-flushes before removing its superblock for exactly this reason.
Expand Down Expand Up @@ -362,9 +355,9 @@ class CraftReplDev : public std::enable_shared_from_this< CraftReplDev > {
void set_peer_fetch_timeout_ms(uint32_t ms) { peer_fetch_timeout_ms_ = ms; }

// Wires the HomeStore checkpoint trigger used by apply_sync_rs_commit_lsn's periodic checkpoint
// (SDSTOR-22888). One CraftCheckpointTrigger instance is shared by every volume's CraftReplDev;
// tests inject a mock.
void set_checkpoint_trigger(CraftCheckpointTrigger* t) { checkpoint_trigger_ = t; }
// (SDSTOR-22888). One checkpoint_trigger_fn_t is shared by every volume's CraftReplDev; tests
// inject a mock.
void set_checkpoint_trigger(checkpoint_trigger_fn_t fn) { checkpoint_trigger_ = std::move(fn); }

// Overrides the commit_lsn delta between checkpoint triggers (default matches
// sync_rs_commit_lsn_interval's own default of 128, tying checkpoint cadence to the periodic
Expand Down Expand Up @@ -509,8 +502,7 @@ class CraftReplDev : public std::enable_shared_from_this< CraftReplDev > {
// operations (write_index_fn_t / delete_index_fn_t, declared at the top of this class) so tests
// can exercise it against a fake index instead of a real VolumeIndexTable. commit() binds these
// to indx_tbl_'s real methods; commit_with() (test-only) binds test doubles.
async_result< int64_t > commit_impl(int64_t upto_lsn, write_index_fn_t const& write_fn,
delete_index_fn_t const& delete_fn);
async_result< int64_t > commit_impl(int64_t upto_lsn, write_index_fn_t write_fn, delete_index_fn_t delete_fn);

// Must be called with state_mu_ held. Returns true if commit_lsn_snapshot has crossed
// checkpoint_lsn_interval_ since last_checkpoint_lsn_ -- and if so, updates last_checkpoint_lsn_ to
Expand Down Expand Up @@ -666,9 +658,9 @@ class CraftReplDev : public std::enable_shared_from_this< CraftReplDev > {
#ifdef _PRERELEASE
std::atomic< int > watchdog_fire_count_{0}; // test-only; never present in production binaries
#endif
CraftCheckpointTrigger* checkpoint_trigger_{nullptr}; // null until production wiring; unit tests inject a mock
int64_t checkpoint_lsn_interval_{128}; // commit_lsn delta between checkpoint triggers; see
// set_checkpoint_lsn_interval()
checkpoint_trigger_fn_t checkpoint_trigger_{}; // empty until production wiring; unit tests inject a mock
int64_t checkpoint_lsn_interval_{128}; // commit_lsn delta between checkpoint triggers; see
// set_checkpoint_lsn_interval()
int64_t last_checkpoint_lsn_{-1}; // commit_lsn as of the last triggered checkpoint (guarded by state_mu_)
// FIXME(S8/SDSTOR-22745): defaults to -1 in lockstep with state_.commit_lsn
// When S8 wires recovering commit_lsn from the journal/superblock on restart,
Expand Down
40 changes: 40 additions & 0 deletions src/lib/craft/tests/test_craft_commit.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -496,6 +496,46 @@ TEST_F(CraftCommitTest, WriteFnErrorAbortsCommit) {
EXPECT_EQ(dev_->commit_lsn(), -1); // nothing applied
}

// ── checkpoint trigger on error exit (PR #182 review) ─────────────────────────

struct MockCraftCheckpointTrigger {
int call_count{0};
async_status operator()(bool) {
++call_count;
co_return ok();
}
};

// The checkpoint-interval check now runs from RunningGuard's destructor, so it must fire even when
// a later slot in the same commit_impl run fails -- not just on the happy-path tail. Slots 0-2
// succeed (crossing a checkpoint_lsn_interval_ of 2), then slot 3's write_fn fails; the checkpoint
// trigger must still fire once, at the partially-advanced commit_lsn (2).
TEST_F(CraftCommitTest, CheckpointTriggerFiresOnErrorExit) {
dev_->seed_lsns(3, {});
journal_->add_data_slot(0, 0, 1, 100, {11});
journal_->add_data_slot(1, 1, 1, 101, {12});
journal_->add_data_slot(2, 2, 1, 102, {13});
journal_->add_data_slot(3, 3, 1, 103, {14});

MockCraftCheckpointTrigger trigger;
dev_->set_checkpoint_trigger(std::ref(trigger));
dev_->set_checkpoint_lsn_interval(2);

auto r = homeblocks::detail::sync_get(dev_->commit_with(
3,
[this](lba_t s, lba_t e, std::unordered_map< lba_t, BlockInfo >& info) -> status {
if (s == 3) return std::unexpected(volume_error::INDEX_ERROR); // fail on the 4th slot
return index_.write_to_index(s, e, info);
},
[this](lba_t s, lba_t e, std::vector< homestore::blk_id >& freed) {
return index_.delete_lba_range(s, e, freed);
}));

ASSERT_FALSE(r.has_value());
EXPECT_EQ(dev_->commit_lsn(), 2); // lsns 0-2 applied; lsn 3 aborted the run
EXPECT_EQ(trigger.call_count, 1); // checkpoint still fired despite the error exit
}

// commit() clamps to last_append_lsn even if asked to commit further than what has actually been
// appended locally (e.g. the client's own view of commit_lsn is ahead of this replica).
TEST_F(CraftCommitTest, ClampsToLastAppendLsn) {
Expand Down
22 changes: 10 additions & 12 deletions src/lib/craft/tests/test_craft_homestore_backend.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@
// backend directly rather than through CraftReplDev or a volume -- the narrowest test that still
// runs the real completion path.
//
// Also exercises HomeStoreCraftCheckpointTrigger::trigger_cp_flush (SDSTOR-22888) against the REAL
// Also exercises make_homestore_checkpoint_trigger_fn's result (SDSTOR-22888) against the REAL
// homestore::cp_mgr() -- same rationale: MockCraftCheckpointTrigger (test_craft_raft_entries.cpp)
// covers CraftReplDev's own gating logic, but the wrapper's factory -> cp_mgr().trigger_cp_flush()
// -> async_status conversion chain had never been compiled and run against a live CPManager.
Expand Down Expand Up @@ -387,38 +387,36 @@ TEST_F(CraftHomeStoreBackendTest, FreeSlotRejectsCorruptEntry) {

// force=false: the value apply_sync_rs_commit_lsn's periodic trigger actually passes today.
TEST_F(CraftHomeStoreBackendTest, CheckpointTriggerFlushesRealCPManager) {
auto trigger = make_homestore_checkpoint_trigger();
ASSERT_TRUE(trigger != nullptr);
auto trigger = make_homestore_checkpoint_trigger_fn();

auto r = homeblocks::detail::sync_get(trigger->trigger_cp_flush(/* force = */ false));
auto r = homeblocks::detail::sync_get(trigger(/* force = */ false));
ASSERT_TRUE(r.has_value());
}

// force=true: untested until now -- this is the value truncate()'s FIXME (craft_repl_dev.hpp) says
// a future correctness-critical call site will need, but the passthrough itself had never been
// exercised against the real cp_mgr() for either value.
TEST_F(CraftHomeStoreBackendTest, CheckpointTriggerHonorsForceFlag) {
auto trigger = make_homestore_checkpoint_trigger();
ASSERT_TRUE(trigger != nullptr);
auto trigger = make_homestore_checkpoint_trigger_fn();

auto r = homeblocks::detail::sync_get(trigger->trigger_cp_flush(/* force = */ true));
auto r = homeblocks::detail::sync_get(trigger(/* force = */ true));
ASSERT_TRUE(r.has_value());
}

// cp_mgr()'s in-flight-flush gate (m_in_flush_phase) is a synchronous check-and-set at entry: the
// first of several concurrent trigger_cp_flush(false) calls holds it for the duration of the real
// flush, so every other concurrent call observes it already set and gets HomeStore's synchronous
// false back -- not a failure. HomeStoreCraftCheckpointTrigger must map that to ok() when
// force=false; every one of these concurrent calls must succeed, not just the one that actually flushed.
// false back -- not a failure. make_homestore_checkpoint_trigger_fn's result must map that to ok()
// when force=false; every one of these concurrent calls must succeed, not just the one that
// actually flushed.
TEST_F(CraftHomeStoreBackendTest, CheckpointTriggerConcurrentForceFalseNeverFails) {
auto trigger = make_homestore_checkpoint_trigger();
ASSERT_TRUE(trigger != nullptr);
auto trigger = make_homestore_checkpoint_trigger_fn();

constexpr int k_concurrent = 20;
std::vector< homestore::async_status > futs; // async_status is itself a sisl::async::task<result<...>>
futs.reserve(k_concurrent);
for (int i = 0; i < k_concurrent; ++i) {
futs.push_back(trigger->trigger_cp_flush(/* force = */ false));
futs.push_back(trigger(/* force = */ false));
}

auto const results = homeblocks::detail::sync_get(sisl::async::when_all(std::move(futs)));
Expand Down
Loading
Loading