#[cfg(not(feature = "std"))]
use alloc::vec::Vec;
#[cfg(feature = "std")]
use std::vec::Vec;
use core::time::Duration;
use crate::{
anchor::{
Anchor,
NoAnchor,
},
batch::Batch,
engine::{
ApplySummary,
Engine,
EngineConfig,
},
error::{
ApplyError,
ConfigError,
DivergenceCause,
DurabilityLost,
EngineStatus,
RollbackError,
},
fold::Fold,
position::{
BlockRef,
Position,
},
sink::{
NoSink,
SnapshotSink,
},
source::{
EventSource,
ProbeSource,
ReplayHorizon,
},
};
const DEFAULT_POLL_INTERVAL: Duration = Duration::from_secs(1);
const DEFAULT_BACKOFF_BASE: Duration = Duration::from_millis(200);
const DEFAULT_BACKOFF_MAX: Duration = Duration::from_secs(30);
const DEFAULT_MAX_DIVERGENCE_RETRIES: u32 = 1;
pub trait Tickable {
fn tick(&mut self) -> Tick;
fn checkpoint(&mut self);
fn status(&self) -> DriverStatus;
fn next_delay(&self) -> Duration;
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Tick {
Progressed(ApplySummary),
Idle,
RolledBack {
to: Option<Position>,
},
Resynced,
SourceError,
DurabilityLost,
Terminal(EngineStatus),
}
#[derive(Debug, Clone, Copy, PartialEq)]
pub struct DriverStatus {
pub cursor: Option<Position>,
pub last_verified: Option<BlockRef>,
pub engine: EngineStatus,
pub caught_up: bool,
pub skips: u64,
pub durable_cursor: Option<Position>,
pub durability_lost: bool,
pub generation: u64,
}
impl DriverStatus {
pub fn is_terminal(&self) -> bool {
!self.engine.is_active()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct DriverConfig {
pub start_block: u64,
pub poll_interval: Duration,
pub backoff_base: Duration,
pub backoff_max: Duration,
pub checkpoint_interval: Option<u64>,
pub snapshot_interval: Option<u64>,
pub max_divergence_retries: u32,
}
impl Default for DriverConfig {
fn default() -> Self {
Self {
start_block: 0,
poll_interval: DEFAULT_POLL_INTERVAL,
backoff_base: DEFAULT_BACKOFF_BASE,
backoff_max: DEFAULT_BACKOFF_MAX,
checkpoint_interval: None,
snapshot_interval: None,
max_divergence_retries: DEFAULT_MAX_DIVERGENCE_RETRIES,
}
}
}
pub struct Driver<F, S, A = NoAnchor<<F as Fold>::View>, K = NoSink>
where
F: Fold,
S: EventSource<Event = F::Event>,
A: Anchor<View = F::View>,
K: SnapshotSink<F>,
{
engine: Engine<F>,
source: S,
anchor: Option<A>,
sink: K,
config: DriverConfig,
batch: Batch<F::Event>,
initial: F,
consecutive_errors: u32,
caught_up: bool,
generation: u64,
divergence_retries: u32,
last_checkpoint_block: Option<u64>,
last_snapshot_block: Option<u64>,
durability_lost: bool,
advanced: bool,
}
impl<F, S> Driver<F, S>
where
F: Fold + Clone,
S: EventSource<Event = F::Event>,
{
pub fn new(
fold: F,
source: S,
engine: EngineConfig,
config: DriverConfig,
) -> Result<Self, ConfigError> {
Self::build(fold, source, None, NoSink, engine, config)
}
pub fn resume(
engine: Engine<F>,
source: S,
genesis: F,
config: DriverConfig,
) -> Result<Self, ConfigError> {
Self::around(engine, source, None, NoSink, genesis, config)
}
}
impl<F, S, K> Driver<F, S, NoAnchor<<F as Fold>::View>, K>
where
F: Fold + Clone,
S: EventSource<Event = F::Event>,
K: SnapshotSink<F>,
{
pub fn with_sink(
fold: F,
source: S,
sink: K,
engine: EngineConfig,
config: DriverConfig,
) -> Result<Self, ConfigError> {
Self::build(fold, source, None, sink, engine, config)
}
pub fn resume_with_sink(
engine: Engine<F>,
source: S,
sink: K,
genesis: F,
config: DriverConfig,
) -> Result<Self, ConfigError> {
Self::around(engine, source, None, sink, genesis, config)
}
}
impl<F, S, A> Driver<F, S, A, NoSink>
where
F: Fold + Clone,
S: EventSource<Event = F::Event>,
A: Anchor<View = F::View>,
{
pub fn with_anchor(
fold: F,
source: S,
anchor: A,
engine: EngineConfig,
config: DriverConfig,
) -> Result<Self, ConfigError> {
Self::build(fold, source, Some(anchor), NoSink, engine, config)
}
}
impl<F, S, A, K> Driver<F, S, A, K>
where
F: Fold + Clone,
S: EventSource<Event = F::Event>,
A: Anchor<View = F::View>,
K: SnapshotSink<F>,
{
fn build(
fold: F,
source: S,
anchor: Option<A>,
sink: K,
engine_config: EngineConfig,
driver_config: DriverConfig,
) -> Result<Self, ConfigError> {
let initial = fold.clone();
let engine = Engine::new(fold, engine_config)?;
Self::around(engine, source, anchor, sink, initial, driver_config)
}
fn around(
engine: Engine<F>,
source: S,
anchor: Option<A>,
sink: K,
initial: F,
driver_config: DriverConfig,
) -> Result<Self, ConfigError> {
if let ReplayHorizon::FromBlock(horizon) = source.horizon()
&& horizon > driver_config.start_block
{
return Err(ConfigError::HorizonExceedsStart {
start: driver_config.start_block,
horizon,
});
}
Ok(Self {
engine,
source,
anchor,
sink,
config: driver_config,
batch: Batch::new(),
initial,
consecutive_errors: 0,
caught_up: false,
generation: 0,
divergence_retries: 0,
last_checkpoint_block: None,
last_snapshot_block: None,
durability_lost: false,
advanced: false,
})
}
pub fn with_anchor_and_sink(
fold: F,
source: S,
anchor: A,
sink: K,
engine: EngineConfig,
config: DriverConfig,
) -> Result<Self, ConfigError> {
Self::build(fold, source, Some(anchor), sink, engine, config)
}
pub fn resume_with_anchor_and_sink(
engine: Engine<F>,
source: S,
anchor: A,
sink: K,
genesis: F,
config: DriverConfig,
) -> Result<Self, ConfigError> {
Self::around(engine, source, Some(anchor), sink, genesis, config)
}
pub fn sink(&self) -> &K {
&self.sink
}
pub fn into_sink(self) -> K {
self.sink
}
pub fn engine(&self) -> &Engine<F> {
&self.engine
}
pub fn engine_mut(&mut self) -> &mut Engine<F> {
&mut self.engine
}
pub fn source_mut(&mut self) -> &mut S {
&mut self.source
}
pub fn is_caught_up(&self) -> bool {
self.caught_up
}
fn auto_checkpoint(&mut self) -> Option<Tick> {
let interval = self.config.checkpoint_interval?;
let cursor = self.engine.cursor()?;
let due = self
.last_checkpoint_block
.is_none_or(|last| cursor.block >= last.saturating_add(interval));
if due { self.run_checkpoint() } else { None }
}
fn run_checkpoint(&mut self) -> Option<Tick> {
self.engine.checkpoint();
if let Some(cursor) = self.engine.cursor() {
self.last_checkpoint_block = Some(cursor.block);
}
self.check_anchor()
}
fn offer_snapshot(&mut self) -> Option<Tick> {
if self.durability_lost {
return None;
}
let interval = self.config.snapshot_interval?;
let point = self.engine.durable_point()?;
let due = self
.last_snapshot_block
.is_none_or(|last| point.block >= last.saturating_add(interval));
if !due {
return None;
}
match self.sink.offer(&self.engine) {
Ok(()) => {
self.last_snapshot_block = Some(point.block);
None
}
Err(DurabilityLost) => Some(self.lose_durability()),
}
}
#[cold]
fn lose_durability(&mut self) -> Tick {
self.durability_lost = true;
Tick::DurabilityLost
}
#[cold]
fn clamp_snapshot_mark(&mut self, to: Option<Position>) {
self.last_snapshot_block = self
.last_snapshot_block
.zip(to)
.map(|(last, point)| last.min(point.block));
}
fn check_anchor(&mut self) -> Option<Tick> {
let anchor = self.anchor.as_ref()?;
let at = self.engine.last_verified()?;
let expected = anchor.expected(&at)?;
if self.engine.view() == expected {
self.divergence_retries = 0;
None
} else {
Some(self.handle_anchor_divergence(at))
}
}
#[cold]
fn handle_anchor_divergence(&mut self, at: BlockRef) -> Tick {
if self.divergence_retries >= self.config.max_divergence_retries {
self.engine
.mark_unrecoverable(DivergenceCause::AnchorDivergence { at: at.number });
return Tick::Terminal(self.engine.status());
}
let tick = self.roll_back_to(at.number.saturating_sub(1));
if matches!(tick, Tick::RolledBack { .. }) {
self.divergence_retries = self.divergence_retries.saturating_add(1);
}
tick
}
#[cold]
fn roll_back_to(&mut self, ancestor: u64) -> Tick {
match self.engine.rollback_at_or_below(ancestor) {
Ok(to) => {
self.caught_up = false;
self.clamp_snapshot_mark(to);
Tick::RolledBack { to }
}
Err(RollbackError::NoCheckpointAtOrBelow { .. }) => self.resync_or_terminal(),
Err(RollbackError::Unrecoverable { cause }) => {
Tick::Terminal(EngineStatus::Unrecoverable { cause })
}
}
}
#[cold]
fn resync_or_terminal(&mut self) -> Tick {
match self.source.horizon() {
ReplayHorizon::FromBlock(horizon) if horizon > self.config.start_block => {
self.engine
.mark_unrecoverable(DivergenceCause::HorizonExceeded {
needed: self.config.start_block,
horizon,
});
Tick::Terminal(self.engine.status())
}
_ => self.resync(),
}
}
#[cold]
fn resync(&mut self) -> Tick {
self.engine.reset(self.initial.clone());
self.caught_up = false;
self.consecutive_errors = 0;
self.divergence_retries = 0;
self.last_checkpoint_block = None;
self.last_snapshot_block = None;
Tick::Resynced
}
fn step<Fork>(&mut self, on_fork: Fork) -> Tick
where
Fork: FnOnce(&mut Self) -> Tick,
{
let tick = self.poll_apply(on_fork);
self.advanced = matches!(
tick,
Tick::Progressed(summary) if summary.applied > 0 || summary.skipped > 0
);
tick
}
fn poll_apply<Fork>(&mut self, on_fork: Fork) -> Tick
where
Fork: FnOnce(&mut Self) -> Tick,
{
self.generation = self.generation.wrapping_add(1);
if !self.engine.status().is_active() {
return Tick::Terminal(self.engine.status());
}
if self
.source
.next_batch(self.engine.cursor(), &mut self.batch)
.is_err()
{
self.consecutive_errors = self.consecutive_errors.saturating_add(1);
return Tick::SourceError;
}
self.consecutive_errors = 0;
match self.engine.apply_batch(&self.batch) {
Ok(summary) => {
self.caught_up = self.batch.is_empty();
if let Some(tick) =
self.auto_checkpoint().or_else(|| self.offer_snapshot())
{
return tick;
}
if self.batch.is_empty() {
Tick::Idle
} else {
Tick::Progressed(summary)
}
}
Err(
ApplyError::ForkSuspected { .. }
| ApplyError::MissingBoundary
| ApplyError::CursorBlockUnobserved { .. },
) => on_fork(self),
Err(ApplyError::Halted { .. } | ApplyError::Poisoned { .. }) => {
Tick::Terminal(self.engine.status())
}
Err(ApplyError::Shape(_) | ApplyError::BoundaryNumberMismatch { .. }) => {
self.consecutive_errors = self.consecutive_errors.saturating_add(1);
Tick::SourceError
}
Err(ApplyError::NotActive { .. }) => {
unreachable!(
"engine status was checked active before this apply_batch call"
)
}
}
}
}
impl<F, S, A, K> Tickable for Driver<F, S, A, K>
where
F: Fold + Clone,
S: EventSource<Event = F::Event>,
A: Anchor<View = F::View>,
K: SnapshotSink<F>,
{
fn tick(&mut self) -> Tick {
self.step(Self::resync_or_terminal)
}
fn checkpoint(&mut self) {
self.run_checkpoint();
}
fn status(&self) -> DriverStatus {
DriverStatus {
cursor: self.engine.cursor(),
last_verified: self.engine.last_verified(),
engine: self.engine.status(),
caught_up: self.caught_up,
skips: self.engine.skip_count(),
durable_cursor: self.sink.durable_cursor(),
durability_lost: self.durability_lost,
generation: self.generation,
}
}
fn next_delay(&self) -> Duration {
if self.consecutive_errors == 0 {
return if self.advanced {
Duration::ZERO
} else {
self.config.poll_interval
};
}
let factor = 1u32
.checked_shl(self.consecutive_errors - 1)
.unwrap_or(u32::MAX);
self.config
.backoff_base
.saturating_mul(factor)
.min(self.config.backoff_max)
}
}
impl<F, S, A, K> Driver<F, S, A, K>
where
F: Fold + Clone,
S: ProbeSource<Event = F::Event>,
A: Anchor<View = F::View>,
K: SnapshotSink<F>,
{
#[cold]
fn recover_via_bisection(&mut self) -> Tick {
let observed: Vec<BlockRef> = self.engine.observed().collect();
let mut lo = 0usize;
let mut hi = observed.len();
while lo < hi {
let mid = lo + (hi - lo) / 2;
match self.source.header_at(observed[mid].number) {
Ok(Some(header)) if header.hash == observed[mid].hash => lo = mid + 1,
Ok(_) => hi = mid,
Err(_) => {
self.consecutive_errors = self.consecutive_errors.saturating_add(1);
return Tick::SourceError;
}
}
}
if lo == 0 {
return self.resync_or_terminal();
}
self.roll_back_to(observed[lo - 1].number)
}
}
pub struct Probed<F, S, A = NoAnchor<<F as Fold>::View>, K = NoSink>
where
F: Fold,
S: ProbeSource<Event = F::Event>,
A: Anchor<View = F::View>,
K: SnapshotSink<F>,
{
inner: Driver<F, S, A, K>,
}
impl<F, S> Probed<F, S>
where
F: Fold + Clone,
S: ProbeSource<Event = F::Event>,
{
pub fn new(
fold: F,
source: S,
engine: EngineConfig,
config: DriverConfig,
) -> Result<Self, ConfigError> {
Ok(Self {
inner: Driver::new(fold, source, engine, config)?,
})
}
pub fn resume(
engine: Engine<F>,
source: S,
genesis: F,
config: DriverConfig,
) -> Result<Self, ConfigError> {
Ok(Self {
inner: Driver::resume(engine, source, genesis, config)?,
})
}
}
impl<F, S, K> Probed<F, S, NoAnchor<<F as Fold>::View>, K>
where
F: Fold + Clone,
S: ProbeSource<Event = F::Event>,
K: SnapshotSink<F>,
{
pub fn with_sink(
fold: F,
source: S,
sink: K,
engine: EngineConfig,
config: DriverConfig,
) -> Result<Self, ConfigError> {
Ok(Self {
inner: Driver::with_sink(fold, source, sink, engine, config)?,
})
}
}
impl<F, S, A, K> Probed<F, S, A, K>
where
F: Fold + Clone,
S: ProbeSource<Event = F::Event>,
A: Anchor<View = F::View>,
K: SnapshotSink<F>,
{
pub fn sink(&self) -> &K {
self.inner.sink()
}
pub fn into_sink(self) -> K {
self.inner.into_sink()
}
pub fn engine(&self) -> &Engine<F> {
self.inner.engine()
}
pub fn engine_mut(&mut self) -> &mut Engine<F> {
self.inner.engine_mut()
}
pub fn source_mut(&mut self) -> &mut S {
self.inner.source_mut()
}
pub fn is_caught_up(&self) -> bool {
self.inner.is_caught_up()
}
}
impl<F, S, A, K> Tickable for Probed<F, S, A, K>
where
F: Fold + Clone,
S: ProbeSource<Event = F::Event>,
A: Anchor<View = F::View>,
K: SnapshotSink<F>,
{
fn tick(&mut self) -> Tick {
self.inner.step(Driver::recover_via_bisection)
}
fn checkpoint(&mut self) {
self.inner.checkpoint();
}
fn status(&self) -> DriverStatus {
self.inner.status()
}
fn next_delay(&self) -> Duration {
self.inner.next_delay()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[cfg(not(feature = "std"))]
use alloc::{
vec,
vec::Vec,
};
#[cfg(feature = "std")]
use std::{
vec,
vec::Vec,
};
use crate::{
batch::{
BlockSpan,
LogEvent,
},
test_util::{
FailKind,
PollFailure,
RecordingFold,
ScriptedChain,
WatermarkSink,
},
};
struct Probe {
inner: ScriptedChain,
calls: u32,
fail_next: u32,
}
impl Probe {
fn new(inner: ScriptedChain) -> Self {
Self {
inner,
calls: 0,
fail_next: 0,
}
}
fn fail_next_probes(&mut self, n: u32) {
self.fail_next = n;
}
}
impl EventSource for Probe {
type Event = u64;
type Error = PollFailure;
fn next_batch(
&mut self,
cursor: Option<Position>,
out: &mut Batch<u64>,
) -> Result<(), PollFailure> {
self.inner.next_batch(cursor, out)
}
fn horizon(&self) -> ReplayHorizon {
self.inner.horizon()
}
}
impl ProbeSource for Probe {
fn header_at(&mut self, number: u64) -> Result<Option<BlockRef>, PollFailure> {
self.calls += 1;
if self.fail_next > 0 {
self.fail_next -= 1;
return Err(PollFailure);
}
self.inner.header_at(number)
}
}
struct Stuck {
block: BlockRef,
}
impl EventSource for Stuck {
type Event = u64;
type Error = PollFailure;
fn next_batch(
&mut self,
cursor: Option<Position>,
out: &mut Batch<u64>,
) -> Result<(), PollFailure> {
out.clear();
out.boundary = cursor.map(|_| self.block);
out.spans.push(BlockSpan {
block: self.block,
start: 0,
end: 1,
});
out.events.push(LogEvent {
log_index: 0,
event: 1,
});
Ok(())
}
fn horizon(&self) -> ReplayHorizon {
ReplayHorizon::Genesis
}
}
struct ExactAnchor;
impl Anchor for ExactAnchor {
type View = Vec<(Position, u64)>;
fn expected(&self, at: &BlockRef) -> Option<Self::View> {
Some((1..=at.number).map(|n| (Position::new(n, 0), n)).collect())
}
}
struct DisagreeingAnchor;
impl Anchor for DisagreeingAnchor {
type View = Vec<(Position, u64)>;
fn expected(&self, _at: &BlockRef) -> Option<Self::View> {
Some(Vec::new())
}
}
fn engine_config(checkpoint_slots: usize) -> EngineConfig {
EngineConfig {
ring_capacity: 8,
checkpoint_slots,
}
}
fn new_driver(
chain: ScriptedChain,
engine: EngineConfig,
config: DriverConfig,
) -> Driver<RecordingFold, ScriptedChain> {
Driver::new(RecordingFold::default(), chain, engine, config).unwrap()
}
fn run_to_idle<D: Tickable>(driver: &mut D) -> Tick {
let mut outcome = driver.tick();
while !matches!(outcome, Tick::Idle) {
outcome = driver.tick();
}
outcome
}
fn collect_to_idle<D: Tickable>(driver: &mut D) -> Vec<Tick> {
let mut ticks = vec![driver.tick()];
while !matches!(ticks.last(), Some(Tick::Idle)) {
ticks.push(driver.tick());
}
ticks
}
fn one_event_chain(blocks: u64) -> ScriptedChain {
let mut chain = ScriptedChain::new(1);
for value in 1..=blocks {
chain.push_block(&[value]);
}
chain.set_batch_blocks(1);
chain
}
fn cadence_config(checkpoint: u64, snapshot: u64) -> DriverConfig {
DriverConfig {
checkpoint_interval: Some(checkpoint),
snapshot_interval: Some(snapshot),
..DriverConfig::default()
}
}
type SinkDriver = Driver<
RecordingFold,
ScriptedChain,
NoAnchor<Vec<(Position, u64)>>,
WatermarkSink,
>;
type ProbedSinkDriver = Probed<
RecordingFold,
ScriptedChain,
NoAnchor<Vec<(Position, u64)>>,
WatermarkSink,
>;
fn sink_driver(blocks: u64, slots: usize, config: DriverConfig) -> SinkDriver {
Driver::with_sink(
RecordingFold::default(),
one_event_chain(blocks),
WatermarkSink::default(),
engine_config(slots),
config,
)
.unwrap()
}
fn probed_sink_driver(
blocks: u64,
slots: usize,
config: DriverConfig,
) -> ProbedSinkDriver {
Probed::with_sink(
RecordingFold::default(),
one_event_chain(blocks),
WatermarkSink::default(),
engine_config(slots),
config,
)
.unwrap()
}
#[test]
fn snapshot_interval_offers_on_durable_point_cadence() {
let mut driver = sink_driver(12, 3, cadence_config(2, 4));
run_to_idle(&mut driver);
assert_eq!(
driver.sink().offered,
vec![Position::new(1, 0), Position::new(5, 0)]
);
}
#[test]
fn no_sink_never_offers_and_reports_no_durable_cursor() {
let mut driver =
new_driver(one_event_chain(4), engine_config(2), cadence_config(1, 1));
let ticks = collect_to_idle(&mut driver);
assert!(!ticks.contains(&Tick::DurabilityLost));
assert_eq!(driver.status().durable_cursor, None);
assert!(!driver.status().durability_lost);
}
#[test]
fn zero_checkpoint_slots_never_offers() {
let mut driver = sink_driver(4, 0, cadence_config(1, 1));
run_to_idle(&mut driver);
assert!(driver.sink().offered.is_empty());
assert_eq!(driver.status().durable_cursor, None);
}
#[test]
fn durable_cursor_trails_the_live_cursor_by_checkpoint_coverage() {
let mut driver = sink_driver(12, 3, cadence_config(2, 1));
run_to_idle(&mut driver);
let status = driver.status();
assert_eq!(status.cursor, Some(Position::new(12, 0)));
assert_eq!(status.durable_cursor, Some(Position::new(7, 0)));
assert!(status.cursor.unwrap().block - status.durable_cursor.unwrap().block >= 4);
}
#[test]
fn reorg_across_the_live_cursor_leaves_the_durable_cursor_untouched() {
let mut driver = probed_sink_driver(8, 4, cadence_config(1, 1));
run_to_idle(&mut driver);
let before = driver.status().durable_cursor;
assert_eq!(before, Some(Position::new(5, 0)));
driver.source_mut().reorg(2, &[&[70], &[80]]);
let outcome = driver.tick();
let status = driver.status();
assert_eq!(
outcome,
Tick::RolledBack {
to: Some(Position::new(6, 0)),
}
);
assert_eq!(status.durable_cursor, before);
assert!(status.durable_cursor.unwrap() <= Position::new(6, 0));
}
#[test]
fn sink_failure_is_reported_once_then_folding_continues() {
let mut driver = Driver::with_sink(
RecordingFold::default(),
one_event_chain(6),
WatermarkSink {
offered: Vec::new(),
fail_next_offers: 1,
},
engine_config(2),
cadence_config(1, 1),
)
.unwrap();
let ticks = collect_to_idle(&mut driver);
let lost = ticks
.iter()
.filter(|tick| **tick == Tick::DurabilityLost)
.count();
assert_eq!(lost, 1);
assert!(driver.status().durability_lost);
assert!(driver.sink().offered.is_empty());
let expected: Vec<(Position, u64)> = (1..=6u64)
.map(|value| (Position::new(value, 0), value))
.collect();
assert_eq!(driver.engine().view(), expected);
}
#[test]
fn rollback_does_not_suppress_the_next_offer() {
let mut driver = probed_sink_driver(8, 4, cadence_config(1, 1));
run_to_idle(&mut driver);
driver.source_mut().reorg(3, &[&[60], &[70], &[80], &[90]]);
let outcome = driver.tick();
let Tick::RolledBack { to } = outcome else {
panic!("expected RolledBack, got {outcome:?}");
};
let restore = to.expect("the rollback restores a cursor").block;
run_to_idle(&mut driver);
let offered = &driver.sink().offered;
assert!(offered.windows(2).all(|pair| pair[0].block < pair[1].block));
assert!(offered.last().expect("offers were made").block > restore);
}
#[test]
fn resync_lowers_the_reported_durable_cursor() {
let mut driver = sink_driver(8, 4, cadence_config(1, 1));
run_to_idle(&mut driver);
let before = driver.status().durable_cursor.expect("offers were made");
assert_eq!(before, Position::new(5, 0));
driver.source_mut().reorg(8, &[&[10], &[20], &[30], &[40]]);
let outcome = driver.tick();
run_to_idle(&mut driver);
assert_eq!(outcome, Tick::Resynced);
let after = driver.status().durable_cursor.expect("offers resumed");
assert!(after < before);
}
#[test]
fn driver_folds_to_tip_and_reports_caught_up() {
let mut chain = ScriptedChain::new(1);
for value in 1..=10u64 {
chain.push_block(&[value]);
}
let mut driver = new_driver(chain, engine_config(0), DriverConfig::default());
run_to_idle(&mut driver);
let expected: Vec<(Position, u64)> = (1..=10u64)
.map(|value| (Position::new(value, 0), value))
.collect();
assert_eq!(driver.engine().view(), expected);
assert!(driver.is_caught_up());
}
#[test]
fn empty_poll_is_idle() {
let mut chain = ScriptedChain::new(1);
chain.push_block(&[1]);
chain.push_block(&[2]);
let mut driver = new_driver(chain, engine_config(0), DriverConfig::default());
driver.tick();
let outcome = driver.tick();
assert_eq!(outcome, Tick::Idle);
}
#[test]
fn source_errors_back_off_exponentially() {
let mut chain = ScriptedChain::new(1);
chain.push_block(&[1]);
chain.fail_next_polls(3);
let mut driver = new_driver(chain, engine_config(0), DriverConfig::default());
let first = driver.tick();
let first_delay = driver.next_delay();
let second = driver.tick();
let second_delay = driver.next_delay();
let third = driver.tick();
let third_delay = driver.next_delay();
let fourth = driver.tick();
assert_eq!(first, Tick::SourceError);
assert_eq!(first_delay, Duration::from_millis(200));
assert_eq!(second, Tick::SourceError);
assert_eq!(second_delay, Duration::from_millis(400));
assert_eq!(third, Tick::SourceError);
assert_eq!(third_delay, Duration::from_millis(800));
assert!(matches!(fourth, Tick::Progressed(_)));
assert_eq!(driver.next_delay(), Duration::ZERO);
}
#[test]
fn backoff_caps_at_max() {
let mut chain = ScriptedChain::new(1);
chain.fail_next_polls(u32::MAX);
let mut driver = new_driver(chain, engine_config(0), DriverConfig::default());
for _ in 0..10 {
driver.tick();
}
assert_eq!(driver.next_delay(), Duration::from_secs(30));
}
#[test]
fn catch_up_ticks_ask_for_no_delay_until_the_tip() {
let mut driver = new_driver(
one_event_chain(10),
engine_config(0),
DriverConfig::default(),
);
driver.tick();
let while_behind = driver.next_delay();
run_to_idle(&mut driver);
let at_tip = driver.next_delay();
assert_eq!(while_behind, Duration::ZERO);
assert_eq!(at_tip, Duration::from_secs(1));
}
#[test]
fn a_fully_deduped_batch_keeps_the_poll_interval() {
let block = BlockRef {
number: 1,
hash: [7u8; 32],
};
let mut driver = Driver::new(
RecordingFold::default(),
Stuck { block },
engine_config(0),
DriverConfig::default(),
)
.unwrap();
let applying = driver.tick();
let after_apply = driver.next_delay();
let deduping = driver.tick();
let after_dedup = driver.next_delay();
assert_eq!(
applying,
Tick::Progressed(ApplySummary {
applied: 1,
deduped: 0,
skipped: 0,
})
);
assert_eq!(after_apply, Duration::ZERO);
assert_eq!(
deduping,
Tick::Progressed(ApplySummary {
applied: 0,
deduped: 1,
skipped: 0,
})
);
assert_eq!(after_dedup, Duration::from_secs(1));
}
#[test]
fn fork_without_probe_resyncs_from_start() {
let mut chain = ScriptedChain::new(1);
for value in 1..=5u64 {
chain.push_block(&[value]);
}
let mut driver = new_driver(chain, engine_config(0), DriverConfig::default());
driver.tick();
driver.source_mut().reorg(3, &[&[10], &[20], &[30]]);
let outcome = driver.tick();
assert_eq!(outcome, Tick::Resynced);
run_to_idle(&mut driver);
let expected = vec![
(Position::new(1, 0), 1),
(Position::new(2, 0), 2),
(Position::new(3, 0), 10),
(Position::new(4, 0), 20),
(Position::new(5, 0), 30),
];
assert_eq!(driver.engine().view(), expected);
}
#[test]
fn resync_with_moved_horizon_is_terminal() {
let mut chain = ScriptedChain::new(1);
for value in 1..=5u64 {
chain.push_block(&[value]);
}
let mut driver = new_driver(chain, engine_config(0), DriverConfig::default());
driver.tick();
driver.source_mut().reorg(3, &[&[10], &[20], &[30]]);
driver.source_mut().set_horizon(ReplayHorizon::FromBlock(1));
let outcome = driver.tick();
assert_eq!(
outcome,
Tick::Terminal(EngineStatus::Unrecoverable {
cause: DivergenceCause::HorizonExceeded {
needed: 0,
horizon: 1,
},
})
);
}
#[test]
fn construction_refuses_horizon_above_start() {
let mut chain = ScriptedChain::new(1);
chain.set_horizon(ReplayHorizon::FromBlock(100));
let result = Driver::new(
RecordingFold::default(),
chain,
engine_config(0),
DriverConfig::default(),
);
assert_eq!(
result.err(),
Some(ConfigError::HorizonExceedsStart {
start: 0,
horizon: 100,
})
);
}
#[test]
fn auto_checkpoint_follows_interval() {
let mut chain = ScriptedChain::new(1);
for value in 1..=12u64 {
chain.push_block(&[value]);
}
chain.set_batch_blocks(1);
let config = DriverConfig {
checkpoint_interval: Some(4),
..DriverConfig::default()
};
let mut driver = new_driver(chain, engine_config(8), config);
run_to_idle(&mut driver);
assert!(driver.engine().checkpoint_count() >= 3);
}
#[test]
fn halt_is_terminal_and_recoverable_via_engine_mut() {
let mut chain = ScriptedChain::new(1);
chain.push_block(&[1]);
chain.push_block(&[2]);
chain.push_block(&[3]);
chain.push_block(&[4]);
chain.push_block(&[5]);
chain.set_batch_blocks(1);
let halt_pos = Position::new(3, 0);
let fold = RecordingFold {
applied: Vec::new(),
fail_at: Some((halt_pos, FailKind::Halt)),
};
let mut driver =
Driver::new(fold, chain, engine_config(2), DriverConfig::default()).unwrap();
driver.tick();
driver.tick();
driver.checkpoint();
let outcome = driver.tick();
assert_eq!(
outcome,
Tick::Terminal(EngineStatus::Halted { at: halt_pos })
);
driver.source_mut().reorg(3, &[&[], &[40], &[50]]);
let restored = driver.engine_mut().rollback_at_or_below(2).unwrap();
assert_eq!(restored, Some(Position::new(2, 0)));
assert_eq!(driver.engine().status(), EngineStatus::Active);
run_to_idle(&mut driver);
assert!(driver.is_caught_up());
assert_eq!(driver.engine().cursor(), Some(Position::new(5, 0)));
}
#[test]
fn generation_increments_every_tick() {
let chain = ScriptedChain::new(1);
let mut driver = new_driver(chain, engine_config(0), DriverConfig::default());
let start = driver.status().generation;
driver.tick();
driver.tick();
driver.tick();
assert_eq!(driver.status().generation, start + 3);
}
#[test]
fn status_snapshot_reflects_engine() {
let mut chain = ScriptedChain::new(1);
chain.push_block(&[1]);
chain.push_block(&[2]);
let mut driver = new_driver(chain, engine_config(0), DriverConfig::default());
driver.tick();
let status = driver.status();
assert_eq!(status.cursor, driver.engine().cursor());
assert_eq!(status.last_verified, driver.engine().last_verified());
assert_eq!(status.engine, driver.engine().status());
assert_eq!(status.caught_up, driver.is_caught_up());
assert_eq!(status.skips, driver.engine().skip_count());
}
#[test]
fn reorged_content_produces_typed_fork_then_recovery() {
let mut chain = ScriptedChain::new(1);
for value in 1..=6u64 {
chain.push_block(&[value]);
}
chain.set_batch_blocks(1);
let mut driver = Probed::new(
RecordingFold::default(),
chain,
engine_config(4),
DriverConfig::default(),
)
.unwrap();
driver.tick();
driver.tick();
driver.tick();
driver.checkpoint();
driver.tick();
driver.tick();
driver.tick();
driver.source_mut().reorg(3, &[&[40], &[50], &[60]]);
let outcome = driver.tick();
assert_eq!(
outcome,
Tick::RolledBack {
to: Some(Position::new(3, 0)),
}
);
run_to_idle(&mut driver);
let expected = vec![
(Position::new(1, 0), 1),
(Position::new(2, 0), 2),
(Position::new(3, 0), 3),
(Position::new(4, 0), 40),
(Position::new(5, 0), 50),
(Position::new(6, 0), 60),
];
assert_eq!(driver.engine().view(), expected);
}
#[test]
fn eventless_fork_point_is_still_detected() {
let mut chain = ScriptedChain::new(1);
chain.push_block(&[]);
chain.push_block(&[2]);
chain.push_block(&[]);
chain.push_block(&[]);
chain.push_block(&[]);
chain.push_block(&[]);
chain.push_block(&[7]);
chain.set_batch_blocks(1);
let mut driver = Probed::new(
RecordingFold::default(),
chain,
engine_config(2),
DriverConfig::default(),
)
.unwrap();
driver.tick();
driver.checkpoint();
driver.tick();
driver.source_mut().reorg(3, &[&[], &[], &[70]]);
let outcome = driver.tick();
assert!(matches!(outcome, Tick::RolledBack { .. }));
run_to_idle(&mut driver);
let expected = vec![(Position::new(2, 0), 2), (Position::new(7, 0), 70)];
assert_eq!(driver.engine().view(), expected);
}
#[test]
fn shorter_chain_fork_is_suspected_not_retried() {
let mut chain = ScriptedChain::new(1);
for value in 1..=5u64 {
chain.push_block(&[value]);
}
chain.set_batch_blocks(1);
let mut driver = Probed::new(
RecordingFold::default(),
chain,
engine_config(2),
DriverConfig::default(),
)
.unwrap();
driver.tick();
driver.checkpoint();
for _ in 0..4 {
driver.tick();
}
driver.source_mut().reorg(4, &[&[99]]);
let outcome = driver.tick();
assert_ne!(outcome, Tick::SourceError);
assert!(matches!(outcome, Tick::RolledBack { .. }));
run_to_idle(&mut driver);
let expected = vec![(Position::new(1, 0), 1), (Position::new(2, 0), 99)];
assert_eq!(driver.engine().view(), expected);
}
#[test]
fn bisection_finds_deepest_canonical_block() {
let mut chain = ScriptedChain::new(1);
for value in 1..=8u64 {
chain.push_block(&[value]);
}
chain.set_batch_blocks(1);
let mut driver = Probed::new(
RecordingFold::default(),
Probe::new(chain),
engine_config(2),
DriverConfig::default(),
)
.unwrap();
driver.tick();
driver.checkpoint();
for _ in 0..7 {
driver.tick();
}
driver.source_mut().inner.reorg(3, &[&[60], &[70], &[80]]);
let outcome = driver.tick();
match outcome {
Tick::RolledBack { to } => {
let landed = to.map_or(0, |pos| pos.block);
assert!(landed <= 5);
}
other => panic!("expected RolledBack, got {other:?}"),
}
assert!(driver.source_mut().calls <= 4);
}
#[test]
fn fork_deeper_than_ring_escalates() {
fn build(horizon: ReplayHorizon) -> Probed<RecordingFold, ScriptedChain> {
let mut chain = ScriptedChain::new(1);
for value in 1..=6u64 {
chain.push_block(&[value]);
}
chain.set_batch_blocks(1);
let mut driver = Probed::new(
RecordingFold::default(),
chain,
EngineConfig {
ring_capacity: 4,
checkpoint_slots: 0,
},
DriverConfig::default(),
)
.unwrap();
run_to_idle(&mut driver);
driver
.source_mut()
.reorg(6, &[&[10], &[20], &[30], &[40], &[50], &[60]]);
driver.source_mut().set_horizon(horizon);
driver
}
let mut resyncable = build(ReplayHorizon::Genesis);
let resync_outcome = resyncable.tick();
let mut terminal = build(ReplayHorizon::FromBlock(1));
let terminal_outcome = terminal.tick();
assert_eq!(resync_outcome, Tick::Resynced);
assert_eq!(
terminal_outcome,
Tick::Terminal(EngineStatus::Unrecoverable {
cause: DivergenceCause::HorizonExceeded {
needed: 0,
horizon: 1,
},
})
);
}
#[test]
fn no_checkpoint_below_ancestor_escalates() {
fn build(horizon: ReplayHorizon) -> Probed<RecordingFold, ScriptedChain> {
let mut chain = ScriptedChain::new(1);
for value in 1..=8u64 {
chain.push_block(&[value]);
}
chain.set_batch_blocks(1);
let mut driver = Probed::new(
RecordingFold::default(),
chain,
engine_config(1),
DriverConfig::default(),
)
.unwrap();
for _ in 0..7 {
driver.tick();
}
driver.checkpoint();
driver.tick();
driver.source_mut().reorg(3, &[&[60], &[70], &[80]]);
driver.source_mut().set_horizon(horizon);
driver
}
let mut resyncable = build(ReplayHorizon::Genesis);
let resync_outcome = resyncable.tick();
let mut terminal = build(ReplayHorizon::FromBlock(1));
let terminal_outcome = terminal.tick();
assert_eq!(resync_outcome, Tick::Resynced);
assert_eq!(
terminal_outcome,
Tick::Terminal(EngineStatus::Unrecoverable {
cause: DivergenceCause::HorizonExceeded {
needed: 0,
horizon: 1,
},
})
);
}
#[test]
fn probe_failure_retries_without_state_damage() {
let mut chain = ScriptedChain::new(1);
for value in 1..=6u64 {
chain.push_block(&[value]);
}
chain.set_batch_blocks(1);
let mut driver = Probed::new(
RecordingFold::default(),
Probe::new(chain),
engine_config(2),
DriverConfig::default(),
)
.unwrap();
driver.tick();
driver.checkpoint();
for _ in 0..5 {
driver.tick();
}
driver.source_mut().inner.reorg(3, &[&[40], &[50], &[60]]);
driver.source_mut().fail_next_probes(1);
let first = driver.tick();
assert_eq!(first, Tick::SourceError);
let second = driver.tick();
assert!(matches!(second, Tick::RolledBack { .. }));
run_to_idle(&mut driver);
let expected = vec![
(Position::new(1, 0), 1),
(Position::new(2, 0), 2),
(Position::new(3, 0), 3),
(Position::new(4, 0), 40),
(Position::new(5, 0), 50),
(Position::new(6, 0), 60),
];
assert_eq!(driver.engine().view(), expected);
}
#[test]
fn anchor_match_resets_divergence_budget() {
let mut chain = ScriptedChain::new(1);
for value in 1..=5u64 {
chain.push_block(&[value]);
}
let mut driver = Driver::with_anchor(
RecordingFold::default(),
chain,
ExactAnchor,
engine_config(2),
DriverConfig::default(),
)
.unwrap();
driver.tick();
let cursor_before = driver.engine().cursor();
driver.checkpoint();
driver.checkpoint();
assert_eq!(driver.engine().cursor(), cursor_before);
assert_eq!(driver.engine().status(), EngineStatus::Active);
assert_eq!(driver.engine().checkpoint_count(), 2);
}
#[test]
fn anchor_divergence_rolls_back_then_terminates() {
let mut chain = ScriptedChain::new(1);
for value in 1..=3u64 {
chain.push_block(&[value]);
}
let mut driver = Driver::with_anchor(
RecordingFold::default(),
chain,
DisagreeingAnchor,
EngineConfig {
ring_capacity: 8,
checkpoint_slots: 4,
},
DriverConfig {
max_divergence_retries: 1,
..DriverConfig::default()
},
)
.unwrap();
driver.checkpoint();
driver.tick();
driver.checkpoint();
let after_first_mismatch = (driver.engine().status(), driver.engine().cursor());
driver.tick();
driver.checkpoint();
let after_second_mismatch = driver.engine().status();
assert_eq!(after_first_mismatch, (EngineStatus::Active, None));
assert_eq!(
after_second_mismatch,
EngineStatus::Unrecoverable {
cause: DivergenceCause::AnchorDivergence { at: 3 },
}
);
}
}