use async_trait::async_trait;
use openraft::Raft;
use openraft::ReadPolicy;
use openraft::error::{ClientWriteError, LinearizableReadError, RaftError};
use tsoracle_consensus::ConsensusError;
use crate::host::OpenraftHighWaterHost;
use crate::log_entry::HighWaterCommand;
use crate::state_machine::HighWaterStateMachine;
use crate::type_config::TypeConfig;
pub struct StandaloneHost {
raft: Raft<TypeConfig, HighWaterStateMachine>,
state_machine: HighWaterStateMachine,
}
impl StandaloneHost {
pub fn new(
raft: Raft<TypeConfig, HighWaterStateMachine>,
state_machine: HighWaterStateMachine,
) -> Self {
Self {
raft,
state_machine,
}
}
}
#[async_trait]
impl OpenraftHighWaterHost for StandaloneHost {
type Config = TypeConfig;
type StateMachine = HighWaterStateMachine;
fn raft(&self) -> &Raft<Self::Config, Self::StateMachine> {
&self.raft
}
async fn current_high_water(&self) -> Result<u64, ConsensusError> {
if let Err(e) = self.raft.ensure_linearizable(ReadPolicy::ReadIndex).await {
return match e {
RaftError::APIError(LinearizableReadError::ForwardToLeader(_)) => {
Err(ConsensusError::NotLeader { observed: None })
}
_ => Err(ConsensusError::TransientDriver(Box::new(e))),
};
}
Ok(self.state_machine.current_value().await)
}
async fn submit_advance(&self, at_least: u64) -> Result<u64, ConsensusError> {
match self
.raft
.client_write(HighWaterCommand::Bump { target: at_least })
.await
{
Ok(resp) => Ok(resp.data.value),
Err(RaftError::APIError(ClientWriteError::ForwardToLeader(_))) => {
Err(ConsensusError::NotLeader { observed: None })
}
Err(e) => Err(ConsensusError::TransientDriver(Box::new(e))),
}
}
}