1use 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}