use std::pin::Pin;
use std::sync::Arc;
use async_trait::async_trait;
use futures::{Stream, StreamExt};
use openraft::RaftTypeConfig;
use tsoracle_consensus::{ConsensusDriver, ConsensusError, LeaderState};
use tsoracle_core::Epoch;
use tsoracle_openraft_toolkit::LeadershipState;
use tsoracle_openraft_toolkit::lifecycle::leader::stream_from_receiver;
use crate::host::OpenraftHighWaterHost;
pub struct OpenraftDriver<H: OpenraftHighWaterHost> {
host: Arc<H>,
}
impl<H: OpenraftHighWaterHost> OpenraftDriver<H> {
pub fn new(host: H) -> Arc<Self> {
Arc::new(Self {
host: Arc::new(host),
})
}
pub fn from_arc(host: Arc<H>) -> Arc<Self> {
Arc::new(Self { host })
}
}
#[async_trait]
impl<H: OpenraftHighWaterHost> ConsensusDriver for OpenraftDriver<H> {
fn leadership_events(&self) -> Pin<Box<dyn Stream<Item = LeaderState> + Send>> {
let host = Arc::clone(&self.host);
Box::pin(owned_leadership_stream::<H>(host))
}
async fn load_high_water(&self) -> Result<u64, ConsensusError> {
self.host.current_high_water().await
}
async fn persist_high_water(
&self,
at_least: u64,
_epoch: Epoch,
) -> Result<u64, ConsensusError> {
self.host.submit_advance(at_least).await
}
}
fn owned_leadership_stream<H: OpenraftHighWaterHost>(
host: Arc<H>,
) -> impl Stream<Item = LeaderState> + Send + 'static {
let rx = host.raft().metrics();
let inner: Pin<Box<dyn Stream<Item = LeaderState> + Send>> =
Box::pin(stream_from_receiver::<H::Config>(rx).map(map_leader_state::<H::Config>));
KeepAlive { _host: host, inner }
}
struct KeepAlive<H: OpenraftHighWaterHost> {
_host: Arc<H>,
inner: Pin<Box<dyn Stream<Item = LeaderState> + Send>>,
}
impl<H: OpenraftHighWaterHost> Stream for KeepAlive<H> {
type Item = LeaderState;
fn poll_next(
mut self: Pin<&mut Self>,
cx: &mut std::task::Context<'_>,
) -> std::task::Poll<Option<Self::Item>> {
self.inner.as_mut().poll_next(cx)
}
}
fn map_leader_state<C: RaftTypeConfig>(s: LeadershipState<C>) -> LeaderState {
match s {
LeadershipState::Leader { term } => LeaderState::Leader {
epoch: Epoch(u128::from(term)),
},
LeadershipState::Follower { .. } => LeaderState::Follower {
leader_endpoint: None,
},
LeadershipState::Candidate { .. }
| LeadershipState::Learner
| LeadershipState::Shutdown => LeaderState::Unknown,
}
}
#[cfg(test)]
mod tests {
use super::map_leader_state;
use crate::type_config::TypeConfig;
use tsoracle_consensus::LeaderState;
use tsoracle_core::Epoch;
use tsoracle_openraft_toolkit::LeadershipState;
#[test]
fn leader_maps_to_leader_with_epoch() {
let s = map_leader_state::<TypeConfig>(LeadershipState::Leader { term: 7 });
assert_eq!(s, LeaderState::Leader { epoch: Epoch(7) });
}
#[test]
fn follower_with_no_leader_maps_to_follower() {
let s = map_leader_state::<TypeConfig>(LeadershipState::Follower {
term: 3,
leader: None,
});
assert_eq!(
s,
LeaderState::Follower {
leader_endpoint: None
}
);
}
#[test]
fn follower_with_known_leader_still_maps_without_endpoint() {
let s = map_leader_state::<TypeConfig>(LeadershipState::Follower {
term: 4,
leader: Some((
2u64,
crate::type_config::OpenraftPeer {
addr: "ignored".into(),
},
)),
});
assert_eq!(
s,
LeaderState::Follower {
leader_endpoint: None
}
);
}
#[test]
fn candidate_maps_to_unknown() {
let s = map_leader_state::<TypeConfig>(LeadershipState::Candidate { term: 5 });
assert_eq!(s, LeaderState::Unknown);
}
#[test]
fn learner_maps_to_unknown() {
let s = map_leader_state::<TypeConfig>(LeadershipState::Learner);
assert_eq!(s, LeaderState::Unknown);
}
#[test]
fn shutdown_maps_to_unknown() {
let s = map_leader_state::<TypeConfig>(LeadershipState::Shutdown);
assert_eq!(s, LeaderState::Unknown);
}
}