use std::collections::HashMap;
use reifydb_core::interface::catalog::flow::{FlowId, OperatorId};
use reifydb_value::value::datetime::DateTime;
use crate::timer::TimerDue;
#[derive(Default)]
pub struct TimerRegistry {
flows: HashMap<FlowId, HashMap<OperatorId, DateTime>>,
}
impl TimerRegistry {
pub fn due_before(&mut self, armed: Vec<TimerDue>, flow: FlowId, watermark: DateTime) -> Vec<TimerDue> {
self.fold(flow, armed);
let Some(operators) = self.flows.get(&flow) else {
return Vec::new();
};
if operators.values().min().is_none_or(|earliest| *earliest > watermark) {
return Vec::new();
}
operators
.iter()
.filter(|(_, due)| **due <= watermark)
.map(|(operator_id, due)| TimerDue {
operator_id: *operator_id,
due: *due,
})
.collect()
}
pub fn refresh(&mut self, flow: FlowId, operator: OperatorId, next: Option<TimerDue>) {
let Some(operators) = self.flows.get_mut(&flow) else {
return;
};
match next {
Some(next) => {
operators.insert(operator, next.due);
}
None => {
operators.remove(&operator);
}
}
}
pub fn rebuild(&mut self, flow: FlowId, armed: Vec<TimerDue>) {
self.flows.remove(&flow);
self.fold(flow, armed);
}
pub fn remove_operator(&mut self, flow: FlowId, operator: OperatorId) {
let Some(operators) = self.flows.get_mut(&flow) else {
return;
};
operators.remove(&operator);
if operators.is_empty() {
self.flows.remove(&flow);
}
}
pub fn remove_flow(&mut self, flow: FlowId) {
self.flows.remove(&flow);
}
pub fn clear(&mut self) {
self.flows.clear();
}
fn fold(&mut self, flow: FlowId, armed: Vec<TimerDue>) {
if armed.is_empty() {
return;
}
let operators = self.flows.entry(flow).or_default();
for entry in armed {
operators
.entry(entry.operator_id)
.and_modify(|earliest| *earliest = (*earliest).min(entry.due))
.or_insert(entry.due);
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn operator(id: u64) -> OperatorId {
OperatorId(id)
}
fn flow(id: u64) -> FlowId {
FlowId(id)
}
fn armed(operator_id: u64, millis: u64) -> TimerDue {
TimerDue {
operator_id: operator(operator_id),
due: DateTime::from_millis(millis),
}
}
fn entries(registry: &TimerRegistry, flow: FlowId) -> Vec<TimerDue> {
let Some(operators) = registry.flows.get(&flow) else {
return Vec::new();
};
let mut held: Vec<TimerDue> = operators
.iter()
.map(|(operator_id, due)| TimerDue {
operator_id: *operator_id,
due: *due,
})
.collect();
held.sort();
held
}
fn due_before(registry: &mut TimerRegistry, armed: Vec<TimerDue>, flow: FlowId, millis: u64) -> Vec<TimerDue> {
let mut found = registry.due_before(armed, flow, DateTime::from_millis(millis));
found.sort();
found
}
#[test]
fn folding_an_arm_keeps_the_earliest_instant_for_an_operator() {
let mut registry = TimerRegistry::default();
due_before(&mut registry, vec![armed(1, 9_000), armed(1, 5_000)], flow(1), 0);
due_before(&mut registry, vec![armed(2, 5_000), armed(2, 9_000)], flow(1), 0);
assert_eq!(entries(®istry, flow(1)), vec![armed(1, 5_000), armed(2, 5_000)]);
}
#[test]
fn a_folded_arm_stays_armed_until_it_is_due() {
let mut registry = TimerRegistry::default();
assert!(due_before(&mut registry, vec![armed(1, 5_000)], flow(1), 0).is_empty());
assert_eq!(due_before(&mut registry, Vec::new(), flow(1), 5_000), vec![armed(1, 5_000)]);
}
#[test]
fn an_entry_beyond_the_watermark_is_withheld_but_kept() {
let mut registry = TimerRegistry::default();
assert!(due_before(&mut registry, vec![armed(1, 9_000)], flow(1), 5_000).is_empty());
assert_eq!(entries(®istry, flow(1)), vec![armed(1, 9_000)]);
}
#[test]
fn an_entry_exactly_at_the_watermark_is_due() {
let mut registry = TimerRegistry::default();
assert_eq!(due_before(&mut registry, vec![armed(1, 5_000)], flow(1), 5_000), vec![armed(1, 5_000)]);
}
#[test]
fn an_operator_past_the_watermark_is_withheld_while_a_sibling_fires() {
let mut registry = TimerRegistry::default();
let found = due_before(&mut registry, vec![armed(1, 5_000), armed(2, 9_000)], flow(1), 5_000);
assert_eq!(found, vec![armed(1, 5_000)]);
assert_eq!(entries(®istry, flow(1)), vec![armed(1, 5_000), armed(2, 9_000)]);
}
#[test]
fn one_flow_never_surfaces_another_flows_operators() {
let mut registry = TimerRegistry::default();
due_before(&mut registry, vec![armed(1, 5_000)], flow(1), 0);
due_before(&mut registry, vec![armed(2, 5_000)], flow(2), 0);
assert_eq!(due_before(&mut registry, Vec::new(), flow(1), 5_000), vec![armed(1, 5_000)]);
assert_eq!(due_before(&mut registry, Vec::new(), flow(2), 5_000), vec![armed(2, 5_000)]);
}
#[test]
fn a_flow_that_has_never_armed_anything_yields_no_candidates() {
let mut registry = TimerRegistry::default();
assert!(due_before(&mut registry, Vec::new(), flow(1), 5_000).is_empty());
assert_eq!(entries(®istry, flow(1)), Vec::new());
}
#[test]
fn refresh_overwrites_an_entry_with_its_next_instant() {
let mut registry = TimerRegistry::default();
due_before(&mut registry, vec![armed(1, 5_000)], flow(1), 0);
registry.refresh(flow(1), operator(1), Some(armed(1, 9_000)));
assert!(due_before(&mut registry, Vec::new(), flow(1), 5_000).is_empty());
assert_eq!(entries(®istry, flow(1)), vec![armed(1, 9_000)]);
}
#[test]
fn refresh_with_no_next_timer_drops_the_entry() {
let mut registry = TimerRegistry::default();
due_before(&mut registry, vec![armed(1, 5_000)], flow(1), 0);
registry.refresh(flow(1), operator(1), None);
assert_eq!(entries(®istry, flow(1)), Vec::new());
}
#[test]
fn rebuild_replaces_a_flows_entries_rather_than_merging_them() {
let mut registry = TimerRegistry::default();
due_before(&mut registry, vec![armed(1, 5_000)], flow(1), 0);
registry.rebuild(flow(1), vec![armed(2, 7_000)]);
assert_eq!(entries(®istry, flow(1)), vec![armed(2, 7_000)]);
}
#[test]
fn rebuilding_with_nothing_armed_clears_the_flow() {
let mut registry = TimerRegistry::default();
due_before(&mut registry, vec![armed(1, 5_000)], flow(1), 0);
registry.rebuild(flow(1), Vec::new());
assert_eq!(entries(®istry, flow(1)), Vec::new());
}
#[test]
fn removing_an_operator_leaves_its_siblings_armed() {
let mut registry = TimerRegistry::default();
due_before(&mut registry, vec![armed(1, 5_000), armed(2, 5_000)], flow(1), 0);
registry.remove_operator(flow(1), operator(1));
assert_eq!(entries(®istry, flow(1)), vec![armed(2, 5_000)]);
}
#[test]
fn removing_a_flow_drops_every_entry_it_held() {
let mut registry = TimerRegistry::default();
due_before(&mut registry, vec![armed(1, 5_000), armed(2, 5_000)], flow(1), 0);
due_before(&mut registry, vec![armed(3, 5_000)], flow(2), 0);
registry.remove_flow(flow(1));
assert_eq!(entries(®istry, flow(1)), Vec::new());
assert_eq!(entries(®istry, flow(2)), vec![armed(3, 5_000)], "retiring one flow must spare the rest");
}
#[test]
fn clearing_the_registry_drops_every_flow() {
let mut registry = TimerRegistry::default();
due_before(&mut registry, vec![armed(1, 5_000)], flow(1), 0);
due_before(&mut registry, vec![armed(2, 5_000)], flow(2), 0);
registry.clear();
assert_eq!(entries(®istry, flow(1)), Vec::new());
assert_eq!(entries(®istry, flow(2)), Vec::new());
}
}