openraft 0.10.0-alpha.35

Advanced Raft consensus
Documentation
use std::collections::BTreeMap;
use std::fmt;

use display_more::DisplayOptionExt;
use display_more::DisplaySliceExt;

use crate::ChangeMembers;
use crate::RaftState;
use crate::RaftTypeConfig;
use crate::base::BoxOnce;
use crate::core::raft_msg::external_command::ExternalCommand;
use crate::core::raft_msg::membership_payloads::MembershipPayloads;
#[cfg(feature = "runtime-stats")]
use crate::core::runtime_stats::RuntimeStats;
use crate::display_ext::DisplayBTreeMapDebugValueExt;
use crate::errors::Infallible;
use crate::errors::InitializeError;
use crate::errors::LinearizableReadError;
use crate::impls::ProgressResponder;
use crate::raft::AppendEntriesRequest;
use crate::raft::ClientWriteResult;
use crate::raft::Precondition;
use crate::raft::VoteRequest;
use crate::raft::VoteResponse;
use crate::raft::linearizable_read::Linearizer;
use crate::raft::linearizable_read::LinearizerOption;
use crate::raft::responder::core_responder::CoreResponder;
use crate::raft::stream_append::StreamAppendResult;
use crate::type_config::alias::BatchOf;
use crate::type_config::alias::CommittedLeaderIdOf;
#[cfg(feature = "runtime-stats")]
use crate::type_config::alias::InstantOf;
use crate::type_config::alias::LogIdOf;
use crate::type_config::alias::OneshotSenderOf;
use crate::type_config::alias::PayloadOf;
use crate::type_config::alias::VoteOf;

pub(crate) mod external_command;
pub(crate) mod install_full_snapshot_request;
pub(crate) mod membership_payloads;
mod raft_msg_name;

pub use raft_msg_name::ExternalCommandName;
pub use raft_msg_name::RaftMsgName;

/// A oneshot TX to send result from `RaftCore` to external caller, e.g. `Raft::append_entries`.
pub(crate) type ResultSender<C, T, E = Infallible> = OneshotSenderOf<C, Result<T, E>>;

/// TX for Vote Response
pub(crate) type VoteTx<C> = OneshotSenderOf<C, VoteResponse<C>>;

/// TX for Append Entries Response
pub(crate) type AppendEntriesTx<C> = OneshotSenderOf<C, StreamAppendResult<C>>;

/// TX for Linearizable Read Response
pub(crate) type ClientReadTx<C> = ResultSender<C, Linearizer<C>, LinearizableReadError<C>>;

/// A message sent by application to the [`RaftCore`].
///
/// [`RaftCore`]: crate::core::RaftCore
pub(crate) enum RaftMsg<C>
where C: RaftTypeConfig
{
    AppendEntries {
        rpc: AppendEntriesRequest<C>,
        tx: AppendEntriesTx<C>,
    },

    RequestVote {
        rpc: VoteRequest<C>,
        tx: VoteTx<C>,
    },

    /// A Pre-Vote request: probe whether a quorum would grant a vote without changing any state.
    RequestPreVote {
        rpc: VoteRequest<C>,
        tx: VoteTx<C>,
    },

    ClientWrite {
        payloads: BatchOf<C, PayloadOf<C>>,
        responders: BatchOf<C, Option<CoreResponder<C>>>,
        expected_leader: Option<CommittedLeaderIdOf<C>>,
        #[cfg(feature = "runtime-stats")]
        proposed_at: InstantOf<C>,
    },

    GetLinearizer {
        linearizer_option: LinearizerOption,
        tx: ClientReadTx<C>,
    },

    Initialize {
        members: BTreeMap<C::NodeId, C::Node>,
        tx: ResultSender<C, (), InitializeError<C>>,
    },

    ChangeMembership {
        changes: ChangeMembers<C::NodeId, C::Node>,

        /// Payloads whose membership OpenRaft replaces with the computed membership.
        ///
        /// `RaftCore` computes the membership, so it also picks the payload that matches the
        /// shape of that membership.
        payloads: MembershipPayloads<C>,

        /// If `retain` is `true`, then the voters that are not in the new
        /// config will be converted into learners, otherwise they will be removed.
        retain: bool,

        preconditions: BatchOf<C, Precondition<C>>,

        tx: ProgressResponder<C, ClientWriteResult<C>>,
    },

    /// Append a caller-built membership as one log entry, with no intermediate joint membership.
    ///
    /// Unlike [`RaftMsg::ChangeMembership`], `RaftCore` does not compute the membership here. It
    /// reads back the membership the caller already bound into `payload`, validates that exact
    /// value, and appends `payload` unchanged.
    AppendMembership {
        /// The caller's payload, after [`RaftPayload::with_membership()`] bound the proposed
        /// membership into it.
        ///
        /// This payload is the single source of truth for the entry: `RaftCore` reads the
        /// membership to validate out of it, and writes this same value to the log. Carrying the
        /// membership in a second field would let the validated value and the written value drift
        /// apart.
        ///
        /// [`RaftPayload::with_membership()`]: crate::entry::RaftPayload::with_membership
        payload: C::Payload,

        preconditions: BatchOf<C, Precondition<C>>,

        tx: ProgressResponder<C, ClientWriteResult<C>>,
    },

    WithRaftState {
        req: BoxOnce<'static, RaftState<C>>,
    },

    /// Transfer Leader to another node.
    ///
    /// If this node is `to`, reset Leader lease and start election.
    /// Otherwise, just reset Leader lease so that the node `to` can become Leader.
    HandleTransferLeader {
        /// The vote of the Leader that is transferring the leadership.
        from: VoteOf<C>,
        /// The assigned node to be the next Leader.
        to: C::NodeId,
        /// The last log id the target must have locally before starting election.
        last_log_id: Option<LogIdOf<C>>,
    },

    ExternalCommand {
        cmd: ExternalCommand<C>,
    },

    /// Get runtime statistics from RaftCore.
    ///
    /// Returns a copy of the current runtime stats.
    #[cfg(feature = "runtime-stats")]
    GetRuntimeStats {
        tx: OneshotSenderOf<C, RuntimeStats<C>>,
    },
}

impl<C> RaftMsg<C>
where C: RaftTypeConfig
{
    /// Returns the name of this message variant.
    pub fn name(&self) -> RaftMsgName {
        match self {
            RaftMsg::AppendEntries { .. } => RaftMsgName::AppendEntries,
            RaftMsg::RequestVote { .. } => RaftMsgName::RequestVote,
            RaftMsg::RequestPreVote { .. } => RaftMsgName::RequestPreVote,
            RaftMsg::ClientWrite { .. } => RaftMsgName::ClientWrite,
            RaftMsg::GetLinearizer { .. } => RaftMsgName::GetLinearizer,
            RaftMsg::Initialize { .. } => RaftMsgName::Initialize,
            RaftMsg::ChangeMembership { .. } => RaftMsgName::ChangeMembership,
            RaftMsg::AppendMembership { .. } => RaftMsgName::AppendMembership,
            RaftMsg::HandleTransferLeader { .. } => RaftMsgName::HandleTransferLeader,
            RaftMsg::WithRaftState { .. } => RaftMsgName::WithRaftState,
            RaftMsg::ExternalCommand { cmd } => RaftMsgName::ExternalCommand(cmd.name()),
            #[cfg(feature = "runtime-stats")]
            RaftMsg::GetRuntimeStats { .. } => RaftMsgName::GetRuntimeStats,
        }
    }
}

impl<C> fmt::Display for RaftMsg<C>
where C: RaftTypeConfig
{
    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
        match self {
            RaftMsg::AppendEntries { rpc, .. } => {
                write!(f, "AppendEntries: {}", rpc)
            }
            RaftMsg::RequestVote { rpc, .. } => {
                write!(f, "RequestVote: {}", rpc)
            }
            RaftMsg::RequestPreVote { rpc, .. } => {
                write!(f, "RequestPreVote: {}", rpc)
            }
            RaftMsg::ClientWrite { .. } => write!(f, "ClientWrite"),
            RaftMsg::GetLinearizer { linearizer_option, .. } => {
                write!(f, "GetLinearizer: {}", linearizer_option)
            }
            RaftMsg::Initialize { members, .. } => {
                write!(f, "Initialize: {}", members.display())
            }
            RaftMsg::ChangeMembership {
                changes,
                retain,
                preconditions,
                ..
            } => {
                write!(
                    f,
                    "ChangeMembership: {}, retain: {}, preconditions: {}",
                    changes,
                    retain,
                    preconditions.as_ref().display()
                )
            }
            RaftMsg::AppendMembership { preconditions, .. } => {
                write!(
                    f,
                    "AppendMembership: preconditions: {}",
                    preconditions.as_ref().display()
                )
            }
            RaftMsg::WithRaftState { .. } => write!(f, "WithRaftState"),
            RaftMsg::HandleTransferLeader { from, to, last_log_id } => {
                write!(
                    f,
                    "TransferLeader: from_leader: vote={}, to: {}, last_log_id: {}",
                    from,
                    to,
                    last_log_id.display()
                )
            }
            RaftMsg::ExternalCommand { cmd } => {
                write!(f, "ExternalCommand: {}", cmd)
            }
            #[cfg(feature = "runtime-stats")]
            RaftMsg::GetRuntimeStats { .. } => {
                write!(f, "GetRuntimeStats")
            }
        }
    }
}

#[cfg(test)]
mod tests {
    use super::*;
    use crate::batch::Batch;
    use crate::engine::testing::UTConfig;
    use crate::engine::testing::log_id;
    use crate::entry::EntryPayload;

    /// `AppendMembership` reports its own name, and prints the preconditions the caller supplied.
    #[test]
    fn test_append_membership_name_and_display() {
        let (tx, _rx) = ProgressResponder::<UTConfig, ClientWriteResult<UTConfig>>::complete_only();

        let precondition = Precondition::LastMembershipLogId {
            last_membership_log_id: Some(log_id(1, 2, 3)),
        };
        let msg = RaftMsg::<UTConfig>::AppendMembership {
            payload: EntryPayload::Blank,
            preconditions: BatchOf::<UTConfig, _>::of([precondition]),
            tx,
        };

        assert_eq!(RaftMsgName::AppendMembership, msg.name());
        assert_eq!(
            "AppendMembership: preconditions: [LastMembershipLogId(T1-N2.3)]",
            msg.to_string()
        );
    }
}