use std::sync::Arc;
use actix::{
dev::ToEnvelope,
prelude::*,
};
use futures::sync::{mpsc::UnboundedReceiver, oneshot::Sender};
use serde::{Serialize, Deserialize};
use crate::{
AppData, AppDataResponse, AppError, NodeId,
messages,
};
pub struct GetInitialState<E: AppError> {
marker: std::marker::PhantomData<E>,
}
impl<E: AppError> GetInitialState<E> {
pub fn new() -> Self {
Self{marker: std::marker::PhantomData}
}
}
impl<E: AppError> Message for GetInitialState<E> {
type Result = Result<InitialState, E>;
}
#[derive(Clone, Debug)]
pub struct InitialState {
pub last_log_index: u64,
pub last_log_term: u64,
pub last_applied_log: u64,
pub hard_state: HardState,
}
pub struct GetLogEntries<D: AppData, E: AppError> {
pub start: u64,
pub stop: u64,
marker_data: std::marker::PhantomData<D>,
marker_error: std::marker::PhantomData<E>,
}
impl<D: AppData, E: AppError> GetLogEntries<D, E> {
pub fn new(start: u64, stop: u64) -> Self {
Self{start, stop, marker_data: std::marker::PhantomData, marker_error: std::marker::PhantomData}
}
}
impl<D: AppData, E: AppError> Message for GetLogEntries<D, E> {
type Result = Result<Vec<messages::Entry<D>>, E>;
}
pub struct AppendEntryToLog<D: AppData, E: AppError> {
pub entry: Arc<messages::Entry<D>>,
marker: std::marker::PhantomData<E>,
}
impl<D: AppData, E: AppError> AppendEntryToLog<D, E> {
pub fn new(entry: Arc<messages::Entry<D>>) -> Self {
Self{entry, marker: std::marker::PhantomData}
}
}
impl<D: AppData, E: AppError> Message for AppendEntryToLog<D, E> {
type Result = Result<(), E>;
}
pub struct ReplicateToLog<D: AppData, E: AppError> {
pub entries: Arc<Vec<messages::Entry<D>>>,
marker: std::marker::PhantomData<E>,
}
impl<D: AppData, E: AppError> ReplicateToLog<D, E> {
pub fn new(entries: Arc<Vec<messages::Entry<D>>>) -> Self {
Self{entries, marker: std::marker::PhantomData}
}
}
impl<D: AppData, E: AppError> Message for ReplicateToLog<D, E> {
type Result = Result<(), E>;
}
pub struct ApplyEntryToStateMachine<D: AppData, R: AppDataResponse, E: AppError> {
pub payload: Arc<messages::Entry<D>>,
marker0: std::marker::PhantomData<R>,
marker1: std::marker::PhantomData<E>,
}
impl<D: AppData, R: AppDataResponse, E: AppError> ApplyEntryToStateMachine<D, R, E> {
pub fn new(payload: Arc<messages::Entry<D>>) -> Self {
Self{payload, marker0: std::marker::PhantomData, marker1: std::marker::PhantomData}
}
}
impl<D: AppData, R: AppDataResponse, E: AppError> Message for ApplyEntryToStateMachine<D, R, E> {
type Result = Result<R, E>;
}
pub struct ReplicateToStateMachine<D: AppData, E: AppError> {
pub payload: Vec<messages::Entry<D>>,
marker: std::marker::PhantomData<E>,
}
impl<D: AppData, E: AppError> ReplicateToStateMachine<D, E> {
pub fn new(payload: Vec<messages::Entry<D>>) -> Self {
Self{payload, marker: std::marker::PhantomData}
}
}
impl<D: AppData, E: AppError> Message for ReplicateToStateMachine<D, E> {
type Result = Result<(), E>;
}
pub struct CreateSnapshot<E: AppError> {
pub through: u64,
marker: std::marker::PhantomData<E>,
}
impl<E: AppError> CreateSnapshot<E> {
pub fn new(through: u64) -> Self {
Self{through, marker: std::marker::PhantomData}
}
}
impl<E: AppError> Message for CreateSnapshot<E> {
type Result = Result<CurrentSnapshotData, E>;
}
pub struct InstallSnapshot<E: AppError> {
pub term: u64,
pub index: u64,
pub stream: UnboundedReceiver<InstallSnapshotChunk>,
marker: std::marker::PhantomData<E>,
}
impl<E: AppError> InstallSnapshot<E> {
pub fn new(term: u64, index: u64, stream: UnboundedReceiver<InstallSnapshotChunk>) -> Self {
Self{term, index, stream, marker: std::marker::PhantomData}
}
}
impl<E: AppError> Message for InstallSnapshot<E> {
type Result = Result<(), E>;
}
pub struct InstallSnapshotChunk {
pub offset: u64,
pub data: Vec<u8>,
pub done: bool,
pub cb: Sender<()>,
}
pub struct GetCurrentSnapshot<E: AppError> {
marker: std::marker::PhantomData<E>,
}
impl<E: AppError> GetCurrentSnapshot<E> {
pub fn new() -> Self {
Self{marker: std::marker::PhantomData}
}
}
impl<E: AppError> Message for GetCurrentSnapshot<E> {
type Result = Result<Option<CurrentSnapshotData>, E>;
}
#[derive(Clone, Debug, PartialEq)]
pub struct CurrentSnapshotData {
pub term: u64,
pub index: u64,
pub membership: messages::MembershipConfig,
pub pointer: messages::EntrySnapshotPointer,
}
pub struct SaveHardState<E: AppError>{
pub hs: HardState,
marker: std::marker::PhantomData<E>,
}
impl<E: AppError> SaveHardState<E> {
pub fn new(hs: HardState) -> Self {
Self{hs, marker: std::marker::PhantomData}
}
}
impl<E: AppError> Message for SaveHardState<E> {
type Result = Result<(), E>;
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct HardState {
pub current_term: u64,
pub voted_for: Option<NodeId>,
pub membership: messages::MembershipConfig,
}
pub trait RaftStorage<D, R, E>: 'static
where
D: AppData,
R: AppDataResponse,
E: AppError,
{
type Actor: Actor<Context=Self::Context> +
Handler<GetInitialState<E>> +
Handler<SaveHardState<E>> +
Handler<GetLogEntries<D, E>> +
Handler<AppendEntryToLog<D, E>> +
Handler<ReplicateToLog<D, E>> +
Handler<ApplyEntryToStateMachine<D, R, E>> +
Handler<ReplicateToStateMachine<D, E>> +
Handler<CreateSnapshot<E>> +
Handler<InstallSnapshot<E>> +
Handler<GetCurrentSnapshot<E>>;
type Context: ActorContext +
ToEnvelope<Self::Actor, GetInitialState<E>> +
ToEnvelope<Self::Actor, SaveHardState<E>> +
ToEnvelope<Self::Actor, GetLogEntries<D, E>> +
ToEnvelope<Self::Actor, AppendEntryToLog<D, E>> +
ToEnvelope<Self::Actor, ReplicateToLog<D, E>> +
ToEnvelope<Self::Actor, ApplyEntryToStateMachine<D, R, E>> +
ToEnvelope<Self::Actor, ReplicateToStateMachine<D, E>> +
ToEnvelope<Self::Actor, CreateSnapshot<E>> +
ToEnvelope<Self::Actor, InstallSnapshot<E>> +
ToEnvelope<Self::Actor, GetCurrentSnapshot<E>>;
}