use reifydb_core::interface::change::Change;
#[cfg(feature = "runtime")]
use reifydb_core::{
interface::{catalog::flow::OperatorId, flow::OperatorCapability},
metrics::heap::OperatorSample,
value::column::columns::Columns,
};
#[cfg(feature = "runtime")]
use reifydb_value::Result;
use reifydb_value::value::datetime::DateTime;
#[cfg(any(feature = "runtime", all(reifydb_target = "host", not(reifydb_dst))))]
use reifydb_value::value::duration::Duration;
#[cfg(feature = "runtime")]
use crate::{operator::host::HostContext, timer::Timer};
#[cfg(all(reifydb_target = "host", not(reifydb_dst)))]
pub fn scale_from_millis(span: Option<u64>) -> Option<Duration> {
span.filter(|millis| *millis > 0)
.and_then(|millis| i64::try_from(millis).ok())
.and_then(|millis| Duration::from_milliseconds(millis).ok())
}
#[cfg(feature = "runtime")]
pub mod aggregation;
#[cfg(feature = "runtime")]
pub mod append;
#[cfg(feature = "runtime")]
pub mod apply;
#[cfg(feature = "runtime")]
pub mod distinct;
#[cfg(feature = "runtime")]
pub mod drops;
#[cfg(feature = "runtime")]
pub mod extend;
#[cfg(feature = "runtime")]
pub mod filter;
#[cfg(feature = "runtime")]
pub mod gate;
#[cfg(feature = "runtime")]
pub mod guard;
#[cfg(feature = "runtime")]
pub mod host;
#[cfg(feature = "runtime")]
pub mod join;
#[cfg(feature = "runtime")]
pub mod map;
#[cfg(feature = "runtime")]
pub mod metrics;
#[cfg(feature = "runtime")]
pub mod provider;
#[cfg(feature = "runtime")]
pub mod scan;
#[cfg(feature = "runtime")]
pub mod sink;
#[cfg(feature = "runtime")]
pub mod sort;
pub mod state;
pub mod state_access;
#[cfg(feature = "runtime")]
pub mod take;
#[cfg(feature = "runtime")]
pub mod window;
#[cfg(feature = "runtime")]
pub trait HostOperator: Send {
fn id(&self) -> OperatorId;
fn capabilities(&self) -> &[OperatorCapability];
fn apply(&mut self, host: &mut dyn HostContext, change: Change) -> Result<Change>;
fn on_timer(&mut self, _host: &mut dyn HostContext, _timer: Timer) -> Result<Option<Change>> {
Ok(None)
}
fn seal_span(&self) -> Option<Duration> {
None
}
fn sample(&self) -> Option<OperatorSample> {
None
}
fn output_schema(&self) -> Option<Columns> {
None
}
}
#[cfg(feature = "runtime")]
pub type BoxedHostOperator = Box<dyn HostOperator>;
pub fn max_input_time(change: &Change) -> Option<DateTime> {
change.diffs
.iter()
.filter_map(|diff| diff.post().or_else(|| diff.pre()))
.flat_map(|columns| columns.time().iter().copied())
.max()
}
#[cfg_attr(not(feature = "runtime"), allow(dead_code))]
pub(crate) fn stamp_output_time(change: &mut Change, inherited: Option<DateTime>) {
let Some(inherited) = inherited else {
return;
};
for diff in change.diffs.iter_mut() {
for columns in diff.columns_mut() {
let stamped: Vec<DateTime> = columns.time().iter().map(|own| (*own).min(inherited)).collect();
columns.system.set_time(stamped);
}
}
}
#[cfg(test)]
mod substrate_stamping_tests {
use reifydb_core::{
common::CommitVersion,
interface::{
catalog::flow::OperatorId,
change::{Diff, Diffs},
},
value::column::{ColumnWithName, buffer::ColumnBuffer, columns::Columns},
};
use reifydb_value::{
factory::time::at_millis,
fragment::Fragment,
value::{row_number::RowNumber, system_columns::SystemColumns},
};
use super::*;
fn columns(times: &[DateTime]) -> Columns {
let n = times.len();
Columns::with_system(
vec![ColumnWithName::new(
Fragment::internal("v"),
ColumnBuffer::int4((0..n as i32).collect::<Vec<_>>()),
)],
SystemColumns::new(
(1..=n as u64).map(RowNumber).collect(),
Vec::new(),
vec![at_millis(0); n],
vec![at_millis(0); n],
times.to_vec(),
),
)
}
fn untimed_columns(n: usize) -> Columns {
Columns::with_system(
vec![ColumnWithName::new(
Fragment::internal("v"),
ColumnBuffer::int4((0..n as i32).collect::<Vec<_>>()),
)],
SystemColumns::new(
(1..=n as u64).map(RowNumber).collect(),
Vec::new(),
vec![at_millis(0); n],
vec![at_millis(0); n],
Vec::new(),
),
)
}
fn change(diffs: Diffs) -> Change {
Change::from_flow(OperatorId(1), CommitVersion(1), diffs, at_millis(0))
}
#[test]
fn the_substrate_stamps_output_with_the_max_input_time() {
let mut diffs = Diffs::new();
diffs.push(Diff::insert(columns(&[at_millis(1_000), at_millis(9_000), at_millis(5_000)])));
assert_eq!(max_input_time(&change(diffs)), Some(at_millis(9_000)));
}
#[test]
fn an_operator_cannot_influence_its_own_output_time() {
let mut produced = Diffs::new();
produced.push(Diff::insert(columns(&[at_millis(999_999), at_millis(1_000)])));
let mut out = change(produced);
stamp_output_time(&mut out, Some(at_millis(4_000)));
assert_eq!(
out.diffs[0].post().unwrap().time().to_vec(),
vec![at_millis(4_000), at_millis(1_000)],
"a row above the inherited instant is pulled down; one below keeps its own"
);
}
#[test]
fn both_sides_of_an_update_are_stamped() {
let mut produced = Diffs::new();
produced.push(Diff::update(columns(&[at_millis(9_000)]), columns(&[at_millis(10_000)])));
let mut out = change(produced);
stamp_output_time(&mut out, Some(at_millis(7_000)));
assert_eq!(out.diffs[0].pre().unwrap().time().to_vec(), vec![at_millis(7_000)]);
assert_eq!(out.diffs[0].post().unwrap().time().to_vec(), vec![at_millis(7_000)]);
}
#[test]
fn a_fan_out_operator_has_every_emitted_row_stamped() {
let mut produced = Diffs::new();
produced.push(Diff::insert(columns(&[
at_millis(9_000),
at_millis(10_000),
at_millis(11_000),
at_millis(12_000),
at_millis(13_000),
])));
let mut out = change(produced);
stamp_output_time(&mut out, Some(at_millis(8_000)));
assert_eq!(out.diffs[0].post().unwrap().time().to_vec(), vec![at_millis(8_000); 5]);
}
#[test]
fn an_empty_input_leaves_the_output_untouched() {
let empty = change(Diffs::new());
assert_eq!(max_input_time(&empty), None);
let mut produced = Diffs::new();
produced.push(Diff::insert(columns(&[at_millis(3_000)])));
let mut out = change(produced);
stamp_output_time(&mut out, None);
assert_eq!(out.diffs[0].post().unwrap().time().to_vec(), vec![at_millis(3_000)]);
}
#[test]
fn an_operator_stamping_above_its_inputs_is_still_overwritten() {
let inherited = at_millis(5_000);
let one_nano_above = DateTime::from_nanos(inherited.to_nanos() + 1);
let mut produced = Diffs::new();
produced.push(Diff::insert(columns(&[one_nano_above])));
let mut out = change(produced);
stamp_output_time(&mut out, Some(inherited));
assert_eq!(out.diffs[0].post().unwrap().time().to_vec(), vec![inherited]);
}
#[test]
fn an_operator_stamping_at_or_below_its_inputs_keeps_its_stamp() {
let inherited = at_millis(5_000);
let mut produced = Diffs::new();
produced.push(Diff::insert(columns(&[at_millis(1_000), inherited])));
let mut out = change(produced);
stamp_output_time(&mut out, Some(inherited));
assert_eq!(
out.diffs[0].post().unwrap().time().to_vec(),
vec![at_millis(1_000), inherited],
"below survives, and equal counts as below"
);
}
#[test]
fn a_row_stamped_at_the_epoch_keeps_its_own_instant() {
let mut produced = Diffs::new();
produced.push(Diff::insert(columns(&[DateTime::default()])));
let mut out = change(produced);
stamp_output_time(&mut out, Some(at_millis(6_000)));
assert_eq!(out.diffs[0].post().unwrap().time().to_vec(), vec![DateTime::default()]);
}
#[test]
fn a_time_less_batch_stays_time_less_through_stamping() {
let mut produced = Diffs::new();
produced.push(Diff::insert(untimed_columns(3)));
let mut out = change(produced);
stamp_output_time(&mut out, Some(at_millis(6_000)));
assert!(out.diffs[0].post().unwrap().time().is_empty(), "#time must stay absent");
}
#[test]
fn a_window_row_carries_its_window_start_through_the_apply_wrapper() {
let window_start = at_millis(60_000);
let newest_event_in_bucket = at_millis(119_000);
let seal_fires_at = at_millis(120_001);
let mut on_apply = change({
let mut d = Diffs::new();
d.push(Diff::insert(columns(&[window_start])));
d
});
stamp_output_time(&mut on_apply, Some(newest_event_in_bucket));
assert_eq!(on_apply.diffs[0].post().unwrap().time().to_vec(), vec![window_start]);
let mut on_timer = change({
let mut d = Diffs::new();
d.push(Diff::insert(columns(&[window_start])));
d
});
stamp_output_time(&mut on_timer, Some(seal_fires_at));
assert_eq!(on_timer.diffs[0].post().unwrap().time().to_vec(), vec![window_start]);
}
#[test]
fn the_max_spans_every_diff_in_the_batch() {
let mut diffs = Diffs::new();
diffs.push(Diff::insert(columns(&[at_millis(1_000)])));
diffs.push(Diff::insert(columns(&[at_millis(12_000)])));
diffs.push(Diff::insert(columns(&[at_millis(3_000)])));
assert_eq!(max_input_time(&change(diffs)), Some(at_millis(12_000)));
}
}