pub struct MultiRaft { /* private fields */ }Expand description
Multi-Raft coordinator managing multiple Raft groups on a single node.
This coordinator:
- Manages all Raft groups hosted on this node
- Batches heartbeats across groups sharing the same leader
- Routes incoming RPCs to the correct group
- Collects
Readyoutput from all groups for the caller to execute
Implementations§
Source§impl MultiRaft
impl MultiRaft
Sourcepub fn propose_conf_change(
&mut self,
group_id: u64,
change: &ConfChange,
) -> Result<(u64, u64)>
pub fn propose_conf_change( &mut self, group_id: u64, change: &ConfChange, ) -> Result<(u64, u64)>
Propose a configuration change to a Raft group.
The change is serialized into the group’s Raft log as a
regular entry with a distinguishing prefix byte. It
replicates through the normal AppendEntries path and is
applied by every follower replica when the entry commits
(see apply_conf_change).
§Single-voter vs. multi-voter groups
Single-voter groups commit inside node.propose itself
(see nodedb_raft::node::RaftNode::propose single-voter
branch). In that case the commit has already happened by
the time we return, so we safely apply the change inline:
any caller that reads routing immediately after the
propose sees the final state.
Multi-voter groups commit asynchronously once enough
followers have replicated the entry. The apply then
happens on the tick loop after it observes the updated
commit_index. We MUST NOT inline-apply in that case —
if the leader steps down before replication completes, a
new leader may truncate the log entry and the local state
would be permanently ahead of the committed state with no
rollback path. Callers that need to wait for the apply
should poll the routing table (see
raft_loop::join::wait_for_routing_contains_learner).
Returns (group_id, log_index) on success.
Sourcepub fn apply_conf_change(
&mut self,
group_id: u64,
change: &ConfChange,
) -> Result<()>
pub fn apply_conf_change( &mut self, group_id: u64, change: &ConfChange, ) -> Result<()>
Apply a committed configuration change to this node’s view of the given Raft group.
This is called from the tick loop for every committed entry
detected as a conf-change (via ConfChange::from_entry_data). It
must be idempotent with respect to no-op changes so replaying the
log after a crash does not double-apply.
Source§impl MultiRaft
impl MultiRaft
Sourcepub fn new(node_id: u64, routing: RoutingTable, data_dir: PathBuf) -> Self
pub fn new(node_id: u64, routing: RoutingTable, data_dir: PathBuf) -> Self
Construct a MultiRaft owning its routing table by value.
Wraps the table in a fresh Arc<RwLock<_>>. Used by tests that do not
need to share the routing handle with a ClusterState. Production
construction sites use MultiRaft::new_with_shared_routing so the
data plane and Raft state machine read/write the SAME table.
Construct a MultiRaft sharing the given routing handle.
The passed Arc<RwLock<RoutingTable>> MUST be the same handle stored
in ClusterState.routing so committed conf-changes converge the
data-plane routing view.
Sourcepub fn with_election_timeout(self, min: Duration, max: Duration) -> Self
pub fn with_election_timeout(self, min: Duration, max: Duration) -> Self
Configure election timeout range.
Sourcepub fn with_heartbeat_interval(self, interval: Duration) -> Self
pub fn with_heartbeat_interval(self, interval: Duration) -> Self
Configure heartbeat interval.
Sourcepub fn with_log_compaction_threshold(self, threshold: Option<u64>) -> Self
pub fn with_log_compaction_threshold(self, threshold: Option<u64>) -> Self
Configure the auto-compaction threshold for every group created on
this node. None disables auto-compaction (the default). See
RaftConfig::log_compaction_threshold.
Sourcepub fn add_group(&mut self, group_id: u64, peers: Vec<u64>) -> Result<()>
pub fn add_group(&mut self, group_id: u64, peers: Vec<u64>) -> Result<()>
Initialize a Raft group on this node as a voting member.
peers is the list of other voters in the group (excluding self).
For a learner-start group, use add_group_as_learner instead.
Sourcepub fn add_group_as_learner(
&mut self,
group_id: u64,
voters: Vec<u64>,
learners: Vec<u64>,
) -> Result<()>
pub fn add_group_as_learner( &mut self, group_id: u64, voters: Vec<u64>, learners: Vec<u64>, ) -> Result<()>
Initialize a Raft group on this node as a non-voting learner.
The local node boots in the Learner role and will not stand for
election until it is promoted by a PromoteLearner conf change.
voters is the full voter set of the group (excluding self).
learners is the learner set of the group excluding self — usually
empty unless multiple learners are being admitted in the same round.
Sourcepub fn tick(&mut self) -> Result<MultiRaftReady>
pub fn tick(&mut self) -> Result<MultiRaftReady>
Tick all Raft groups. Returns aggregated ready output.
Any HardState staged by a tick (an election term bump + self-vote from
an election timeout) is durably persisted BEFORE the aggregated Ready
— and therefore the vote requests it carries — is returned for
dispatch. A persist failure aborts the tick so the caller never sends
vote requests for a term that was not made durable.
Sourcepub fn routing(&self) -> Arc<RwLock<RoutingTable>> ⓘ
pub fn routing(&self) -> Arc<RwLock<RoutingTable>> ⓘ
Clone of the shared routing handle.
Returns an Arc clone pointing at the same RwLock<RoutingTable> the
data plane reads. Callers that need a RoutingTable value take a tight
read guard and clone it out.
pub fn node_id(&self) -> u64
Sourcepub fn in_flight_snapshots(&self) -> Arc<InFlightSnapshots> ⓘ
pub fn in_flight_snapshots(&self) -> Arc<InFlightSnapshots> ⓘ
Clone of the in-flight InstallSnapshot tracker.
The tick loop clones this to mark snapshot transfers active for their
lifetime; maybe_compact_group reads it to defer compaction while a
transfer is in flight.
pub fn group_count(&self) -> usize
Sourcepub fn contains_group(&self, group_id: u64) -> bool
pub fn contains_group(&self, group_id: u64) -> bool
Whether this node hosts the given Raft group.
Sourcepub fn group_ids(&self) -> Vec<u64>
pub fn group_ids(&self) -> Vec<u64>
IDs of every Raft group hosted on this node, including groups that do not own vShards (for example the Calvin sequencer).
Sourcepub fn group_membership(&self, group_id: u64) -> Option<GroupMembership>
pub fn group_membership(&self, group_id: u64) -> Option<GroupMembership>
Snapshot the actual Raft membership rather than the vShard routing view.
Sourcepub fn groups_mut(&mut self) -> &mut HashMap<u64, RaftNode<RedbLogStorage>>
pub fn groups_mut(&mut self) -> &mut HashMap<u64, RaftNode<RedbLogStorage>>
Mutable access to the underlying Raft groups (for testing / bootstrap).
Sourcepub fn group_statuses(&self) -> Vec<GroupStatus>
pub fn group_statuses(&self) -> Vec<GroupStatus>
Snapshot of all Raft group states for observability.
Sourcepub fn leader_for_vshard(&self, vshard_id: u32) -> Result<Option<u64>>
pub fn leader_for_vshard(&self, vshard_id: u32) -> Result<Option<u64>>
Get the leader for a given vShard (from local group state).
Sourcepub fn vshard_role_is_leader(&self, vshard_id: u32) -> bool
pub fn vshard_role_is_leader(&self, vshard_id: u32) -> bool
Whether THIS node is currently the leader of the data-group that owns
vshard_id.
Maps the vshard to its Raft group via the routing table and reuses the
existing local leader-role check — no new election. Returns false when
the vshard has no group mapping or this node is a follower/learner for
the owning group. Used by the Calvin scheduler to stamp the per-node,
non-replicated is_group_leader dispatch flag so the OLLP optimistic-lock
verification runs only on the leader while every replica applies the same
predicted write-set (determinism).
Sourcepub fn propose(&mut self, vshard_id: u32, data: Vec<u8>) -> Result<(u64, u64)>
pub fn propose(&mut self, vshard_id: u32, data: Vec<u8>) -> Result<(u64, u64)>
Propose a command to the Raft group that owns the given vShard.
Returns (group_id, log_index) on success.
Sourcepub fn is_group_leader(&self, group_id: u64) -> bool
pub fn is_group_leader(&self, group_id: u64) -> bool
Returns true if this node is currently the leader of group_id.
Returns false when the group does not exist on this node or when the
node is a follower, candidate, or learner in the group.
Sourcepub fn propose_to_group(&mut self, group_id: u64, data: Vec<u8>) -> Result<u64>
pub fn propose_to_group(&mut self, group_id: u64, data: Vec<u8>) -> Result<u64>
Propose a command directly to a specific Raft group (e.g. the metadata group, which has no vShard mapping).
Returns the committed log index on success.
Sourcepub fn read_committed_entries(
&self,
group_id: u64,
lo: u64,
hi: u64,
) -> Result<Vec<LogEntry>>
pub fn read_committed_entries( &self, group_id: u64, lo: u64, hi: u64, ) -> Result<Vec<LogEntry>>
Read committed log entries for a Raft group in the inclusive index
range [lo, hi].
hi is clamped to the group’s commit_index so callers that pass
u64::MAX never read uncommitted entries.
Used by the Calvin scheduler’s rebuild path to replay sequenced transactions from the sequencer Raft log after a restart.
Returns Err(ClusterError::Raft(RaftError::LogCompacted)) if lo
has been compacted into a snapshot (caller must install a snapshot
instead of replaying from log).
Sourcepub fn first_available_index(&self, group_id: u64) -> Option<u64>
pub fn first_available_index(&self, group_id: u64) -> Option<u64>
The lowest committed index still available in group_id’s retained log
(snapshot_index + 1), or None when the group is absent on this node.
Used to arm a Calvin scheduler catch-up from the earliest replayable sequencer index so its drain reads exactly the retained log and never faults on a compacted range.
Sourcepub fn maybe_compact_group(
&mut self,
group_id: u64,
applied_index: u64,
) -> Result<bool>
pub fn maybe_compact_group( &mut self, group_id: u64, applied_index: u64, ) -> Result<bool>
Auto-compact a group’s log if its configured threshold has been
reached, given the DATA-PLANE applied watermark applied_index.
applied_index MUST be the index the data-plane state machine has
durably applied to (NOT raft’s commit index). Compacting past an
unapplied index would let the SnapshotBuilder serialize
incomplete state and corrupt a lagging follower’s snapshot.
No-op (returns Ok(false)) when the group is absent on this node,
the threshold is None, or the retained-entry count is below the
threshold. Returns Ok(true) when a compaction was performed.
Source§impl MultiRaft
impl MultiRaft
Sourcepub fn group_contains_node(&self, group_id: u64, node_id: u64) -> Option<bool>
pub fn group_contains_node(&self, group_id: u64, node_id: u64) -> Option<bool>
Whether a node is already admitted as a voter or learner.
Sourcepub fn commit_index_for(&self, group_id: u64) -> Option<u64>
pub fn commit_index_for(&self, group_id: u64) -> Option<u64>
Current commit index for a group, or None if the group is not
hosted on this node.
Sourcepub fn ready_learners(&self, group_id: u64) -> Vec<u64>
pub fn ready_learners(&self, group_id: u64) -> Vec<u64>
Learners in group_id whose match_index on this leader has
caught up to the current commit_index — safe to promote.
Returns an empty vec if this node is not the leader of the group or the group is not hosted here.
Sourcepub fn group_leader(&self, group_id: u64) -> u64
pub fn group_leader(&self, group_id: u64) -> u64
Observed leader id for a group (0 = unknown / no election yet).
Sourcepub fn group_role_is_leader(&self, group_id: u64) -> bool
pub fn group_role_is_leader(&self, group_id: u64) -> bool
Whether this node is currently the leader of group_id.
Sourcepub fn transfer_leadership(&mut self, group_id: u64, target: u64) -> Result<()>
pub fn transfer_leadership(&mut self, group_id: u64, target: u64) -> Result<()>
Initiate a leadership transfer for group_id to target.
Delegates to RaftNode::transfer_leadership. Returns
ClusterError::GroupNotFound if the group is not hosted on this node.
The outbound TimeoutNow trigger is emitted into the group’s Ready
output and dispatched by the next tick.
Source§impl MultiRaft
impl MultiRaft
Sourcepub fn handle_append_entries(
&mut self,
req: &AppendEntriesRequest,
) -> Result<AppendEntriesResponse>
pub fn handle_append_entries( &mut self, req: &AppendEntriesRequest, ) -> Result<AppendEntriesResponse>
Route an AppendEntries RPC to the correct group.
Sourcepub fn handle_request_vote(
&mut self,
req: &RequestVoteRequest,
) -> Result<RequestVoteResponse>
pub fn handle_request_vote( &mut self, req: &RequestVoteRequest, ) -> Result<RequestVoteResponse>
Route a RequestVote RPC to the correct group.
Sourcepub fn handle_install_snapshot(
&mut self,
req: &InstallSnapshotRequest,
) -> Result<InstallSnapshotResponse>
pub fn handle_install_snapshot( &mut self, req: &InstallSnapshotRequest, ) -> Result<InstallSnapshotResponse>
Route an InstallSnapshot RPC to the correct group.
Sourcepub fn handle_timeout_now(&mut self, req: &TimeoutNowRequest)
pub fn handle_timeout_now(&mut self, req: &TimeoutNowRequest)
Route a TimeoutNow RPC to the correct group.
One-way — no response is produced. Silently ignored if the group is
not mounted on this node (mirrors handle_request_vote for absent
groups). The term+leader_id guard inside RaftNode::handle_timeout_now
remains in place as an additional correctness check.
Sourcepub fn persist_group_hard_state(&mut self, group_id: u64) -> Result<()>
pub fn persist_group_hard_state(&mut self, group_id: u64) -> Result<()>
Durably persist a group’s HardState (current_term/voted_for) if it
changed since the last persist. Must run under the MultiRaft lock
before an RPC reply that granted a vote or bumped the term leaves this
node, so a restart cannot forget the vote and let two leaders form.
No-op when the group is not mounted on this node.
Sourcepub fn snapshot_metadata(&self, group_id: u64) -> Result<(u64, u64, u64)>
pub fn snapshot_metadata(&self, group_id: u64) -> Result<(u64, u64, u64)>
Get the current term and snapshot metadata for a group (for building InstallSnapshot RPCs).
Sourcepub fn handle_append_entries_response(
&mut self,
group_id: u64,
peer: u64,
resp: &AppendEntriesResponse,
) -> Result<()>
pub fn handle_append_entries_response( &mut self, group_id: u64, peer: u64, resp: &AppendEntriesResponse, ) -> Result<()>
Handle AppendEntries response for a specific group.
Sourcepub fn handle_request_vote_response(
&mut self,
group_id: u64,
peer: u64,
resp: &RequestVoteResponse,
) -> Result<()>
pub fn handle_request_vote_response( &mut self, group_id: u64, peer: u64, resp: &RequestVoteResponse, ) -> Result<()>
Handle RequestVote response for a specific group.
Sourcepub fn advance_applied(&mut self, group_id: u64, applied_to: u64) -> Result<()>
pub fn advance_applied(&mut self, group_id: u64, applied_to: u64) -> Result<()>
Advance applied index for a group after processing committed entries.
This is the DELIVERY watermark. See Self::save_applied_index for the
durable floor a restart resumes from.
Sourcepub fn save_applied_index(
&mut self,
group_id: u64,
applied_to: u64,
) -> Result<()>
pub fn save_applied_index( &mut self, group_id: u64, applied_to: u64, ) -> Result<()>
Durably record applied_to as the group’s applied floor.
applied_to MUST name an entry whose state-machine effects are already
durable — for data groups, one whose redo record the WAL has fsynced.
The next boot resumes delivery at applied_to + 1, so this is what
keeps WAL replay and Raft replay from applying the same entry twice.
Monotonic per group: an index at or below the current floor is a no-op.
Sourcepub fn match_index_for(&self, group_id: u64, peer: u64) -> Option<u64>
pub fn match_index_for(&self, group_id: u64, peer: u64) -> Option<u64>
Query a peer’s match_index from a specific Raft group’s leader state.
Sourcepub fn last_applied(&self, group_id: u64) -> Option<u64>
pub fn last_applied(&self, group_id: u64) -> Option<u64>
Read the locally-applied index for a Raft group hosted on this
node. Returns None if the group is not mounted here.
Used by the tick loop to mirror last_applied into the
per-group crate::applied_watcher::AppliedIndexWatcher —
covers both the regular apply path and the snapshot-install
path (which sets last_applied = last_included_index
directly without producing committed entries).
Sourcepub fn applied_indices(&self) -> Vec<(u64, u64)>
pub fn applied_indices(&self) -> Vec<(u64, u64)>
(group_id, last_applied) pairs for every locally-mounted
group. Cheap O(groups) snapshot — groups are few (one
metadata + handful of vshard groups per node).
Auto Trait Implementations§
impl !RefUnwindSafe for MultiRaft
impl !UnwindSafe for MultiRaft
impl Freeze for MultiRaft
impl Send for MultiRaft
impl Sync for MultiRaft
impl Unpin for MultiRaft
impl UnsafeUnpin for MultiRaft
Blanket Implementations§
Source§impl<T> ArchivePointee for T
impl<T> ArchivePointee for T
Source§type ArchivedMetadata = ()
type ArchivedMetadata = ()
Source§fn pointer_metadata(
_: &<T as ArchivePointee>::ArchivedMetadata,
) -> <T as Pointee>::Metadata
fn pointer_metadata( _: &<T as ArchivePointee>::ArchivedMetadata, ) -> <T as Pointee>::Metadata
Source§impl<'a, T, E> AsTaggedExplicit<'a, E> for Twhere
T: 'a,
impl<'a, T, E> AsTaggedExplicit<'a, E> for Twhere
T: 'a,
Source§impl<'a, T, E> AsTaggedExplicit<'a, E> for Twhere
T: 'a,
impl<'a, T, E> AsTaggedExplicit<'a, E> for Twhere
T: 'a,
Source§impl<'a, T, E> AsTaggedImplicit<'a, E> for Twhere
T: 'a,
impl<'a, T, E> AsTaggedImplicit<'a, E> for Twhere
T: 'a,
Source§impl<'a, T, E> AsTaggedImplicit<'a, E> for Twhere
T: 'a,
impl<'a, T, E> AsTaggedImplicit<'a, E> for Twhere
T: 'a,
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
Source§fn in_current_span(self) -> Instrumented<Self>
fn in_current_span(self) -> Instrumented<Self>
Source§impl<T> LayoutRaw for T
impl<T> LayoutRaw for T
Source§fn layout_raw(_: <T as Pointee>::Metadata) -> Result<Layout, LayoutError>
fn layout_raw(_: <T as Pointee>::Metadata) -> Result<Layout, LayoutError>
Source§impl<T, N1, N2> Niching<NichedOption<T, N1>> for N2
impl<T, N1, N2> Niching<NichedOption<T, N1>> for N2
Source§unsafe fn is_niched(niched: *const NichedOption<T, N1>) -> bool
unsafe fn is_niched(niched: *const NichedOption<T, N1>) -> bool
Source§fn resolve_niched(out: Place<NichedOption<T, N1>>)
fn resolve_niched(out: Place<NichedOption<T, N1>>)
out indicating that a T is niched.impl<T> Read<Exclusive, BecauseExclusive> for Twhere
T: ?Sized,
Source§impl<SS, SP> SupersetOf<SS> for SPwhere
SS: SubsetOf<SP>,
impl<SS, SP> SupersetOf<SS> for SPwhere
SS: SubsetOf<SP>,
Source§fn to_subset(&self) -> Option<SS>
fn to_subset(&self) -> Option<SS>
self from the equivalent element of its
superset. Read moreSource§fn is_in_subset(&self) -> bool
fn is_in_subset(&self) -> bool
self is actually part of its subset T (and can be converted to it).Source§fn to_subset_unchecked(&self) -> SS
fn to_subset_unchecked(&self) -> SS
self.to_subset but without any property checks. Always succeeds.Source§fn from_subset(element: &SS) -> SP
fn from_subset(element: &SS) -> SP
self to the equivalent element of its superset.