use core::pin::Pin;
use futures::Stream;
use tsoracle_core::Epoch;
use crate::error::ConsensusError;
use crate::leadership::LeaderState;
#[async_trait::async_trait]
pub trait ConsensusDriver: Send + Sync + 'static {
fn leadership_events(&self) -> Pin<Box<dyn Stream<Item = LeaderState> + Send>>;
async fn load_high_water(&self) -> Result<u64, ConsensusError>;
async fn persist_high_water(&self, at_least: u64, epoch: Epoch) -> Result<u64, ConsensusError>;
async fn load_dense_seq(&self, key: &tsoracle_core::SeqKey) -> Result<u64, ConsensusError> {
let _ = key;
Err(ConsensusError::DenseUnsupported)
}
async fn advance_dense(
&self,
key: &tsoracle_core::SeqKey,
count: u32,
expected_epoch: Epoch,
) -> Result<u64, ConsensusError> {
let (_, _, _) = (key, count, expected_epoch);
Err(ConsensusError::DenseUnsupported)
}
async fn advance_dense_batch(
&self,
entries: &[(tsoracle_core::SeqKey, u32)],
expected_epoch: Epoch,
) -> Result<Vec<u64>, ConsensusError> {
let (_, _) = (entries, expected_epoch);
Err(ConsensusError::DenseUnsupported)
}
async fn load_leases(&self) -> Result<Vec<tsoracle_core::LeaseRecord>, ConsensusError> {
Err(ConsensusError::LeasesUnsupported)
}
async fn persist_leases(
&self,
live: &[tsoracle_core::LeaseRecord],
epoch: Epoch,
) -> Result<(), ConsensusError> {
let (_, _) = (live, epoch);
Err(ConsensusError::LeasesUnsupported)
}
}
#[cfg(test)]
mod tests {
use super::*;
use futures::StreamExt;
use futures::stream;
struct Dummy;
#[async_trait::async_trait]
impl ConsensusDriver for Dummy {
fn leadership_events(&self) -> Pin<Box<dyn Stream<Item = LeaderState> + Send>> {
Box::pin(stream::once(async { LeaderState::Unknown }))
}
async fn load_high_water(&self) -> Result<u64, ConsensusError> {
Ok(0)
}
async fn persist_high_water(
&self,
at_least: u64,
_epoch: Epoch,
) -> Result<u64, ConsensusError> {
Ok(at_least)
}
}
#[test]
fn dummy_is_object_safe() {
let _: Box<dyn ConsensusDriver> = Box::new(Dummy);
}
#[test]
fn dummy_driver_methods_return_documented_defaults() {
let driver: Box<dyn ConsensusDriver> = Box::new(Dummy);
futures::executor::block_on(async {
let mut events = driver.leadership_events();
assert_eq!(
events.next().await,
Some(LeaderState::Unknown),
"first item must be the current state, synchronously",
);
assert!(events.next().await.is_none(), "no transitions after");
assert_eq!(driver.load_high_water().await.unwrap(), 0);
assert_eq!(
driver.persist_high_water(42, Epoch(7)).await.unwrap(),
42,
"persist returns the at_least argument unchanged",
);
let key = tsoracle_core::SeqKey::try_new("k").unwrap();
assert!(
matches!(
driver.load_dense_seq(&key).await,
Err(ConsensusError::DenseUnsupported)
),
"default load_dense_seq is DenseUnsupported",
);
assert!(
matches!(
driver.advance_dense(&key, 1, Epoch(1)).await,
Err(ConsensusError::DenseUnsupported)
),
"default advance_dense is DenseUnsupported",
);
let entries = vec![(tsoracle_core::SeqKey::try_new("k").unwrap(), 1u32)];
assert!(
matches!(
driver.advance_dense_batch(&entries, Epoch(1)).await,
Err(ConsensusError::DenseUnsupported)
),
"default advance_dense_batch is DenseUnsupported",
);
assert!(
matches!(
driver.load_leases().await,
Err(ConsensusError::LeasesUnsupported)
),
"default load_leases is LeasesUnsupported",
);
let record = tsoracle_core::LeaseRecord {
lease_id: 1,
holder: b"g1".to_vec(),
holder_epoch: 1,
ttl_ms: 10_000,
ts_upper_bound: 1,
expires_at_ms: 10_001,
superseded: false,
};
assert!(
matches!(
driver.persist_leases(&[record], Epoch(1)).await,
Err(ConsensusError::LeasesUnsupported)
),
"default persist_leases is LeasesUnsupported",
);
});
}
}