Skip to main content

reifydb_sub_flow/operator/window/
rolling.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use std::collections::{BTreeMap, HashMap, HashSet};
5
6use reifydb_core::{
7	common::{TimeDomain, WindowKind},
8	interface::change::{Change, Diff},
9	value::column::columns::Columns,
10	window::{
11		accumulator::WindowAccumulator,
12		engine::{
13			AccumulatorEvent, EmitKind,
14			rolling::{
15				RollingBuckets, RollingBuffer, RollingEngine, RollingEviction, RollingExpiry,
16				RollingResult,
17			},
18		},
19		store::WindowStore,
20	},
21};
22use reifydb_engine::flow::aggregate::SlotKind;
23use reifydb_value::{
24	Result,
25	util::hash::Hash128,
26	value::{Value, duration::Duration, row_number::RowNumber},
27};
28use serde::{Deserialize, Serialize};
29
30use super::{
31	accumulator::{RowAccumulator, StampedAccumulator, WindowSlotKey},
32	operator::WindowOperator,
33	store::FlowWindowStore,
34	tumbling::slot_coord,
35};
36use crate::transaction::FlowTransaction;
37
38impl WindowOperator {
39	pub fn rolling_lag_ms(&self) -> u64 {
40		match &self.kind {
41			WindowKind::Rolling {
42				lag: Some(lag),
43				..
44			} => lag.milliseconds().unwrap_or(0) as u64,
45			_ => 0,
46		}
47	}
48}
49
50impl WindowOperator {
51	pub fn is_rolling_processing(&self) -> bool {
52		matches!(self.kind, WindowKind::Rolling { .. })
53			&& !self.is_count_based()
54			&& self.kind.time() == TimeDomain::Processing
55	}
56}
57
58#[derive(Default, Serialize, Deserialize)]
59struct RollingWindowMeta {
60	group_hash: u128,
61	row_number: u64,
62	group_values: Vec<Value>,
63	last_value: Vec<Value>,
64}
65
66type RollingEngineBuckets = RollingBuckets<Hash128, u64, (WindowSlotKey, Vec<Option<Value>>)>;
67
68fn combine_rolling(
69	buffer: &RollingBuffer<u64, RowAccumulator>,
70	kinds: &[SlotKind],
71	lag_ms: u64,
72	lateness: Option<Duration>,
73) -> Option<Vec<Value>> {
74	let (&newest, _) = buffer.iter().next_back()?;
75	let aggregate_cutoff = newest.saturating_sub(lag_ms);
76	let mut merged = RowAccumulator::new(kinds, lateness);
77	let mut any = false;
78	for (_coord, accumulator) in buffer.range(..=aggregate_cutoff) {
79		merged.merge(accumulator);
80		any = true;
81	}
82	if any {
83		merged.finalize()
84	} else {
85		None
86	}
87}
88
89#[allow(clippy::too_many_arguments)]
90fn route_rolling_columns(
91	operator: &WindowOperator,
92	columns: &Columns,
93	is_add: bool,
94	is_count: bool,
95	buckets: &mut RollingEngineBuckets,
96	group_values: &mut HashMap<Hash128, Vec<Value>>,
97	touched: &mut Vec<Hash128>,
98	touched_set: &mut HashSet<Hash128>,
99) -> Result<()> {
100	let row_count = columns.row_count();
101	if row_count == 0 {
102		return Ok(());
103	}
104	let groups = operator.core.compute_groups(columns)?;
105	let timestamps = if is_count {
106		Vec::new()
107	} else {
108		operator.resolve_event_timestamps(columns, row_count)?
109	};
110	let slot_cols = operator.core.evaluate_slot_inputs(columns)?;
111	for row_idx in 0..row_count {
112		let (hash, gvals) = &groups[row_idx];
113		let coord = if is_count {
114			columns.row_numbers[row_idx].0
115		} else {
116			timestamps[row_idx]
117		};
118		let slot_key = slot_coord(is_count, coord, columns.row_numbers[row_idx].0);
119		let contribution = (slot_key, operator.core.build_contribution(columns, &slot_cols, row_idx));
120		let event = if is_add {
121			AccumulatorEvent::Add(contribution)
122		} else {
123			AccumulatorEvent::Remove(contribution)
124		};
125		buckets.entry((*hash, coord)).or_default().push(event);
126		group_values.entry(*hash).or_insert_with(|| gvals.clone());
127		if touched_set.insert(*hash) {
128			touched.push(*hash);
129		}
130	}
131	Ok(())
132}
133
134pub fn apply_rolling_engine(operator: &WindowOperator, txn: &mut FlowTransaction, change: Change) -> Result<Change> {
135	let kinds = operator.core.slot_kinds.clone().expect("engine mode requires slot kinds");
136	let is_count = operator.is_count_based();
137	let lateness = operator.sealing_lateness();
138	let lag_ms = operator.rolling_lag_ms();
139	let is_event_time = !is_count && operator.kind.time() == TimeDomain::Event;
140	let size_ms = operator.size_duration().map(|d| d.milliseconds().unwrap_or(0) as u64).unwrap_or(0);
141
142	let mut buckets: RollingEngineBuckets = BTreeMap::new();
143	let mut group_values: HashMap<Hash128, Vec<Value>> = HashMap::new();
144	let mut touched: Vec<Hash128> = Vec::new();
145	let mut touched_set: HashSet<Hash128> = HashSet::new();
146	for diff in change.diffs.iter() {
147		match diff {
148			Diff::Insert {
149				post,
150				..
151			} => route_rolling_columns(
152				operator,
153				post,
154				true,
155				is_count,
156				&mut buckets,
157				&mut group_values,
158				&mut touched,
159				&mut touched_set,
160			)?,
161			Diff::Remove {
162				pre,
163				..
164			} => route_rolling_columns(
165				operator,
166				pre,
167				false,
168				is_count,
169				&mut buckets,
170				&mut group_values,
171				&mut touched,
172				&mut touched_set,
173			)?,
174			Diff::Update {
175				pre,
176				post,
177				..
178			} => {
179				route_rolling_columns(
180					operator,
181					pre,
182					false,
183					is_count,
184					&mut buckets,
185					&mut group_values,
186					&mut touched,
187					&mut touched_set,
188				)?;
189				route_rolling_columns(
190					operator,
191					post,
192					true,
193					is_count,
194					&mut buckets,
195					&mut group_values,
196					&mut touched,
197					&mut touched_set,
198				)?;
199			}
200		}
201	}
202
203	if buckets.is_empty() {
204		return Ok(Change::from_flow(operator.core.node, change.version, Vec::new(), change.changed_at));
205	}
206
207	let eviction = if is_count {
208		RollingEviction::Capacity(operator.size_count().unwrap_or(0) as usize)
209	} else if is_event_time {
210		let batch_max = buckets.keys().map(|&(_, coord)| coord).max().unwrap_or(0);
211		operator.advance_event_watermark(txn, batch_max)?;
212		RollingEviction::Before(operator.event_time_cutoff(txn, size_ms + lag_ms)?)
213	} else {
214		RollingEviction::Before(operator.core.current_timestamp().saturating_sub(size_ms + lag_ms))
215	};
216
217	let results = {
218		let mut store = FlowWindowStore::new(txn, operator.core.node);
219		for hash in &touched {
220			let key = operator.core.create_window_key(*hash, 0);
221			store.get_or_create_row_number(&key)?;
222		}
223		let mut engine = RollingEngine::<Hash128, u64, RowAccumulator>::new(operator.engine_config());
224		let res = engine.apply_evicting(
225			&mut store,
226			buckets,
227			eviction,
228			|hash| operator.core.create_window_key(*hash, 0),
229			|| RowAccumulator::new(&kinds, lateness),
230			|_g, buffer| combine_rolling(buffer, &kinds, lag_ms, lateness),
231		)?;
232		engine.flush(&mut store)?;
233		res
234	};
235
236	let diffs = finish_rolling_results(operator, txn, &change, &results, &group_values, &touched)?;
237	Ok(Change::from_flow(operator.core.node, change.version, diffs, change.changed_at))
238}
239
240fn finish_rolling_results(
241	operator: &WindowOperator,
242	txn: &mut FlowTransaction,
243	change: &Change,
244	results: &[RollingResult<Hash128, Vec<Value>>],
245	group_values: &HashMap<Hash128, Vec<Value>>,
246	touched: &[Hash128],
247) -> Result<Vec<Diff>> {
248	let ts_nanos = change.changed_at.to_nanos();
249	let mut diffs = Vec::new();
250	let mut emitted: HashSet<Hash128> = HashSet::new();
251	let mut store = FlowWindowStore::new(txn, operator.core.node);
252	for r in results {
253		emitted.insert(r.group);
254		let meta_key = operator.create_rolling_meta_key(r.group);
255		let prior = store.state_get::<RollingWindowMeta>(&meta_key)?;
256		if matches!(r.kind, EmitKind::Remove) {
257			if let Some(m) = prior {
258				let pre = operator.core.build_engine_row(
259					&m.group_values,
260					&m.last_value,
261					RowNumber(m.row_number),
262					ts_nanos,
263				)?;
264				diffs.push(Diff::remove(Columns::from_row(&pre)));
265				store.state_drop(&meta_key)?;
266			}
267			continue;
268		}
269		let gvals = group_values.get(&r.group).cloned().unwrap_or_default();
270		let post = operator.core.build_engine_row(&gvals, &r.value, r.row_number, ts_nanos)?;
271		match prior {
272			Some(m) => {
273				let pre = operator.core.build_engine_row(
274					&gvals,
275					&m.last_value,
276					r.row_number,
277					ts_nanos,
278				)?;
279				diffs.push(Diff::update(Columns::from_row(&pre), Columns::from_row(&post)));
280			}
281			None => diffs.push(Diff::insert(Columns::from_row(&post))),
282		}
283		store.state_set(
284			&meta_key,
285			&RollingWindowMeta {
286				group_hash: r.group.0,
287				row_number: r.row_number.0,
288				group_values: gvals,
289				last_value: r.value.clone(),
290			},
291		)?;
292	}
293	for hash in touched {
294		if emitted.contains(hash) {
295			continue;
296		}
297		let meta_key = operator.create_rolling_meta_key(*hash);
298		if let Some(m) = store.state_get::<RollingWindowMeta>(&meta_key)? {
299			let pre = operator.core.build_engine_row(
300				&m.group_values,
301				&m.last_value,
302				RowNumber(m.row_number),
303				ts_nanos,
304			)?;
305			diffs.push(Diff::remove(Columns::from_row(&pre)));
306			store.state_drop(&meta_key)?;
307		}
308	}
309	Ok(diffs)
310}
311
312pub fn tick_expire_rolling_engine(
313	operator: &WindowOperator,
314	txn: &mut FlowTransaction,
315	current_timestamp: u64,
316) -> Result<Vec<Diff>> {
317	let size_ms = match operator.size_duration() {
318		Some(d) => d.milliseconds().unwrap_or(0) as u64,
319		None => return Ok(Vec::new()),
320	};
321	if size_ms == 0 {
322		return Ok(Vec::new());
323	}
324	let lag_ms = operator.rolling_lag_ms();
325	let lateness = operator.sealing_lateness();
326	let kinds = operator.core.slot_kinds.clone().expect("engine mode requires slot kinds");
327	let cutoff = operator.event_time_cutoff(txn, size_ms + lag_ms)?;
328	let ts_nanos = current_timestamp.saturating_mul(1_000_000);
329
330	let expiries = {
331		let mut store = FlowWindowStore::new(txn, operator.core.node);
332		let mut engine = RollingEngine::<Hash128, u64, RowAccumulator>::new(operator.engine_config());
333		let res = engine.expire_before(&mut store, cutoff, |_g, buffer| {
334			combine_rolling(buffer, &kinds, lag_ms, lateness)
335		})?;
336		engine.flush(&mut store)?;
337		res
338	};
339
340	let mut diffs = Vec::new();
341	let mut store = FlowWindowStore::new(txn, operator.core.node);
342	for expiry in expiries {
343		match expiry {
344			RollingExpiry::Update {
345				row_number,
346				group,
347				value,
348			} => {
349				let meta_key = operator.create_rolling_meta_key(group);
350				let Some(meta) = store.state_get::<RollingWindowMeta>(&meta_key)? else {
351					continue;
352				};
353				let pre = operator.core.build_engine_row(
354					&meta.group_values,
355					&meta.last_value,
356					row_number,
357					ts_nanos,
358				)?;
359				let post = operator.core.build_engine_row(
360					&meta.group_values,
361					&value,
362					row_number,
363					ts_nanos,
364				)?;
365				diffs.push(Diff::update(Columns::from_row(&pre), Columns::from_row(&post)));
366				store.state_set(
367					&meta_key,
368					&RollingWindowMeta {
369						group_hash: meta.group_hash,
370						row_number: meta.row_number,
371						group_values: meta.group_values,
372						last_value: value,
373					},
374				)?;
375			}
376			RollingExpiry::Remove {
377				row_number,
378				group,
379			} => {
380				let meta_key = operator.create_rolling_meta_key(group);
381				let Some(meta) = store.state_get::<RollingWindowMeta>(&meta_key)? else {
382					continue;
383				};
384				let pre = operator.core.build_engine_row(
385					&meta.group_values,
386					&meta.last_value,
387					row_number,
388					ts_nanos,
389				)?;
390				diffs.push(Diff::remove(Columns::from_row(&pre)));
391				store.state_drop(&meta_key)?;
392			}
393		}
394	}
395	Ok(diffs)
396}
397
398type StampedBuckets = RollingBuckets<Hash128, u64, ((WindowSlotKey, Vec<Option<Value>>), u64)>;
399
400fn combine_stamped(buffer: &RollingBuffer<u64, StampedAccumulator>, kinds: &[SlotKind]) -> Option<Vec<Value>> {
401	let mut merged = RowAccumulator::new(kinds, None);
402	let mut any = false;
403	for (_coord, accumulator) in buffer.iter() {
404		merged.merge(accumulator.inner());
405		any = true;
406	}
407	if any {
408		merged.finalize()
409	} else {
410		None
411	}
412}
413
414#[allow(clippy::too_many_arguments)]
415fn route_rolling_processing(
416	operator: &WindowOperator,
417	columns: &Columns,
418	is_add: bool,
419	now: u64,
420	buckets: &mut StampedBuckets,
421	group_values: &mut HashMap<Hash128, Vec<Value>>,
422	touched: &mut Vec<Hash128>,
423	touched_set: &mut HashSet<Hash128>,
424) -> Result<()> {
425	let row_count = columns.row_count();
426	if row_count == 0 {
427		return Ok(());
428	}
429	let groups = operator.core.compute_groups(columns)?;
430	let slot_cols = operator.core.evaluate_slot_inputs(columns)?;
431	for (row_idx, (hash, gvals)) in groups.iter().enumerate() {
432		let coord = columns.row_numbers[row_idx].0;
433		let slot_key = slot_coord(true, 0, coord);
434		let value_contrib = (slot_key, operator.core.build_contribution(columns, &slot_cols, row_idx));
435		let event = if is_add {
436			AccumulatorEvent::Add((value_contrib, now))
437		} else {
438			AccumulatorEvent::Remove((value_contrib, 0))
439		};
440		buckets.entry((*hash, coord)).or_default().push(event);
441		group_values.entry(*hash).or_insert_with(|| gvals.clone());
442		if touched_set.insert(*hash) {
443			touched.push(*hash);
444		}
445	}
446	Ok(())
447}
448
449pub fn apply_rolling_processing_engine(
450	operator: &WindowOperator,
451	txn: &mut FlowTransaction,
452	change: Change,
453) -> Result<Change> {
454	let kinds = operator.core.slot_kinds.clone().expect("engine mode requires slot kinds");
455	let size_ms = operator.size_duration().map(|d| d.milliseconds().unwrap_or(0) as u64).unwrap_or(0);
456	let lag_ms = operator.rolling_lag_ms();
457	let now = operator.core.current_timestamp();
458
459	let mut buckets: StampedBuckets = BTreeMap::new();
460	let mut group_values: HashMap<Hash128, Vec<Value>> = HashMap::new();
461	let mut touched: Vec<Hash128> = Vec::new();
462	let mut touched_set: HashSet<Hash128> = HashSet::new();
463	for diff in change.diffs.iter() {
464		match diff {
465			Diff::Insert {
466				post,
467				..
468			} => route_rolling_processing(
469				operator,
470				post,
471				true,
472				now,
473				&mut buckets,
474				&mut group_values,
475				&mut touched,
476				&mut touched_set,
477			)?,
478			Diff::Remove {
479				pre,
480				..
481			} => route_rolling_processing(
482				operator,
483				pre,
484				false,
485				now,
486				&mut buckets,
487				&mut group_values,
488				&mut touched,
489				&mut touched_set,
490			)?,
491			Diff::Update {
492				pre,
493				post,
494				..
495			} => {
496				route_rolling_processing(
497					operator,
498					pre,
499					false,
500					now,
501					&mut buckets,
502					&mut group_values,
503					&mut touched,
504					&mut touched_set,
505				)?;
506				route_rolling_processing(
507					operator,
508					post,
509					true,
510					now,
511					&mut buckets,
512					&mut group_values,
513					&mut touched,
514					&mut touched_set,
515				)?;
516			}
517		}
518	}
519
520	if buckets.is_empty() {
521		return Ok(Change::from_flow(operator.core.node, change.version, Vec::new(), change.changed_at));
522	}
523
524	let cutoff = now.saturating_sub(size_ms + lag_ms);
525	let results = {
526		let mut store = FlowWindowStore::new(txn, operator.core.node);
527		for hash in &touched {
528			let key = operator.core.create_window_key(*hash, 0);
529			store.get_or_create_row_number(&key)?;
530		}
531		let mut engine = RollingEngine::<Hash128, u64, StampedAccumulator>::new(operator.engine_config());
532		let res = engine.apply_evicting(
533			&mut store,
534			buckets,
535			RollingEviction::BeforeStamp(cutoff),
536			|hash| operator.core.create_window_key(*hash, 0),
537			|| StampedAccumulator::new(&kinds, None),
538			|_g, buffer| combine_stamped(buffer, &kinds),
539		)?;
540		engine.flush(&mut store)?;
541		res
542	};
543
544	let diffs = finish_rolling_results(operator, txn, &change, &results, &group_values, &touched)?;
545	Ok(Change::from_flow(operator.core.node, change.version, diffs, change.changed_at))
546}
547
548pub fn tick_expire_rolling_processing_engine(
549	operator: &WindowOperator,
550	txn: &mut FlowTransaction,
551	current_timestamp: u64,
552) -> Result<Vec<Diff>> {
553	let size_ms = match operator.size_duration() {
554		Some(d) => d.milliseconds().unwrap_or(0) as u64,
555		None => return Ok(Vec::new()),
556	};
557	if size_ms == 0 {
558		return Ok(Vec::new());
559	}
560	let lag_ms = operator.rolling_lag_ms();
561	let kinds = operator.core.slot_kinds.clone().expect("engine mode requires slot kinds");
562	let cutoff = current_timestamp.saturating_sub(size_ms + lag_ms);
563	let ts_nanos = current_timestamp.saturating_mul(1_000_000);
564
565	let expiries = {
566		let mut store = FlowWindowStore::new(txn, operator.core.node);
567		let mut engine = RollingEngine::<Hash128, u64, StampedAccumulator>::new(operator.engine_config());
568		let res =
569			engine.expire_before_stamp(&mut store, cutoff, |_g, buffer| combine_stamped(buffer, &kinds))?;
570		engine.flush(&mut store)?;
571		res
572	};
573
574	let mut diffs = Vec::new();
575	let mut store = FlowWindowStore::new(txn, operator.core.node);
576	for expiry in expiries {
577		match expiry {
578			RollingExpiry::Update {
579				row_number,
580				group,
581				value,
582			} => {
583				let meta_key = operator.create_rolling_meta_key(group);
584				let Some(meta) = store.state_get::<RollingWindowMeta>(&meta_key)? else {
585					continue;
586				};
587				let pre = operator.core.build_engine_row(
588					&meta.group_values,
589					&meta.last_value,
590					row_number,
591					ts_nanos,
592				)?;
593				let post = operator.core.build_engine_row(
594					&meta.group_values,
595					&value,
596					row_number,
597					ts_nanos,
598				)?;
599				diffs.push(Diff::update(Columns::from_row(&pre), Columns::from_row(&post)));
600				store.state_set(
601					&meta_key,
602					&RollingWindowMeta {
603						group_hash: meta.group_hash,
604						row_number: meta.row_number,
605						group_values: meta.group_values,
606						last_value: value,
607					},
608				)?;
609			}
610			RollingExpiry::Remove {
611				row_number,
612				group,
613			} => {
614				let meta_key = operator.create_rolling_meta_key(group);
615				let Some(meta) = store.state_get::<RollingWindowMeta>(&meta_key)? else {
616					continue;
617				};
618				let pre = operator.core.build_engine_row(
619					&meta.group_values,
620					&meta.last_value,
621					row_number,
622					ts_nanos,
623				)?;
624				diffs.push(Diff::remove(Columns::from_row(&pre)));
625				store.state_drop(&meta_key)?;
626			}
627		}
628	}
629	Ok(diffs)
630}