use reifydb_catalog::catalog::Catalog;
use reifydb_codec::key::encoded::EncodedKey;
use reifydb_core::{actors::pending::PendingLayers, interface::catalog::flow::OperatorId, state::timer::TimerKind};
use reifydb_flow::{
timer::{
Timer, TimerDue,
wheel::{DueTimers, MAX_TIMERS_PER_SCAN, TimerWheel},
},
transaction::{
DeferredParams, FlowTransaction,
deferred::DeferredTransaction,
substrate::{FlowSubstrate, apply_operator_state},
},
};
use reifydb_runtime::context::clock::{Clock, MockClock};
use reifydb_test_harness::engine::TestEngine;
use reifydb_transaction::interceptor::interceptors::Interceptors;
use reifydb_value::{factory::time::at_millis, value::identity::IdentityId};
const NODE: OperatorId = OperatorId(1);
const NO_LIMIT: usize = usize::MAX;
fn deferred(engine: &TestEngine) -> DeferredTransaction {
deferred_with_clock(engine, MockClock::from_millis(0))
}
fn deferred_with_clock(engine: &TestEngine, clock: MockClock) -> DeferredTransaction {
let parent = engine.begin_admin(IdentityId::system()).unwrap();
let version = parent.version();
DeferredTransaction::new(DeferredParams {
version,
pending: PendingLayers::empty(),
query: Some(parent.multi.begin_query().unwrap()),
state_query: Some(parent.multi.begin_query().unwrap()),
catalog: Catalog::testing(),
interceptors: Interceptors::new(),
clock: Clock::Mock(clock),
substrate: FlowSubstrate::with_dictionary(
engine.inner().dictionary_allocators(),
engine.inner().operator_state(),
),
})
}
fn commit_pending(engine: &TestEngine, txn: &mut impl FlowTransaction) {
let pending = txn.take_pending();
apply_operator_state(&engine.inner().operator_state(), &pending);
}
fn timer(millis: u64, kind: TimerKind, key: &str) -> Timer {
Timer {
due: at_millis(millis),
kind,
key: EncodedKey::new(key.as_bytes()),
}
}
#[test]
fn a_timer_is_due_exactly_when_the_watermark_reaches_it() {
let engine = TestEngine::new();
let mut txn = deferred(&engine);
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Seal, "bucket")).unwrap();
assert!(
TimerWheel::take_due(NODE, &mut txn, at_millis(4_999), NO_LIMIT, None).unwrap().timers.is_empty(),
"must not fire early"
);
assert_eq!(
TimerWheel::take_due(NODE, &mut txn, at_millis(5_000), NO_LIMIT, None).unwrap().timers,
vec![timer(5_000, TimerKind::Seal, "bucket")]
);
}
#[test]
fn rearming_a_unique_kind_moves_its_deadline_instead_of_minting_a_second_timer() {
let engine = TestEngine::new();
let mut txn = deferred(&engine);
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Maintenance, "m")).unwrap();
TimerWheel::arm(NODE, &mut txn, &timer(6_000, TimerKind::Maintenance, "m")).unwrap();
TimerWheel::arm(NODE, &mut txn, &timer(7_000, TimerKind::Maintenance, "m")).unwrap();
assert_eq!(
TimerWheel::take_due(NODE, &mut txn, at_millis(10_000), NO_LIMIT, None).unwrap().timers,
vec![timer(7_000, TimerKind::Maintenance, "m")],
"re-arming must move the one deadline, not leave the superseded instants armed"
);
}
#[test]
fn a_backlog_kind_still_holds_every_instant_it_was_armed_at() {
let engine = TestEngine::new();
let mut txn = deferred(&engine);
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Seal, "m")).unwrap();
TimerWheel::arm(NODE, &mut txn, &timer(6_000, TimerKind::Seal, "m")).unwrap();
assert_eq!(
TimerWheel::take_due(NODE, &mut txn, at_millis(10_000), NO_LIMIT, None).unwrap().timers,
vec![timer(5_000, TimerKind::Seal, "m"), timer(6_000, TimerKind::Seal, "m")],
"a backlog kind must keep every bucket it was armed for"
);
}
#[test]
fn a_fired_unique_timer_can_be_armed_again_at_a_later_instant() {
let engine = TestEngine::new();
let mut txn = deferred(&engine);
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Maintenance, "m")).unwrap();
assert_eq!(
TimerWheel::take_due(NODE, &mut txn, at_millis(5_000), NO_LIMIT, None).unwrap().timers,
vec![timer(5_000, TimerKind::Maintenance, "m")]
);
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Maintenance, "m")).unwrap();
assert_eq!(
TimerWheel::take_due(NODE, &mut txn, at_millis(10_000), NO_LIMIT, None).unwrap().timers,
vec![timer(5_000, TimerKind::Maintenance, "m")],
"re-arming the instant that just fired must arm a live timer, not be skipped as a duplicate"
);
}
#[test]
fn due_timers_return_in_at_then_kind_then_id_order() {
let engine = TestEngine::new();
let mut txn = deferred(&engine);
TimerWheel::arm(NODE, &mut txn, &timer(7_000, TimerKind::Grace, "a")).unwrap();
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Grace, "a")).unwrap();
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Seal, "z")).unwrap();
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Seal, "a")).unwrap();
assert_eq!(
TimerWheel::take_due(NODE, &mut txn, at_millis(10_000), NO_LIMIT, None).unwrap().timers,
vec![
timer(5_000, TimerKind::Seal, "z"),
timer(5_000, TimerKind::Seal, "a"),
timer(5_000, TimerKind::Grace, "a"),
timer(7_000, TimerKind::Grace, "a"),
]
);
}
#[test]
fn arming_the_same_timer_twice_fires_once() {
let engine = TestEngine::new();
let clock = MockClock::from_millis(0);
let mut txn = deferred_with_clock(&engine, clock.clone());
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Grace, "group")).unwrap();
clock.advance_millis(250);
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Grace, "group")).unwrap();
assert_eq!(TimerWheel::take_due(NODE, &mut txn, at_millis(5_000), NO_LIMIT, None).unwrap().timers.len(), 1);
}
#[test]
fn a_capped_take_drains_the_earliest_first_and_leaves_the_rest_armed() {
let engine = TestEngine::new();
let mut txn = deferred(&engine);
for at_ms in [9_000u64, 5_000, 7_000] {
TimerWheel::arm(NODE, &mut txn, &timer(at_ms, TimerKind::Seal, "b")).unwrap();
}
assert_eq!(
TimerWheel::take_due(NODE, &mut txn, at_millis(10_000), 2, None).unwrap().timers,
vec![timer(5_000, TimerKind::Seal, "b"), timer(7_000, TimerKind::Seal, "b")],
"a capped take must drain in firing order, earliest first"
);
assert_eq!(
TimerWheel::take_due(NODE, &mut txn, at_millis(10_000), NO_LIMIT, None).unwrap().timers,
vec![timer(9_000, TimerKind::Seal, "b")],
"what the cap left behind must still be armed for the next round"
);
}
#[test]
fn a_disarmed_timer_does_not_fire_and_its_replacement_does() {
let engine = TestEngine::new();
let mut txn = deferred(&engine);
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Seal, "session")).unwrap();
TimerWheel::disarm(NODE, &mut txn, &timer(5_000, TimerKind::Seal, "session")).unwrap();
TimerWheel::arm(NODE, &mut txn, &timer(8_000, TimerKind::Seal, "session")).unwrap();
assert_eq!(
TimerWheel::take_due(NODE, &mut txn, at_millis(9_000), NO_LIMIT, None).unwrap().timers,
vec![timer(8_000, TimerKind::Seal, "session")],
"the superseded instant must not fire and the re-armed one must"
);
}
#[test]
fn disarming_either_end_of_the_wheel_leaves_the_other_timer_firing() {
let engine = TestEngine::new();
let mut txn = deferred(&engine);
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Seal, "a")).unwrap();
TimerWheel::arm(NODE, &mut txn, &timer(8_000, TimerKind::Seal, "b")).unwrap();
TimerWheel::disarm(NODE, &mut txn, &timer(8_000, TimerKind::Seal, "b")).unwrap();
assert_eq!(
TimerWheel::take_due(NODE, &mut txn, at_millis(9_000), NO_LIMIT, None).unwrap().timers,
vec![timer(5_000, TimerKind::Seal, "a")],
"disarming the later instant must leave the earlier one armed"
);
let mut txn = deferred(&engine);
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Seal, "a")).unwrap();
TimerWheel::arm(NODE, &mut txn, &timer(8_000, TimerKind::Seal, "b")).unwrap();
TimerWheel::disarm(NODE, &mut txn, &timer(5_000, TimerKind::Seal, "a")).unwrap();
assert_eq!(
TimerWheel::take_due(NODE, &mut txn, at_millis(9_000), NO_LIMIT, None).unwrap().timers,
vec![timer(8_000, TimerKind::Seal, "b")],
"disarming the earliest instant must leave the later one armed"
);
}
#[test]
fn disarming_by_key_cancels_the_instant_the_index_names_and_spares_every_other_key() {
let engine = TestEngine::new();
let mut txn = deferred(&engine);
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Maintenance, "emptied")).unwrap();
TimerWheel::arm(NODE, &mut txn, &timer(8_000, TimerKind::Maintenance, "emptied")).unwrap();
TimerWheel::arm(NODE, &mut txn, &timer(6_000, TimerKind::Maintenance, "neighbour")).unwrap();
TimerWheel::disarm_by_key(NODE, &mut txn, TimerKind::Maintenance, &EncodedKey::new(b"emptied")).unwrap();
assert_eq!(
TimerWheel::take_due(NODE, &mut txn, at_millis(10_000), NO_LIMIT, None).unwrap().timers,
vec![timer(6_000, TimerKind::Maintenance, "neighbour")],
"only the disarmed key's armed instant may go"
);
}
#[test]
fn a_key_disarmed_by_key_can_be_armed_again_at_the_very_instant_it_held() {
let engine = TestEngine::new();
let mut txn = deferred(&engine);
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Maintenance, "refilled")).unwrap();
TimerWheel::disarm_by_key(NODE, &mut txn, TimerKind::Maintenance, &EncodedKey::new(b"refilled")).unwrap();
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Maintenance, "refilled")).unwrap();
assert_eq!(
TimerWheel::take_due(NODE, &mut txn, at_millis(10_000), NO_LIMIT, None).unwrap().timers,
vec![timer(5_000, TimerKind::Maintenance, "refilled")],
"a re-arm after a disarm by key must arm a live timer, not be swallowed as a duplicate"
);
}
#[test]
fn disarming_an_unarmed_key_leaves_the_wheel_untouched() {
let engine = TestEngine::new();
let mut txn = deferred(&engine);
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Maintenance, "armed")).unwrap();
TimerWheel::disarm_by_key(NODE, &mut txn, TimerKind::Maintenance, &EncodedKey::new(b"never-armed")).unwrap();
assert_eq!(
TimerWheel::take_due(NODE, &mut txn, at_millis(10_000), NO_LIMIT, None).unwrap().timers,
vec![timer(5_000, TimerKind::Maintenance, "armed")]
);
}
#[test]
fn a_restart_does_not_fire_a_timer_disarmed_by_key() {
let engine = TestEngine::new();
let mut txn = deferred(&engine);
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Maintenance, "emptied")).unwrap();
TimerWheel::disarm_by_key(NODE, &mut txn, TimerKind::Maintenance, &EncodedKey::new(b"emptied")).unwrap();
commit_pending(&engine, &mut txn);
let mut cold_txn = deferred(&engine);
assert!(TimerWheel::take_due(NODE, &mut cold_txn, at_millis(9_000), NO_LIMIT, None).unwrap().timers.is_empty());
}
#[test]
fn a_restart_does_not_fire_a_disarmed_timer() {
let engine = TestEngine::new();
let mut txn = deferred(&engine);
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Seal, "session")).unwrap();
TimerWheel::disarm(NODE, &mut txn, &timer(5_000, TimerKind::Seal, "session")).unwrap();
commit_pending(&engine, &mut txn);
let mut cold_txn = deferred(&engine);
assert!(
TimerWheel::take_due(NODE, &mut cold_txn, at_millis(9_000), NO_LIMIT, None).unwrap().timers.is_empty(),
"a disarm that only lived in RAM lets the superseded timer survive the restart"
);
}
#[test]
fn take_due_removes_what_it_returns_and_keeps_the_rest() {
let engine = TestEngine::new();
let mut txn = deferred(&engine);
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Seal, "due")).unwrap();
TimerWheel::arm(NODE, &mut txn, &timer(9_000, TimerKind::Seal, "later")).unwrap();
assert_eq!(TimerWheel::take_due(NODE, &mut txn, at_millis(6_000), NO_LIMIT, None).unwrap().timers.len(), 1);
assert!(TimerWheel::take_due(NODE, &mut txn, at_millis(6_000), NO_LIMIT, None).unwrap().timers.is_empty());
assert_eq!(
TimerWheel::take_due(NODE, &mut txn, at_millis(9_000), NO_LIMIT, None).unwrap().timers,
vec![timer(9_000, TimerKind::Seal, "later")]
);
}
#[test]
fn a_restart_still_fires_persisted_timers() {
let engine = TestEngine::new();
let mut txn = deferred(&engine);
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Seal, "bucket")).unwrap();
commit_pending(&engine, &mut txn);
let mut cold_txn = deferred(&engine);
assert_eq!(
TimerWheel::take_due(NODE, &mut cold_txn, at_millis(5_000), NO_LIMIT, None).unwrap().timers,
vec![timer(5_000, TimerKind::Seal, "bucket")]
);
}
fn due(millis: u64) -> TimerDue {
TimerDue {
operator_id: NODE,
due: at_millis(millis),
}
}
#[test]
fn a_take_reports_the_earliest_instant_it_left_behind_not_the_latest() {
let engine = TestEngine::new();
let mut txn = deferred(&engine);
for at_ms in [9_000u64, 5_000, 7_000] {
TimerWheel::arm(NODE, &mut txn, &timer(at_ms, TimerKind::Seal, "b")).unwrap();
}
let DueTimers {
timers: fired,
next,
..
} = TimerWheel::take_due(NODE, &mut txn, at_millis(0), NO_LIMIT, None).unwrap();
assert!(fired.is_empty(), "a watermark below every armed instant must fire nothing");
assert_eq!(next, Some(at_millis(5_000)));
}
#[test]
fn a_take_reports_nothing_armed_once_the_wheel_is_drained() {
let engine = TestEngine::new();
let mut txn = deferred(&engine);
assert_eq!(TimerWheel::take_due(NODE, &mut txn, at_millis(5_000), NO_LIMIT, None).unwrap().next, None);
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Seal, "b")).unwrap();
let DueTimers {
timers: fired,
next,
..
} = TimerWheel::take_due(NODE, &mut txn, at_millis(5_000), NO_LIMIT, None).unwrap();
assert_eq!(fired.len(), 1);
assert_eq!(next, None, "a wheel drained by the take itself must report nothing armed");
}
#[test]
fn a_take_reports_an_instant_no_watermark_has_reached() {
let engine = TestEngine::new();
let mut txn = deferred(&engine);
TimerWheel::arm(NODE, &mut txn, &timer(10_000_000_000, TimerKind::Seal, "distant")).unwrap();
let DueTimers {
timers: fired,
next,
..
} = TimerWheel::take_due(NODE, &mut txn, at_millis(9_000), NO_LIMIT, None).unwrap();
assert!(fired.is_empty());
assert_eq!(next, Some(at_millis(10_000_000_000)));
}
#[test]
fn a_capped_take_reports_a_leftover_that_is_already_due() {
let engine = TestEngine::new();
let mut txn = deferred(&engine);
for at_ms in [9_000u64, 5_000, 7_000] {
TimerWheel::arm(NODE, &mut txn, &timer(at_ms, TimerKind::Seal, "b")).unwrap();
}
let DueTimers {
timers: fired,
next,
..
} = TimerWheel::take_due(NODE, &mut txn, at_millis(10_000), 2, None).unwrap();
assert_eq!(fired.len(), 2);
assert_eq!(next, Some(at_millis(9_000)));
assert!(
next.is_some_and(|due| due <= at_millis(10_000)),
"a capped take must report its leftover as still due"
);
let DueTimers {
timers: rest,
next,
..
} = TimerWheel::take_due(NODE, &mut txn, at_millis(10_000), NO_LIMIT, None).unwrap();
assert_eq!(rest, vec![timer(9_000, TimerKind::Seal, "b")], "the next round must fire what the cap left behind");
assert_eq!(next, None);
}
#[test]
fn next_due_stored_ignores_an_arm_that_has_not_been_committed() {
let engine = TestEngine::new();
let mut txn = deferred(&engine);
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Seal, "b")).unwrap();
assert_eq!(TimerWheel::next_due_stored(NODE, &engine.inner().operator_state()), None);
commit_pending(&engine, &mut txn);
assert_eq!(TimerWheel::next_due_stored(NODE, &engine.inner().operator_state()), Some(due(5_000)));
}
#[test]
fn next_due_stored_reports_the_earliest_committed_instant() {
let engine = TestEngine::new();
let mut txn = deferred(&engine);
TimerWheel::arm(NODE, &mut txn, &timer(9_000, TimerKind::Seal, "b")).unwrap();
TimerWheel::arm(NODE, &mut txn, &timer(5_000, TimerKind::Seal, "b")).unwrap();
TimerWheel::arm(NODE, &mut txn, &timer(7_000, TimerKind::Seal, "b")).unwrap();
commit_pending(&engine, &mut txn);
assert_eq!(TimerWheel::next_due_stored(NODE, &engine.inner().operator_state()), Some(due(5_000)));
}
#[test]
fn an_uncapped_take_still_bounds_what_it_pulls_and_leaves_the_rest_armed() {
let engine = TestEngine::new();
let mut txn = deferred(&engine);
const ARMED: u64 = MAX_TIMERS_PER_SCAN as u64 + 8;
for step in 0..ARMED {
TimerWheel::arm(NODE, &mut txn, &timer(1_000 + step, TimerKind::Seal, "b")).unwrap();
}
let DueTimers {
timers: first,
next,
..
} = TimerWheel::take_due(NODE, &mut txn, at_millis(10_000), NO_LIMIT, None).unwrap();
assert!(
(first.len() as u64) < ARMED,
"an uncapped budget must not turn into an unbounded scan, took {} of {}",
first.len(),
ARMED
);
assert_eq!(
next,
Some(at_millis(1_000 + first.len() as u64)),
"the bound must report the earliest instant it left behind, not none"
);
let mut drained: Vec<u64> = first.iter().map(|timer| timer.due.to_millis() as u64).collect();
loop {
let batch = TimerWheel::take_due(NODE, &mut txn, at_millis(10_000), NO_LIMIT, None).unwrap().timers;
if batch.is_empty() {
break;
}
drained.extend(batch.iter().map(|timer| timer.due.to_millis() as u64));
}
assert_eq!(
drained,
(0..ARMED).map(|step| 1_000 + step).collect::<Vec<u64>>(),
"successive takes must reach every armed instant exactly once, in firing order"
);
}