use async_trait::async_trait;
use openraft::RaftMetrics;
use openraft::RaftTypeConfig;
use openraft::type_config::alias::WatchReceiverOf;
use tsoracle_consensus::ConsensusError;
#[async_trait]
pub trait OpenraftHighWaterHost: Send + Sync + 'static {
type Config: RaftTypeConfig;
fn metrics(&self) -> WatchReceiverOf<Self::Config, RaftMetrics<Self::Config>>;
async fn current_high_water(&self) -> Result<u64, ConsensusError>;
async fn submit_advance(&self, at_least: u64) -> Result<u64, ConsensusError>;
fn active_write_version(&self) -> u8;
async fn submit_advance_dense(
&self,
key: &tsoracle_core::SeqKey,
count: u32,
) -> Result<u64, ConsensusError>;
async fn current_dense_seq(&self, key: &tsoracle_core::SeqKey) -> Result<u64, ConsensusError>;
}
#[cfg(test)]
mod tests {
use super::*;
use crate::type_config::TypeConfig;
use openraft::RaftMetrics;
use openraft::type_config::TypeConfigExt;
use openraft::type_config::alias::WatchReceiverOf;
struct MetricsOnlyHost {
rx: WatchReceiverOf<TypeConfig, RaftMetrics<TypeConfig>>,
}
#[async_trait]
impl OpenraftHighWaterHost for MetricsOnlyHost {
type Config = TypeConfig;
fn metrics(&self) -> WatchReceiverOf<Self::Config, RaftMetrics<Self::Config>> {
self.rx.clone()
}
async fn current_high_water(&self) -> Result<u64, ConsensusError> {
Ok(0)
}
async fn submit_advance(&self, _at_least: u64) -> Result<u64, ConsensusError> {
Ok(0)
}
fn active_write_version(&self) -> u8 {
tsoracle_openraft_toolkit::BASELINE_WRITE_VERSION
}
async fn submit_advance_dense(
&self,
_key: &tsoracle_core::SeqKey,
_count: u32,
) -> Result<u64, ConsensusError> {
Err(ConsensusError::DenseUnsupported)
}
async fn current_dense_seq(
&self,
_key: &tsoracle_core::SeqKey,
) -> Result<u64, ConsensusError> {
Err(ConsensusError::DenseUnsupported)
}
}
fn assert_implements_host<H: OpenraftHighWaterHost>() {}
#[tokio::test]
async fn metrics_only_host_satisfies_and_drives_trait() {
assert_implements_host::<MetricsOnlyHost>();
let metrics: RaftMetrics<TypeConfig> = RaftMetrics::new_initial(1u64);
let (_tx, rx) = <TypeConfig as TypeConfigExt>::watch_channel(metrics);
let host = MetricsOnlyHost { rx };
let _rx = host.metrics();
assert!(host.current_high_water().await.is_ok());
assert!(host.submit_advance(7).await.is_ok());
let _ = host.active_write_version();
let key = tsoracle_core::SeqKey::try_new("k").unwrap();
assert!(host.submit_advance_dense(&key, 1).await.is_err());
assert!(host.current_dense_seq(&key).await.is_err());
}
}