pub mod guest_as_host;
pub mod rolling;
pub mod rolling_incremental;
pub mod rolling_top_k;
pub mod tumbling;
pub mod tumbling_carry;
use std::{collections::HashMap, hash::Hash};
use reifydb_codec::key::encoded::EncodedKey;
use reifydb_core::{
key::operator_state::GroupId,
state::store::{StateStore, TimerKind, TimerStore},
};
use reifydb_flow::{
operator::state::seal::{
coord::Coord,
ledger::{FiredAt, SealLedger},
policy::SEAL_GATE_STEP,
},
window::{engine::config::WindowEngineConfig, span::WindowSpan},
};
use reifydb_value::{
Result,
config::Config,
value::{datetime::DateTime, duration::Duration},
};
use crate::flow::operator::context::GuestContext;
pub(crate) fn seal_frontier(store: &mut (impl StateStore + TimerStore)) -> Result<DateTime> {
let ledger = SealLedger::read_order(store)?.unwrap_or(0);
let watermark = store.flow_watermark()?.map_or(0, |at| at.to_order());
Ok(<DateTime as Coord>::from_order(ledger.max(watermark)))
}
pub(crate) fn advance_seal_frontier(store: &mut (impl StateStore + TimerStore), fired: FiredAt) -> Result<DateTime> {
SealLedger::advance(store, fired)?;
seal_frontier(store)
}
pub(crate) fn bucket_of(coord: DateTime, size: Duration) -> DateTime {
WindowSpan::for_coord(coord, size).start
}
pub(crate) fn seal_horizon_of(frontier: DateTime, lateness: Duration) -> DateTime {
frontier.saturating_sub(lateness)
}
pub(crate) fn arm_seal_timer(store: &mut impl TimerStore, newest_window: DateTime, lateness: Duration) -> Result<()> {
let at = newest_window.saturating_add(lateness).saturating_add(SEAL_GATE_STEP);
store.arm_timer(at, TimerKind::Seal, &EncodedKey::new(Vec::new()))
}
pub(crate) type WindowGroups<G, C> = HashMap<(G, C), GroupId>;
pub(crate) fn intern_window_groups<G, C>(
ctx: &mut impl GuestContext,
windows: impl IntoIterator<Item = ((G, C), EncodedKey)>,
) -> Result<WindowGroups<G, C>>
where
G: Clone + Eq + Hash,
C: Copy + Eq + Hash,
{
let (windows, keys): (Vec<(G, C)>, Vec<EncodedKey>) = windows.into_iter().unzip();
if windows.is_empty() {
return Ok(WindowGroups::new());
}
Ok(windows.into_iter().zip(ctx.intern_groups(&keys)?).map(|(window, (group, _))| (window, group)).collect())
}
pub(crate) fn group_of<G, C>(groups: &WindowGroups<G, C>, group: &G, coord: C) -> GroupId
where
G: Clone + Eq + Hash,
C: Copy + Eq + Hash,
{
groups.get(&(group.clone(), coord)).copied().expect("every routed window is interned before the engine runs")
}
pub(crate) fn window_engine_config(_config: &Config) -> WindowEngineConfig {
WindowEngineConfig::builder().build()
}
#[cfg(test)]
mod tests {
use reifydb_core::key::operator_state::group_data_of_inner;
use crate::flow::operator::state::utils::empty_key;
#[test]
fn state_a_driver_addresses_without_a_group_can_never_be_reclaimed() {
let key = empty_key();
assert!(key.as_bytes().is_empty());
assert_eq!(
group_data_of_inner(key.as_bytes()),
None,
"state addressed without a group must not be attributable to any group"
);
}
}