#[cfg(not(feature = "std"))]
use alloc::vec::Vec;
#[cfg(feature = "std")]
use std::vec::Vec;
use core::time::Duration;
use crate::{
batch::Batch,
engine::{
ApplySummary,
Engine,
EngineConfig,
},
error::{
ApplyError,
ConfigError,
DivergenceCause,
DurabilityLost,
EngineStatus,
RollbackError,
},
fold::Fold,
position::{
BlockRef,
Position,
},
sink::{
NoSink,
SnapshotSink,
},
source::{
ReplayHorizon,
Source,
},
};
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_CHECKPOINT_INTERVAL: u64 = 64;
fn due(last: Option<u64>, block: u64, interval: u64) -> bool {
last.is_none_or(|last| block >= last.saturating_add(interval))
}
#[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>,
}
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: Some(DEFAULT_CHECKPOINT_INTERVAL),
snapshot_interval: Some(DEFAULT_CHECKPOINT_INTERVAL),
}
}
}
impl DriverConfig {
pub fn from_block(start_block: u64) -> Self {
Self {
start_block,
..Self::default()
}
}
}
pub struct Driver<F, S, K = NoSink>
where
F: Fold,
S: Source<Event = F::Event>,
K: SnapshotSink<F>,
{
engine: Engine<F>,
source: S,
sink: K,
config: DriverConfig,
batch: Batch<F::Event>,
scratch: Vec<(BlockRef, u32, F::Event)>,
scanned_to: Option<u64>,
initial: F,
consecutive_errors: u32,
caught_up: bool,
generation: u64,
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: Source<Event = F::Event>,
{
pub fn new(
fold: F,
source: S,
engine: EngineConfig,
config: DriverConfig,
) -> Result<Self, ConfigError> {
Self::build(fold, source, NoSink, engine, config)
}
pub fn resume(
engine: Engine<F>,
source: S,
genesis: F,
config: DriverConfig,
) -> Result<Self, ConfigError> {
Self::around(engine, source, NoSink, genesis, config)
}
}
impl<F, S, K> Driver<F, S, K>
where
F: Fold + Clone,
S: Source<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, 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, sink, genesis, config)
}
}
impl<F, S, K> Driver<F, S, K>
where
F: Fold + Clone,
S: Source<Event = F::Event>,
K: SnapshotSink<F>,
{
fn build(
fold: F,
source: S,
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, sink, initial, driver_config)
}
fn around(
engine: Engine<F>,
source: S,
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,
sink,
config: driver_config,
batch: Batch::new(),
scratch: Vec::new(),
scanned_to: None,
initial,
consecutive_errors: 0,
caught_up: false,
generation: 0,
last_checkpoint_block: None,
last_snapshot_block: None,
durability_lost: false,
advanced: false,
})
}
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) {
let Some(interval) = self.config.checkpoint_interval else {
return;
};
let Some(cursor) = self.engine.cursor() else {
return;
};
if due(self.last_checkpoint_block, cursor.block, interval) {
self.run_checkpoint();
}
}
fn run_checkpoint(&mut self) {
self.engine.checkpoint();
if let Some(cursor) = self.engine.cursor() {
self.last_checkpoint_block = Some(cursor.block);
}
}
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()?;
if !due(self.last_snapshot_block, point.block, interval) {
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));
}
#[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.scanned_to = None;
self.last_checkpoint_block = None;
self.last_snapshot_block = None;
Tick::Resynced
}
fn step(&mut self) -> Tick {
let tick = self.poll_apply();
self.advanced = matches!(
tick,
Tick::Progressed(summary) if summary.applied > 0 || summary.skipped > 0
);
tick
}
fn poll_apply(&mut self) -> Tick {
self.generation = self.generation.wrapping_add(1);
if !self.engine.status().is_active() {
return Tick::Terminal(self.engine.status());
}
if self.scan().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();
self.auto_checkpoint();
if let Some(tick) = self.offer_snapshot() {
return tick;
}
if self.batch.is_empty() {
Tick::Idle
} else {
Tick::Progressed(summary)
}
}
Err(
ApplyError::ForkSuspected { .. }
| ApplyError::MissingBoundary
| ApplyError::CursorBlockUnobserved { .. },
) => self.recover_via_bisection(),
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"
)
}
}
}
fn scan(&mut self) -> Result<(), S::Error> {
self.batch.clear();
let cursor = self.engine.cursor();
self.batch.boundary = match cursor {
Some(cursor) => self.source.header_at(cursor.block)?,
None => None,
};
let head = self.source.head()?;
let mut from = match cursor {
Some(cursor) => self
.scanned_to
.map_or(cursor.block + 1, |to| to.saturating_add(1))
.min(cursor.block + 1),
None => self.config.start_block,
};
let window = self.source.window().max(1);
while from <= head {
let to = head.min(from.saturating_add(window - 1));
self.scratch.clear();
self.source.events_in(from, to, &mut self.scratch)?;
self.scanned_to = Some(to);
from = to.saturating_add(1);
if !self.scratch.is_empty() {
group_into(&mut self.batch, &mut self.scratch);
break;
}
}
Ok(())
}
#[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 fn tick(&mut self) -> Tick {
self.step()
}
pub fn checkpoint(&mut self) {
self.run_checkpoint();
}
pub 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,
}
}
pub 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)
}
}
fn group_into<E>(batch: &mut Batch<E>, entries: &mut Vec<(BlockRef, u32, E)>) {
entries.sort_by_key(|(block, log_index, _)| (block.number, *log_index));
let mut span: Vec<(u32, E)> = Vec::new();
let mut current: Option<BlockRef> = None;
for (block, log_index, event) in entries.drain(..) {
match current {
Some(open) if open.number != block.number => {
batch.push_block(open, span.drain(..));
current = Some(block);
}
None => current = Some(block),
_ => {}
}
span.push((log_index, event));
}
if let Some(open) = current {
batch.push_block(open, span);
}
}
#[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::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 Source for Probe {
type Event = u64;
type Error = PollFailure;
fn head(&mut self) -> Result<u64, PollFailure> {
self.inner.head()
}
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)
}
fn events_in(
&mut self,
from: u64,
to: u64,
out: &mut Vec<(BlockRef, u32, u64)>,
) -> Result<(), PollFailure> {
self.inner.events_in(from, to, out)
}
fn horizon(&self) -> ReplayHorizon {
self.inner.horizon()
}
fn window(&self) -> u64 {
self.inner.window()
}
}
struct Stuck {
block: BlockRef,
}
impl Source for Stuck {
type Event = u64;
type Error = PollFailure;
fn head(&mut self) -> Result<u64, PollFailure> {
Ok(self.block.number + 1)
}
fn header_at(&mut self, _number: u64) -> Result<Option<BlockRef>, PollFailure> {
Ok(Some(self.block))
}
fn events_in(
&mut self,
_from: u64,
_to: u64,
out: &mut Vec<(BlockRef, u32, u64)>,
) -> Result<(), PollFailure> {
out.push((self.block, 0, 1));
Ok(())
}
}
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<F, S, K>(driver: &mut Driver<F, S, K>) -> Tick
where
F: Fold + Clone,
S: Source<Event = F::Event>,
K: SnapshotSink<F>,
{
let mut outcome = driver.tick();
while !matches!(outcome, Tick::Idle) {
outcome = driver.tick();
}
outcome
}
fn collect_to_idle<F, S, K>(driver: &mut Driver<F, S, K>) -> Vec<Tick>
where
F: Fold + Clone,
S: Source<Event = F::Event>,
K: SnapshotSink<F>,
{
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_window(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, 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) -> SinkDriver {
Driver::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().fold().applied, 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().fold().applied, 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 an_empty_window_does_not_report_caught_up_while_behind_head() {
let mut chain = ScriptedChain::new(1);
for _ in 1..10 {
chain.push_block(&[]);
}
chain.push_block(&[42]);
chain.set_window(1);
let mut driver = new_driver(chain, engine_config(0), DriverConfig::from_block(1));
let tick = driver.tick();
assert!(matches!(tick, Tick::Progressed(_)), "got {tick:?}");
assert!(!driver.is_caught_up());
assert_eq!(
driver.engine().fold().applied,
vec![(Position::new(10, 0), 42)]
);
}
#[test]
fn rollback_replays_from_the_restored_cursor_not_the_scan_mark() {
let mut chain = ScriptedChain::new(1);
for value in 1..=10u64 {
chain.push_block(&[value]);
}
chain.set_window(1);
let mut driver = new_driver(
chain,
engine_config(4),
DriverConfig {
checkpoint_interval: Some(1),
..DriverConfig::from_block(1)
},
);
run_to_idle(&mut driver);
driver.source_mut().reorg(3, &[&[80], &[90], &[100]]);
run_to_idle(&mut driver);
let applied: Vec<u64> = driver
.engine()
.fold()
.applied
.iter()
.map(|(_, event)| *event)
.collect();
assert_eq!(applied, vec![1, 2, 3, 4, 5, 6, 7, 80, 90, 100]);
}
#[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().fold().applied, 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_window(1);
let config = DriverConfig {
checkpoint_interval: Some(4),
..DriverConfig::default()
};
let engine = EngineConfig {
ring_capacity: 16,
checkpoint_slots: 8,
};
let mut driver = new_driver(chain, engine, config);
run_to_idle(&mut driver);
assert!(driver.engine().checkpoint_count() >= 3);
}
#[test]
fn checkpoints_expire_once_their_block_leaves_the_ring() {
let mut chain = ScriptedChain::new(1);
for value in 1..=12u64 {
chain.push_block(&[value]);
}
chain.set_window(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_eq!(driver.engine().checkpoint_count(), 2);
assert_eq!(driver.engine().durable_point(), Some(Position::new(5, 0)));
}
#[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_window(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_window(1);
let mut driver = Driver::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().fold().applied, 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_window(1);
let mut driver = Driver::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().fold().applied, 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_window(1);
let mut driver = Driver::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().fold().applied, 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_window(1);
let mut driver = Driver::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]]);
driver.source_mut().calls = 0;
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) -> Driver<RecordingFold, ScriptedChain> {
let mut chain = ScriptedChain::new(1);
for value in 1..=6u64 {
chain.push_block(&[value]);
}
chain.set_window(1);
let mut driver = Driver::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) -> Driver<RecordingFold, ScriptedChain> {
let mut chain = ScriptedChain::new(1);
for value in 1..=8u64 {
chain.push_block(&[value]);
}
chain.set_window(1);
let mut driver = Driver::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_window(1);
let mut driver = Driver::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().fold().applied, expected);
}
}