Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
25 commits
Select commit Hold shift + click to select a range
e4e3ad3
add raft for peer communications
raakella1 Jul 18, 2026
551e983
Add wire ops and structures for non raft peer to peer communication
raakella1 Jul 20, 2026
de7d47b
add true login logic in replica. Remove non tcp code from tcp_server
raakella1 Jul 22, 2026
ec64e94
Add the login logic according to the craft design
raakella1 Jul 25, 2026
791f45b
fix test cases
raakella1 Jul 26, 2026
bae9f61
add session fencing login in helo
raakella1 Jul 26, 2026
93a8466
do truncate above rs commit lsn in the internal login commit
raakella1 Jul 29, 2026
6427c79
add python scripts to create replicas and craft disk and perform io …
raakella1 Jul 28, 2026
bfdeefb
add basic fio test
raakella1 Jul 30, 2026
fc9e621
Add more helper scripts for the testing. Bug fixes from running the test
raakella1 Aug 1, 2026
f3d573d
review_comments
raakella1 Aug 5, 2026
71767de
rc2
raakella1 Aug 6, 2026
352963f
remove server dependency on craft_reference library. Add raft replica
raakella1 Aug 12, 2026
3d1ef6e
Add a static registry class to store type erased key value pairs. Thi…
raakella1 Aug 14, 2026
74e185c
remove replica manager and use the registry
raakella1 Aug 14, 2026
971831d
add support torestart replica by moving journal and index into registry
raakella1 Aug 15, 2026
a06aab1
Move replica journal, index volume and peer info to registry
raakella1 Aug 19, 2026
c89ed00
add raft persistence to registry manager
raakella1 Aug 19, 2026
3bd3312
use weak_ptr for registry_mgr to avoid cyclic destruction issue
raakella1 Aug 20, 2026
a795909
Merge dev branch
raakella1 Sep 11, 2026
e720e7d
Fix raft init order and the weak_from_init issue is raft service"
raakella1 Sep 13, 2026
073deec
fix code format
raakella1 Sep 13, 2026
ece6698
persist craft partition state
raakella1 Sep 15, 2026
4f4c9dc
add partition recovery logic. Do not persist state manager in registr…
raakella1 Sep 15, 2026
a09fb79
fix code style
raakella1 Sep 15, 2026
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
16 changes: 4 additions & 12 deletions CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -69,31 +69,23 @@ target_include_directories(craft_reference PRIVATE ${CMAKE_CURRENT_SOURCE_DIR}/s
target_link_libraries(craft_reference PUBLIC craft_client)
target_compile_features(craft_reference PUBLIC cxx_std_23)

# ── craft_replica_mgr: the internal replica manager (testing only) ──
add_library(craft_replica_mgr STATIC
src/replica_mgr.cpp
src/net/tcp_peer.cpp)
target_include_directories(craft_replica_mgr PUBLIC ${CMAKE_CURRENT_SOURCE_DIR}/include)
target_include_directories(craft_replica_mgr PRIVATE ${CMAKE_CURRENT_SOURCE_DIR}/src)
target_link_libraries(craft_replica_mgr PUBLIC craft_client)
target_compile_features(craft_replica_mgr PUBLIC cxx_std_23)

add_library(raft_service STATIC
src/raft/raft_service.cpp
src/raft/in_memory_log_store.cpp
src/raft/raft_state_manager.cpp)
target_include_directories(raft_service PUBLIC ${CMAKE_CURRENT_SOURCE_DIR}/include)
target_include_directories(raft_service PRIVATE ${CMAKE_CURRENT_SOURCE_DIR}/src)
target_link_libraries(raft_service PUBLIC nuraft_mesg::proto craft_replica_mgr)
target_link_libraries(raft_service PUBLIC nuraft_mesg::proto)
target_compile_features(raft_service PUBLIC cxx_std_23)

add_library(raft_tcp_server STATIC
src/raft/raft_replica.cpp
src/watchdog.cpp
src/net/tcp_server.cpp)
src/net/tcp_server.cpp
src/net/tcp_peer.cpp)
target_include_directories(raft_tcp_server PUBLIC ${CMAKE_CURRENT_SOURCE_DIR}/include)
target_include_directories(raft_tcp_server PRIVATE ${CMAKE_CURRENT_SOURCE_DIR}/src)
target_link_libraries(raft_tcp_server PUBLIC craft_reference craft_replica_mgr raft_service)
target_link_libraries(raft_tcp_server PUBLIC craft_reference raft_service)
target_compile_features(raft_tcp_server PUBLIC cxx_std_23)

# ── craft_reference_tcp_srv: a STANDALONE single-replica reference server. Run N of them on different ports to
Expand Down
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 CraftClientConan(ConanFile):
name = "craft_client"
version = "0.4.1"
version = "0.4.2"

description = (
"CRAFT reference client + wire protocol -- transport-agnostic, HomeStore-free"
Expand Down
5 changes: 5 additions & 0 deletions src/helper.hpp
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
#pragma once
#include <boost/uuid/random_generator.hpp>
#include <boost/uuid/uuid_io.hpp>
#include <nlohmann/json.hpp>
#include <stdexec/execution.hpp>
#include <fstream>
Expand Down Expand Up @@ -47,4 +48,8 @@ inline std::error_condition jsonObjectFromFile(std::string const& filename, json
return std::error_condition();
}

inline std::string registry_key(std::string const& prefix, boost::uuids::uuid const& id) {
return fmt::format("{}_{}", prefix, boost::uuids::to_string(id));
}

} // namespace craft
70 changes: 40 additions & 30 deletions src/mem/replica.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,11 @@ std::shared_ptr< std::vector< uint8_t > > take_payload(sisl::sg_list const& s) {
} // namespace

MemCraftReplica::MemCraftReplica(replica_endpoint ep, uint32_t page_size, std::shared_ptr< MemTransport > net) :
ep_{std::move(ep)}, page_size_{page_size}, net_{std::move(net)} {
ep_{std::move(ep)},
page_size_{page_size},
net_{std::move(net)},
journal_{std::make_shared< journal_t >()},
index_{std::make_shared< index_t >()} {
// Publish the initial (healthy) fault snapshot before any IO can read it.
auto initial = std::make_unique< replica_faults const >();
faults_.store(initial.get(), std::memory_order_release);
Expand Down Expand Up @@ -234,7 +238,7 @@ result< lsn_pair > MemCraftReplica::do_write(client_hdr hdr, int64_t dlsn, uint6
}
std::lock_guard< std::mutex > g{mu_};
if (hdr.term != state_.term) return fail(craft_error::STALE_TERM);
if (auto it = journal_.find(dlsn); it != journal_.end() && it->second.is_empty) {
if (auto it = journal_->find(dlsn); it != journal_->end() && it->second.is_empty) {
// An Empty verdict is permanent (reconciliation: Empty beats data). A late arrival into the slot is
// REJECTED -- deterministically -- so that write's own ack path concludes the slot is void, matching
// the verdict instead of phantom-acking a write every replica discarded.
Expand All @@ -247,9 +251,10 @@ result< lsn_pair > MemCraftReplica::do_write(client_hdr hdr, int64_t dlsn, uint6
slot.len = static_cast< lba_count_t >(len / page_size_); // byte length -> block count
slot.all_zeros = !bytes; // no payload => zero write; no all_zeros flag
slot.bytes = std::move(bytes); // adopt the buffer; do not copy it again
journal_[dlsn] = std::move(slot);
(*journal_)[dlsn] = std::move(slot);
state_.last_append_lsn = std::max(state_.last_append_lsn, dlsn);
apply_up_to(hdr.commit_lsn); // piggybacked commit: advance the frontier best-effort, in dLSN order
on_state_changed();
// Piggyback the watermarks on the ack (the wire's write_rsp), so any round-trip refreshes the client.
return lsn_pair{state_.commit_lsn, state_.last_append_lsn};
}
Expand Down Expand Up @@ -286,17 +291,18 @@ result< lsn_pair > MemCraftReplica::do_lsns() {

status MemCraftReplica::do_truncate(int64_t lsn) {
std::lock_guard< std::mutex > g{mu_};
journal_.erase(journal_.upper_bound(lsn), journal_.end());
journal_->erase(journal_->upper_bound(lsn), journal_->end());
state_.last_append_lsn = std::min(state_.last_append_lsn, lsn);
on_state_changed();
return ok();
}

result< std::vector< JournalSlot > > MemCraftReplica::do_fetch(std::vector< int64_t > const& lsns) {
std::lock_guard< std::mutex > g{mu_};
std::vector< JournalSlot > out;
for (auto lsn : lsns) {
auto it = journal_.find(lsn);
if (it == journal_.end()) continue; // not-present-here => omit
auto it = journal_->find(lsn);
if (it == journal_->end()) continue; // not-present-here => omit
auto const& s = it->second;
JournalSlot js;
js.lsn = lsn;
Expand All @@ -322,19 +328,20 @@ result< resolution_result > MemCraftReplica::do_resolve_local(client_hdr hdr, in
if (hdr.term != state_.term) return fail(craft_error::STALE_TERM);
resolution_result out{upto, {}};
for (int64_t d = state_.commit_lsn + 1; d <= upto; ++d) {
auto it = journal_.find(d);
if (it == journal_.end()) {
auto it = journal_->find(d);
if (it == journal_->end()) {
MemJournalSlot s;
s.term = state_.term;
s.is_empty = true;
journal_[d] = std::move(s);
(*journal_)[d] = std::move(s);
out.empty_slots.push_back(d);
} else if (it->second.is_empty) {
out.empty_slots.push_back(d); // a prior verdict; re-report it so the client can retire the slot
}
}
state_.last_append_lsn = std::max(state_.last_append_lsn, upto);
apply_up_to(upto);
on_state_changed();
return out;
}

Expand All @@ -343,20 +350,20 @@ result< resolution_result > MemCraftReplica::do_resolve_local(client_hdr hdr, in
void MemCraftReplica::apply_slot(int64_t dlsn, MemJournalSlot const& s) {
if (s.all_zeros) {
for (lba_count_t i = 0; i < s.len; ++i) {
index_.erase(s.lba + i); // unmap => hole
index_->erase(s.lba + i); // unmap => hole
}
} else {
for (lba_count_t i = 0; i < s.len; ++i) {
index_[s.lba + i] = IndexCell{dlsn, s.bytes, static_cast< std::size_t >(i) * page_size_};
(*index_)[s.lba + i] = IndexCell{dlsn, s.bytes, static_cast< std::size_t >(i) * page_size_};
}
}
}

void MemCraftReplica::apply_up_to(int64_t target) {
int64_t next = state_.commit_lsn + 1;
while (next <= target) {
auto it = journal_.find(next);
if (it == journal_.end()) break; // Missing hole -> stall (best-effort)
auto it = journal_->find(next);
if (it == journal_->end()) break; // Missing hole -> stall (best-effort)
if (!it->second.is_empty) apply_slot(next, it->second);
state_.commit_lsn = next; // Empty slots are skipped on apply but still advance the frontier
++next;
Expand All @@ -366,7 +373,7 @@ void MemCraftReplica::apply_up_to(int64_t target) {
// Highest-dLSN journal-tail slot with commit_lsn < dLSN <= H that covers `x` (the journal-tail overlay,
// materialized on demand). Slots above H are never examined -- that is the horizon clamp.
MemCraftReplica::MemJournalSlot const* MemCraftReplica::highest_slot_le(lba_t x, int64_t H) const {
for (auto it = journal_.upper_bound(H); it != journal_.begin();) {
for (auto it = journal_->upper_bound(H); it != journal_->begin();) {
--it;
if (it->first <= state_.commit_lsn) break; // reached the applied prefix (served from index_)
auto const& s = it->second;
Expand Down Expand Up @@ -416,7 +423,7 @@ std::vector< io_extent > MemCraftReplica::read_range(int64_t H, uint64_t addr, u
page = s->bytes->data() + static_cast< std::size_t >(x - s->lba) * page_size_;
hole = false;
} // else: zero write => hole
} else if (auto it = index_.find(x); it != index_.end()) {
} else if (auto it = index_->find(x); it != index_->end()) {
page = it->second.buf->data() + it->second.off;
hole = false;
}
Expand Down Expand Up @@ -450,14 +457,14 @@ replica_stats MemCraftReplica::stats() const {
s.last_append_lsn = state_.last_append_lsn;
s.term = state_.term;
s.client_token = state_.client_token;
s.mapped_blocks = index_.size();
s.mapped_blocks = index_->size();

s.journal_slots = journal_.size();
if (!journal_.empty()) {
s.journal_first_dlsn = journal_.begin()->first;
s.journal_last_dlsn = journal_.rbegin()->first;
s.journal_slots = journal_->size();
if (!journal_->empty()) {
s.journal_first_dlsn = journal_->begin()->first;
s.journal_last_dlsn = journal_->rbegin()->first;
}
for (auto const& [dlsn, slot] : journal_) {
for (auto const& [dlsn, slot] : (*journal_)) {
if (slot.is_empty) {
++s.empty_slots;
} else if (slot.all_zeros) {
Expand All @@ -473,8 +480,8 @@ replica_stats MemCraftReplica::stats() const {
int64_t const lo = state_.commit_lsn + 1;
int64_t const hi = state_.last_append_lsn;
if (hi >= lo) {
auto const first = journal_.lower_bound(lo);
auto const last = journal_.upper_bound(hi);
auto const first = journal_->lower_bound(lo);
auto const last = journal_->upper_bound(hi);
auto const present = static_cast< std::size_t >(std::distance(first, last));
s.missing_count = static_cast< std::size_t >(hi - lo + 1) - present;

Expand Down Expand Up @@ -508,31 +515,34 @@ void MemCraftReplica::cold_apply_login(uint64_t client_token, uint64_t term) {
std::lock_guard< std::mutex > g{mu_};
state_.client_token = client_token;
state_.term = term;
on_state_changed();
}
void MemCraftReplica::cold_apply_logout() {
std::lock_guard< std::mutex > g{mu_};
state_.client_token = 0;
state_.term = 0; // no active session; subsequent IOs with old term fail STALE_TERM
on_state_changed();
}
void MemCraftReplica::cold_truncate_above(int64_t rs_commit_lsn) {
std::lock_guard< std::mutex > g{mu_};
journal_.erase(journal_.upper_bound(rs_commit_lsn), journal_.end());
journal_->erase(journal_->upper_bound(rs_commit_lsn), journal_->end());
state_.last_append_lsn = std::min(state_.last_append_lsn, rs_commit_lsn);
on_state_changed();
}

// ── resolution-round hooks (driven by MemTransport::run_resolution) ──

std::optional< MemCraftReplica::MemJournalSlot > MemCraftReplica::peek_slot(int64_t dlsn) {
std::lock_guard< std::mutex > g{mu_};
auto const it = journal_.find(dlsn);
if (it == journal_.end()) return std::nullopt;
auto const it = journal_->find(dlsn);
if (it == journal_->end()) return std::nullopt;
return it->second; // copies the slot; `bytes` is shared (immutable once appended), so no payload copy
}

void MemCraftReplica::cold_install_slot(int64_t dlsn, MemJournalSlot s) {
std::lock_guard< std::mutex > g{mu_};
if (journal_.contains(dlsn)) return; // already holds it (or a verdict); a fetch never overwrites
journal_[dlsn] = std::move(s);
if (journal_->contains(dlsn)) return; // already holds it (or a verdict); a fetch never overwrites
(*journal_)[dlsn] = std::move(s);
state_.last_append_lsn = std::max(state_.last_append_lsn, dlsn);
}

Expand All @@ -543,14 +553,14 @@ void MemCraftReplica::cold_mark_empty(int64_t dlsn) {
MemJournalSlot s;
s.term = state_.term;
s.is_empty = true;
journal_[dlsn] = std::move(s);
(*journal_)[dlsn] = std::move(s);
state_.last_append_lsn = std::max(state_.last_append_lsn, dlsn);
}

std::vector< int64_t > MemCraftReplica::peek_empties(int64_t upto) {
std::lock_guard< std::mutex > g{mu_};
std::vector< int64_t > out;
for (auto const& [d, s] : journal_) {
for (auto const& [d, s] : (*journal_)) {
if (d > upto) break;
if (s.is_empty) out.push_back(d); // ascending: journal_ is an ordered map
}
Expand Down
9 changes: 7 additions & 2 deletions src/mem/replica.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -225,6 +225,9 @@ class MemCraftReplica : public craft_replica,
std::size_t off{0};
};

using journal_t = std::map< int64_t, MemJournalSlot >;
using index_t = std::map< lba_t, IndexCell >;

// Synchronous cores: the SERVER. Each takes mu_. Deliverability, latency and payload ownership are the
// transport's job (MemTransport::send_*), which is why nothing below consults net_ or copies bytes.
// do_write takes the payload already owned and adopts it; `bytes == nullptr` is a zero write.
Expand Down Expand Up @@ -272,6 +275,7 @@ class MemCraftReplica : public craft_replica,
std::vector< int64_t > peek_empties(int64_t upto); // every is_empty dLSN <= upto

protected:
virtual void on_state_changed() {} // called under mu_ after every state_ mutation; override to persist
void cold_apply_login(uint64_t client_token, uint64_t term);
void cold_truncate_above(int64_t rs_commit_lsn);
void cold_install_slot(int64_t dlsn, MemJournalSlot s); // fill a hole; never overwrites an entry
Expand All @@ -295,8 +299,9 @@ class MemCraftReplica : public craft_replica,
std::shared_ptr< MemTransport > net_;

CraftPartitionState state_;
std::map< int64_t, MemJournalSlot > journal_; // dLSN -> slot (out-of-order arrival tolerated)
std::map< lba_t, IndexCell > index_; // applied prefix (<= commit_lsn); an absent LBA is a hole
// shared_ptr for restart support using registry_manager.
std::shared_ptr< journal_t > journal_; // dLSN -> slot (out-of-order arrival tolerated)
std::shared_ptr< index_t > index_; // applied prefix (<= commit_lsn); an absent LBA is a hole
mutable std::mutex mu_;
std::shared_ptr< Watchdog > watchdog_; // the keepalive and login watchdog
};
Expand Down
38 changes: 20 additions & 18 deletions src/net/tcp_server.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -21,13 +21,13 @@

#include <algorithm>
#include <utility>
#include <boost/uuid/uuid_io.hpp>

#include <sisl/logging/logging.h> // server-side r/w trace (base module; visible with -v trace / when a consumer inits logging)

#include "raft/raft_replica.hpp" // the full RaftReplica (+ sisl::sg_list via sisl/fds/buffer.hpp)
#include <craft/status.hpp> // to_wire_status (the shared wire <-> craft_error bridge)
#include "raft/raft_service.hpp"
#include "replica_mgr.hpp"
#include "registry_mgr.hpp"
#include <craft/status.hpp> // to_wire_status (the shared wire <-> craft_error bridge)
#include "helper.hpp"

namespace craft::net {
Expand All @@ -39,20 +39,22 @@ std::span< uint8_t const > as_bytes(T const& v) {
}
} // namespace

craft_tcp_server::craft_tcp_server(server_geometry geo, std::string const& server_config_file) : geo_{std::move(geo)} {
craft_tcp_server::craft_tcp_server(server_geometry geo, std::string const& server_config_file,
std::shared_ptr< registry_manager > registry_mgr, bool init_raft_service) :
geo_{std::move(geo)},
registry_mgr_{registry_mgr ? std::move(registry_mgr) : std::make_shared< registry_manager >()},
raft_enabled_{init_raft_service} {
auto ep = replica_endpoint{.id = to_uuid(geo_.member.id), .addr = geo_.member.addr};
LOGINFO("craft_tcp_server: starting [id={}] config_file='{}'", boost::uuids::to_string(ep.id), server_config_file);
// net == nullptr: this replica serves exclusively through its srv_* seam (the TCP frontend IS the wire).
// start replica service and raft service if server_config_file is provided
if (!server_config_file.empty()) {
replica_manager::instance()->start_replica_service(server_config_file, ep.id);
raft_service::instance()->start_raft_service(ep.id);
LOGINFO("craft_tcp_server: replica_manager + raft_service started [id={}]", boost::uuids::to_string(ep.id));
} else {
LOGINFO("craft_tcp_server: no server_config_file given -- running in standalone/cold-path mode [id={}]",
boost::uuids::to_string(ep.id));
}
replica_ = std::make_shared< RaftReplica >(std::move(ep), geo_.lba_size, geo_.max_tx);
replica_ = std::make_shared< RaftReplica >(raft_replica_params{
.ep = std::move(ep),
.page_size = geo_.lba_size,
.max_tx = geo_.max_tx,
.replica_config_path = server_config_file,
.watchdog = std::make_shared< Watchdog >(),
.registry_mgr = registry_mgr_,
.init_raft_service = init_raft_service,
});
}

craft_tcp_server::~craft_tcp_server() = default;
Expand Down Expand Up @@ -177,7 +179,7 @@ void craft_tcp_server::on_helo(craft_conn& conn, wire::message const& req) {
auto const hr = wire::decode< wire::helo_req >(req.op_header);
wire::status code = wire::status::ok;

if (bool is_raft_enabled = raft_service::instance()->is_raft_enabled(); !is_raft_enabled) {
if (!raft_enabled_) {
// no raft, follow fake cold path
auto result = replica_->srv_establish(hr.volume_id, hr.client_token, session_term_);
if (!result) {
Expand Down Expand Up @@ -363,8 +365,8 @@ void craft_tcp_server::on_create_volume(craft_conn& conn, wire::message const& r
replica_members.emplace_back(replica_endpoint{.id = craft::to_uuid(m.id), .addr = m.addr});
}

if (auto const r = replica_->srv_create_volume(cr.volume_id, replica_members); !r) {
LOGERROR("craft_srv CREATE_VOLUME [rid:{}]: srv_create_volume failed: {}", req.hdr.request_id,
if (auto const r = replica_->srv_create_partition(cr.volume_id, replica_members); !r) {
LOGERROR("craft_srv CREATE_VOLUME [rid:{}]: srv_create_partition failed: {}", req.hdr.request_id,
r.error().message());
code = to_wire_status(r.error());
}
Expand Down
Loading
Loading