1use std::collections::{BTreeMap, BTreeSet, HashMap};
19use std::time::Duration;
20
21use common_telemetry::debug;
22use common_time::Timestamp;
23use datatypes::value::Value;
24use session::context::QueryContextRef;
25use snafu::{OptionExt, ResultExt, ensure};
26use tokio::sync::oneshot;
27use tokio::time::Instant;
28
29use crate::batching_mode::task::BatchingTask;
30use crate::batching_mode::time_window::TimeWindowExpr;
31use crate::error::{DatatypesSnafu, InternalSnafu, TimeSnafu, UnexpectedSnafu};
32use crate::metrics::{
33 METRIC_FLOW_BATCHING_ENGINE_QUERY_WINDOW_CNT, METRIC_FLOW_BATCHING_ENGINE_QUERY_WINDOW_SIZE,
34 METRIC_FLOW_BATCHING_ENGINE_STALLED_WINDOW_SIZE,
35};
36use crate::{Error, FlowId};
37
38#[derive(Debug)]
40pub struct TaskState {
41 pub(crate) query_ctx: QueryContextRef,
43 last_update_time: Instant,
45 last_query_duration: Duration,
47 last_exec_time_millis: Option<i64>,
49 start_time_millis: Option<i64>,
51 pub(crate) dirty_time_windows: DirtyTimeWindows,
54 checkpoint_mode: CheckpointMode,
55 pending_fenced_repair: Option<FencedRepair>,
56 checkpoints: BTreeMap<u64, u64>,
59 incremental_disabled: bool,
63 exec_state: ExecState,
64 pub(crate) shutdown_rx: oneshot::Receiver<()>,
66 pub(crate) task_handle: Option<tokio::task::JoinHandle<()>>,
68}
69impl TaskState {
70 pub fn new(query_ctx: QueryContextRef, shutdown_rx: oneshot::Receiver<()>) -> Self {
71 Self::with_dirty_time_windows(query_ctx, shutdown_rx, DirtyTimeWindows::default())
72 }
73
74 pub fn with_dirty_time_windows(
75 query_ctx: QueryContextRef,
76 shutdown_rx: oneshot::Receiver<()>,
77 dirty_time_windows: DirtyTimeWindows,
78 ) -> Self {
79 Self {
80 query_ctx,
81 last_update_time: Instant::now(),
82 last_query_duration: Duration::from_secs(0),
83 last_exec_time_millis: None,
84 start_time_millis: None,
85 dirty_time_windows,
86 checkpoint_mode: CheckpointMode::FullSnapshot,
87 pending_fenced_repair: None,
88 checkpoints: Default::default(),
89 incremental_disabled: false,
90 exec_state: ExecState::Idle,
91 shutdown_rx,
92 task_handle: None,
93 }
94 }
95
96 pub fn record_start_time_if_first(&mut self) {
99 if self.start_time_millis.is_none() {
100 self.start_time_millis = Some(common_time::util::current_time_millis());
105 }
106 }
107
108 pub fn after_query_exec(&mut self, elapsed: Duration, is_succ: bool) {
109 self.exec_state = ExecState::Idle;
110 self.last_query_duration = elapsed;
111 self.last_update_time = Instant::now();
112 if is_succ {
113 self.last_exec_time_millis = Some(common_time::util::current_time_millis());
114 }
115 }
116
117 pub fn last_execution_time_millis(&self) -> Option<i64> {
118 self.last_exec_time_millis
119 }
120
121 pub fn start_time_millis(&self) -> Option<i64> {
123 self.start_time_millis
124 }
125
126 pub fn checkpoint_mode(&self) -> CheckpointMode {
127 self.checkpoint_mode
128 }
129
130 pub fn checkpoints(&self) -> &BTreeMap<u64, u64> {
131 &self.checkpoints
132 }
133
134 pub fn dirty_time_windows(&self) -> &DirtyTimeWindows {
136 &self.dirty_time_windows
137 }
138
139 pub fn pending_fenced_repair(&self) -> Option<&FencedRepair> {
142 self.pending_fenced_repair.as_ref()
143 }
144
145 pub fn is_incremental_disabled(&self) -> bool {
146 self.incremental_disabled
147 }
148
149 pub fn disable_incremental(&mut self) {
152 self.incremental_disabled = true;
153 self.mark_full_snapshot();
154 }
155
156 pub fn mark_full_snapshot(&mut self) {
160 self.abandon_fenced_repair();
161 }
162
163 pub fn advance_checkpoints(&mut self, watermark_map: HashMap<u64, u64>) {
166 self.checkpoints = watermark_map.into_iter().collect();
167 self.pending_fenced_repair = None;
168 if !self.incremental_disabled {
169 self.checkpoint_mode = CheckpointMode::Incremental;
170 }
171 }
172
173 pub fn advance_incremental_checkpoints_with_participation(
176 &mut self,
177 participating_regions: &BTreeSet<u64>,
178 watermark_map: HashMap<u64, u64>,
179 ) {
180 for region_id in participating_regions {
181 if let Some(seq) = watermark_map.get(region_id) {
182 self.checkpoints.insert(*region_id, *seq);
183 }
184 }
185 if !self.incremental_disabled {
186 self.checkpoint_mode = CheckpointMode::Incremental;
187 }
188 self.pending_fenced_repair = None;
189 }
190
191 pub fn start_fenced_repair(&mut self, high: BTreeMap<u64, u64>) -> Option<&FencedRepair> {
196 if self.dirty_time_windows.is_empty() {
197 self.pending_fenced_repair = None;
198 return None;
199 }
200
201 let pending_windows = self.dirty_time_windows.clone();
202 self.dirty_time_windows.clean();
203 self.pending_fenced_repair = Some(FencedRepair {
204 high,
205 pending_windows,
206 });
207 self.checkpoint_mode = CheckpointMode::FullSnapshot;
208 self.pending_fenced_repair.as_ref()
209 }
210
211 pub fn start_fenced_repair_windows(
216 &mut self,
217 high: BTreeMap<u64, u64>,
218 windows: Vec<(Timestamp, Timestamp)>,
219 ) {
220 let mut pending_windows = self.dirty_time_windows.clone();
221 pending_windows.clean();
222 pending_windows.add_windows(windows);
223 self.pending_fenced_repair = Some(FencedRepair {
224 high,
225 pending_windows,
226 });
227 self.checkpoint_mode = CheckpointMode::FullSnapshot;
228 }
229
230 pub fn finish_fenced_repair(&mut self) -> Option<BTreeMap<u64, u64>> {
233 let repair = self.pending_fenced_repair.take()?;
234 self.checkpoints = repair.high;
235 if !self.incremental_disabled {
236 self.checkpoint_mode = CheckpointMode::Incremental;
237 }
238 Some(self.checkpoints.clone())
239 }
240
241 pub fn abandon_fenced_repair(&mut self) -> bool {
244 self.checkpoint_mode = CheckpointMode::FullSnapshot;
245 let Some(repair) = self.pending_fenced_repair.take() else {
246 return false;
247 };
248
249 self.dirty_time_windows
250 .add_dirty_windows(&repair.pending_windows);
251 true
252 }
253
254 pub fn restore_scoped_windows(&mut self, filter: &FilterExprInfo) {
258 if let Some(repair) = self.pending_fenced_repair.as_mut() {
259 repair
260 .pending_windows
261 .add_windows(filter.time_ranges.clone());
262 return;
263 }
264
265 self.dirty_time_windows
266 .add_windows(filter.time_ranges.clone());
267 }
268
269 pub fn gen_scoped_filter_exprs(
272 &mut self,
273 col_name: &str,
274 expire_lower_bound: Option<Timestamp>,
275 window_size: chrono::Duration,
276 window_cnt: usize,
277 flow_id: FlowId,
278 task_ctx: Option<&BatchingTask>,
279 ) -> Result<Option<FilterExprInfo>, Error> {
280 if let Some(repair) = self.pending_fenced_repair.as_mut() {
281 return repair.pending_windows.gen_filter_exprs(
286 col_name,
287 None,
288 window_size,
289 window_cnt,
290 flow_id,
291 task_ctx,
292 );
293 }
294
295 self.dirty_time_windows.gen_filter_exprs(
296 col_name,
297 expire_lower_bound,
298 window_size,
299 window_cnt,
300 flow_id,
301 task_ctx,
302 )
303 }
304
305 pub fn fenced_repair_watermarks_match_high(
308 &self,
309 participating_regions: &BTreeSet<u64>,
310 watermark_map: &HashMap<u64, u64>,
311 ) -> bool {
312 let Some(repair) = self.pending_fenced_repair.as_ref() else {
313 return false;
314 };
315
316 !participating_regions.is_empty()
317 && participating_regions.len() == repair.high.len()
318 && watermark_map.len() == repair.high.len()
319 && participating_regions.iter().all(|region_id| {
320 repair
321 .high
322 .get(region_id)
323 .zip(watermark_map.get(region_id))
324 .is_some_and(|(high, watermark)| high == watermark)
325 })
326 }
327
328 pub fn fenced_repair_pending_is_empty(&self) -> bool {
330 self.pending_fenced_repair
331 .as_ref()
332 .is_some_and(|repair| repair.pending_windows.is_empty())
333 }
334
335 pub fn can_advance_full_snapshot_checkpoints(
338 &self,
339 participating_regions: &BTreeSet<u64>,
340 watermark_map: &HashMap<u64, u64>,
341 ) -> bool {
342 !participating_regions.is_empty()
343 && participating_regions.len() == watermark_map.len()
344 && participating_regions
345 .iter()
346 .all(|region_id| watermark_map.contains_key(region_id))
347 }
348
349 pub fn can_advance_incremental_checkpoints_with_participation(
352 &self,
353 participating_regions: &BTreeSet<u64>,
354 watermark_map: &HashMap<u64, u64>,
355 ) -> bool {
356 !self.incremental_disabled
357 && !self.checkpoints.is_empty()
358 && !participating_regions.is_empty()
359 && participating_regions.len() == watermark_map.len()
360 && participating_regions
361 .iter()
362 .all(|region_id| self.checkpoints.contains_key(region_id))
363 && participating_regions.iter().all(|region_id| {
364 let checkpoint = self.checkpoints.get(region_id);
365 watermark_map
366 .get(region_id)
367 .zip(checkpoint)
368 .is_some_and(|(seq, checkpoint)| seq >= checkpoint)
369 })
370 }
371
372 pub fn get_next_start_query_time(
390 &self,
391 flow_id: FlowId,
392 time_window_size: &Option<Duration>,
393 min_refresh_duration: Duration,
394 max_timeout: Option<Duration>,
395 max_filter_num_per_query: usize,
396 prefer_short_incremental_cadence: bool,
397 ) -> Instant {
398 let lower = time_window_size.unwrap_or(min_refresh_duration);
400 let next_duration = self.last_query_duration.max(lower);
401 let next_duration = if let Some(max_timeout) = max_timeout {
402 next_duration.min(max_timeout)
403 } else {
404 next_duration
405 };
406
407 if self
408 .pending_fenced_repair
409 .as_ref()
410 .is_some_and(|repair| !repair.pending_windows().is_empty())
411 {
412 debug!(
413 "Flow id = {}, active fenced repair still has pending windows, execute immediately",
414 flow_id,
415 );
416 return Instant::now();
417 }
418
419 let cur_dirty_window_size = self.dirty_time_windows.window_size();
420 let max_query_update_range = (*time_window_size)
422 .unwrap_or_default()
423 .mul_f64(max_filter_num_per_query as f64);
424 if cur_dirty_window_size < max_query_update_range {
427 if prefer_short_incremental_cadence {
428 let next_duration = self.last_query_duration.max(min_refresh_duration);
432 let next_duration = if let Some(max_timeout) = max_timeout {
433 next_duration.min(max_timeout)
434 } else {
435 next_duration
436 };
437 self.last_update_time + next_duration
438 } else {
439 self.last_update_time + next_duration
440 }
441 } else {
442 debug!(
445 "Flow id = {}, still have too many {} dirty time window({:?}), execute immediately",
446 flow_id,
447 self.dirty_time_windows.windows.len(),
448 self.dirty_time_windows.windows
449 );
450 Instant::now()
451 }
452 }
453}
454
455#[derive(Debug, Clone)]
458pub struct DirtyTimeWindows {
459 windows: BTreeMap<Timestamp, Option<Timestamp>>,
462 max_filter_num_per_query: usize,
464 time_window_merge_threshold: usize,
467}
468
469impl DirtyTimeWindows {
470 pub fn new(max_filter_num_per_query: usize, time_window_merge_threshold: usize) -> Self {
471 Self {
472 windows: BTreeMap::new(),
473 max_filter_num_per_query,
474 time_window_merge_threshold,
475 }
476 }
477
478 #[cfg(test)]
479 pub(crate) fn max_filter_num_per_query(&self) -> usize {
480 self.max_filter_num_per_query
481 }
482
483 #[cfg(test)]
484 pub(crate) fn time_window_merge_threshold(&self) -> usize {
485 self.time_window_merge_threshold
486 }
487}
488
489impl Default for DirtyTimeWindows {
490 fn default() -> Self {
491 Self {
492 windows: BTreeMap::new(),
493 max_filter_num_per_query: 20,
494 time_window_merge_threshold: 3,
495 }
496 }
497}
498
499impl DirtyTimeWindows {
500 pub const MERGE_DIST: i32 = 3;
504
505 pub fn add_lower_bounds(&mut self, lower_bounds: impl Iterator<Item = Timestamp>) {
511 for lower_bound in lower_bounds {
512 let entry = self.windows.entry(lower_bound);
513 entry.or_insert(None);
514 }
515 }
516
517 pub fn window_size(&self) -> Duration {
518 let mut ret = Duration::from_secs(0);
519 for (start, end) in &self.windows {
520 if let Some(end) = end
521 && let Some(duration) = end.sub(start)
522 {
523 ret += duration.to_std().unwrap_or_default();
524 }
525 }
526 ret
527 }
528
529 pub fn add_window(&mut self, start: Timestamp, end: Option<Timestamp>) {
530 self.add_or_merge_window(start, end);
531 }
532
533 pub fn add_windows(&mut self, time_ranges: Vec<(Timestamp, Timestamp)>) {
534 for (start, end) in time_ranges {
535 self.add_or_merge_window(start, Some(end));
536 }
537 }
538
539 pub fn add_dirty_windows(&mut self, dirty_windows: &DirtyTimeWindows) {
541 for (start, end) in &dirty_windows.windows {
542 self.add_or_merge_window(*start, *end);
543 }
544 }
545
546 fn add_or_merge_window(&mut self, start: Timestamp, end: Option<Timestamp>) {
547 self.windows
548 .entry(start)
549 .and_modify(|current_end| {
550 *current_end = Self::union_window_end(*current_end, end);
551 })
552 .or_insert(end);
553 }
554
555 fn union_window_end(
556 current_end: Option<Timestamp>,
557 incoming_end: Option<Timestamp>,
558 ) -> Option<Timestamp> {
559 match (current_end, incoming_end) {
560 (Some(current), Some(incoming)) => Some(current.max(incoming)),
561 (Some(end), None) | (None, Some(end)) => Some(end),
565 (None, None) => None,
566 }
567 }
568
569 pub(crate) fn detach(&mut self) -> Self {
571 Self {
572 windows: std::mem::take(&mut self.windows),
573 max_filter_num_per_query: self.max_filter_num_per_query,
574 time_window_merge_threshold: self.time_window_merge_threshold,
575 }
576 }
577
578 pub fn clean(&mut self) {
580 self.windows.clear();
581 }
582
583 pub fn set_dirty(&mut self) {
586 self.add_or_merge_window(Timestamp::new_second(0), None);
587 }
588
589 pub fn len(&self) -> usize {
591 self.windows.len()
592 }
593
594 pub fn is_empty(&self) -> bool {
595 self.windows.is_empty()
596 }
597
598 pub fn effective_count(&self, window_size: &Duration) -> usize {
601 if self.windows.is_empty() {
602 return 0;
603 }
604 let window_size =
605 chrono::Duration::from_std(*window_size).unwrap_or(chrono::Duration::zero());
606 let total_window_time_range =
607 self.windows
608 .iter()
609 .fold(chrono::Duration::zero(), |acc, (start, end)| {
610 if let Some(end) = end {
611 acc + end.sub(start).unwrap_or(chrono::Duration::zero())
612 } else {
613 acc + window_size
614 }
615 });
616
617 if window_size.num_seconds() == 0 {
619 0
620 } else {
621 (total_window_time_range.num_seconds() / window_size.num_seconds()) as usize
622 }
623 }
624
625 pub fn gen_filter_exprs(
631 &mut self,
632 col_name: &str,
633 expire_lower_bound: Option<Timestamp>,
634 window_size: chrono::Duration,
635 window_cnt: usize,
636 flow_id: FlowId,
637 task_ctx: Option<&BatchingTask>,
638 ) -> Result<Option<FilterExprInfo>, Error> {
639 ensure!(
640 window_size.num_seconds() > 0,
641 UnexpectedSnafu {
642 reason: "window_size is zero, can't generate filter exprs",
643 }
644 );
645
646 debug!(
647 "expire_lower_bound: {:?}, window_size: {:?}",
648 expire_lower_bound.map(|t| t.to_iso8601_string()),
649 window_size
650 );
651 self.merge_dirty_time_windows(window_size, expire_lower_bound)?;
652
653 if self.windows.len() > window_cnt {
654 let first_time_window = self.windows.first_key_value();
655 let last_time_window = self.windows.last_key_value();
656
657 if let Some(task_ctx) = task_ctx {
658 debug!(
659 "Flow id = {:?}, too many time windows: {}, only the first {} are taken for this query, the group by expression might be wrong. Time window expr={:?}, expire_after={:?}, first_time_window={:?}, last_time_window={:?}, the original query: {:?}",
660 task_ctx.config.flow_id,
661 self.windows.len(),
662 window_cnt,
663 task_ctx.config.time_window_expr,
664 task_ctx.config.expire_after,
665 first_time_window,
666 last_time_window,
667 task_ctx.config.query
668 );
669 } else {
670 debug!(
671 "Flow id = {:?}, too many time windows: {}, only the first {} are taken for this query, the group by expression might be wrong. first_time_window={:?}, last_time_window={:?}",
672 flow_id,
673 self.windows.len(),
674 window_cnt,
675 first_time_window,
676 last_time_window
677 )
678 }
679 }
680
681 let max_time_range = window_size * window_cnt as i32;
683
684 let mut to_be_query = BTreeMap::new();
685 let mut new_windows = self.windows.clone();
686 let mut cur_time_range = chrono::Duration::zero();
687 for (idx, (start, end)) in self.windows.iter().enumerate() {
688 let first_end = start
689 .add_duration(window_size.to_std().unwrap())
690 .context(TimeSnafu)?;
691 let end = end.unwrap_or(first_end);
692
693 if cur_time_range >= max_time_range {
695 break;
696 }
697
698 if idx >= window_cnt {
700 break;
701 }
702
703 let Some(x) = end.sub(start) else {
704 continue;
705 };
706 if cur_time_range + x <= max_time_range {
707 to_be_query.insert(*start, Some(end));
708 new_windows.remove(start);
709 cur_time_range += x;
710 } else {
711 let surplus = max_time_range - cur_time_range;
714 if surplus.num_seconds() <= window_size.num_seconds() {
715 break;
717 }
718 let times = surplus.num_seconds() / window_size.num_seconds();
719
720 let split_offset = window_size * times as i32;
721 let split_at = start
722 .add_duration(split_offset.to_std().unwrap())
723 .context(TimeSnafu)?;
724 to_be_query.insert(*start, Some(split_at));
725
726 new_windows.remove(start);
728 new_windows.insert(split_at, Some(end));
729 cur_time_range += split_offset;
730 break;
731 }
732 }
733
734 self.windows = new_windows;
735
736 METRIC_FLOW_BATCHING_ENGINE_QUERY_WINDOW_CNT
737 .with_label_values(&[flow_id.to_string().as_str()])
738 .observe(to_be_query.len() as f64);
739
740 let full_time_range = to_be_query
741 .iter()
742 .fold(chrono::Duration::zero(), |acc, (start, end)| {
743 if let Some(end) = end {
744 acc + end.sub(start).unwrap_or(chrono::Duration::zero())
745 } else {
746 acc + window_size
747 }
748 })
749 .num_seconds() as f64;
750 METRIC_FLOW_BATCHING_ENGINE_QUERY_WINDOW_SIZE
751 .with_label_values(&[flow_id.to_string().as_str()])
752 .observe(full_time_range);
753
754 let stalled_time_range =
755 self.windows
756 .iter()
757 .fold(chrono::Duration::zero(), |acc, (start, end)| {
758 if let Some(end) = end {
759 acc + end.sub(start).unwrap_or(chrono::Duration::zero())
760 } else {
761 acc + window_size
762 }
763 });
764
765 METRIC_FLOW_BATCHING_ENGINE_STALLED_WINDOW_SIZE
766 .with_label_values(&[flow_id.to_string().as_str()])
767 .observe(stalled_time_range.num_seconds() as f64);
768
769 let std_window_size = window_size.to_std().map_err(|e| {
770 InternalSnafu {
771 reason: e.to_string(),
772 }
773 .build()
774 })?;
775
776 let mut expr_lst = vec![];
777 let mut time_ranges = vec![];
778 for (start, end) in to_be_query.into_iter() {
779 let (start, end) = if let Some(ctx) = task_ctx {
781 let Some(time_window_expr) = &ctx.config.time_window_expr else {
782 UnexpectedSnafu {
783 reason: "time_window_expr is not set",
784 }
785 .fail()?
786 };
787 Self::align_time_window(start, end, time_window_expr)?
788 } else {
789 (start, end)
790 };
791 let end = end.unwrap_or(start.add_duration(std_window_size).context(TimeSnafu)?);
792 time_ranges.push((start, end));
793
794 debug!(
795 "Time window start: {:?}, end: {:?}",
796 start.to_iso8601_string(),
797 end.to_iso8601_string()
798 );
799
800 use datafusion_expr::{col, lit};
801 let lower = to_df_literal(start)?;
802 let upper = to_df_literal(end)?;
803 let expr = col(col_name)
804 .gt_eq(lit(lower))
805 .and(col(col_name).lt(lit(upper)));
806 expr_lst.push(expr);
807 }
808 let expr = expr_lst.into_iter().reduce(|a, b| a.or(b));
809 let ret = expr.map(|expr| FilterExprInfo {
810 expr,
811 col_name: col_name.to_string(),
812 time_ranges,
813 window_size,
814 });
815 Ok(ret)
816 }
817
818 pub(crate) fn align_time_window(
821 start: Timestamp,
822 end: Option<Timestamp>,
823 time_window_expr: &TimeWindowExpr,
824 ) -> Result<(Timestamp, Option<Timestamp>), Error> {
825 let align_start = time_window_expr.eval(start)?.0.context(UnexpectedSnafu {
826 reason: format!(
827 "Failed to align start time {:?} with time window expr {:?}",
828 start, time_window_expr
829 ),
830 })?;
831 let align_end = end
832 .and_then(|end| {
833 time_window_expr
834 .eval(end)
835 .map(|r| if r.0 == Some(end) { r.0 } else { r.1 })
837 .transpose()
838 })
839 .transpose()?;
840 Ok((align_start, align_end))
841 }
842
843 pub fn merge_dirty_time_windows(
847 &mut self,
848 window_size: chrono::Duration,
849 expire_lower_bound: Option<Timestamp>,
850 ) -> Result<(), Error> {
851 if self.windows.is_empty() {
852 return Ok(());
853 }
854
855 let mut new_windows = BTreeMap::new();
856
857 let std_window_size = window_size.to_std().map_err(|e| {
858 InternalSnafu {
859 reason: e.to_string(),
860 }
861 .build()
862 })?;
863
864 let mut prev_tw = None;
866 for (mut lower_bound, upper_bound) in std::mem::take(&mut self.windows) {
867 if let Some(expire_lower_bound) = expire_lower_bound {
869 match upper_bound {
870 Some(upper_bound) if upper_bound <= expire_lower_bound => continue,
873 Some(_) if lower_bound < expire_lower_bound => {
878 lower_bound = expire_lower_bound;
879 }
880 None if lower_bound < expire_lower_bound => continue,
882 _ => {}
883 }
884 }
885
886 let Some(prev_tw) = &mut prev_tw else {
887 prev_tw = Some((lower_bound, upper_bound));
888 continue;
889 };
890
891 let prev_upper = prev_tw
894 .1
895 .unwrap_or(prev_tw.0.add_duration(std_window_size).context(TimeSnafu)?);
896 prev_tw.1 = Some(prev_upper);
897
898 let cur_upper = upper_bound.unwrap_or(
899 lower_bound
900 .add_duration(std_window_size)
901 .context(TimeSnafu)?,
902 );
903
904 if lower_bound
905 .sub(&prev_upper)
906 .map(|dist| dist <= window_size * self.time_window_merge_threshold as i32)
907 .unwrap_or(false)
908 {
909 prev_tw.1 = Some(prev_upper.max(cur_upper));
912 } else {
913 new_windows.insert(prev_tw.0, prev_tw.1);
914 *prev_tw = (lower_bound, Some(cur_upper));
915 }
916 }
917
918 if let Some(prev_tw) = prev_tw {
919 new_windows.insert(prev_tw.0, prev_tw.1);
920 }
921
922 self.windows = new_windows;
923
924 Ok(())
925 }
926}
927
928pub(crate) fn to_df_literal(value: Timestamp) -> Result<datafusion_common::ScalarValue, Error> {
929 let value = Value::from(value);
930 let value = value
931 .try_to_scalar_value(&value.data_type())
932 .with_context(|_| DatatypesSnafu {
933 extra: format!("Failed to convert to scalar value: {}", value),
934 })?;
935 Ok(value)
936}
937
938#[derive(Debug, Clone)]
939enum ExecState {
940 Idle,
941 Executing,
942}
943
944#[derive(Debug, Clone, Copy, PartialEq, Eq)]
945pub enum CheckpointMode {
946 FullSnapshot,
947 Incremental,
948}
949
950#[derive(Debug, Clone)]
953pub struct FencedRepair {
954 high: BTreeMap<u64, u64>,
955 pending_windows: DirtyTimeWindows,
956}
957
958impl FencedRepair {
959 pub fn high(&self) -> &BTreeMap<u64, u64> {
961 &self.high
962 }
963
964 pub fn pending_windows(&self) -> &DirtyTimeWindows {
966 &self.pending_windows
967 }
968}
969
970#[derive(Debug, Clone)]
972pub struct FilterExprInfo {
973 pub expr: datafusion_expr::Expr,
974 pub col_name: String,
975 pub time_ranges: Vec<(Timestamp, Timestamp)>,
976 pub window_size: chrono::Duration,
977}
978
979impl FilterExprInfo {
980 pub fn total_window_length(&self) -> chrono::Duration {
981 self.time_ranges
982 .iter()
983 .fold(chrono::Duration::zero(), |acc, (start, end)| {
984 acc + end.sub(start).unwrap_or(chrono::Duration::zero())
985 })
986 }
987
988 pub fn predicate_for_col(
989 &self,
990 col_name: &str,
991 ) -> Result<Option<datafusion_expr::Expr>, Error> {
992 use datafusion_common::Column;
993 use datafusion_expr::{Expr, lit};
994
995 let mut expr_lst = Vec::with_capacity(self.time_ranges.len());
996 for (start, end) in &self.time_ranges {
997 let lower = to_df_literal(*start)?;
998 let upper = to_df_literal(*end)?;
999 let filter_col = || Expr::Column(Column::new_unqualified(col_name));
1000 expr_lst.push(
1001 filter_col()
1002 .gt_eq(lit(lower))
1003 .and(filter_col().lt(lit(upper))),
1004 );
1005 }
1006
1007 Ok(expr_lst.into_iter().reduce(|a, b| a.or(b)))
1008 }
1009}
1010
1011#[cfg(test)]
1012mod test {
1013 use pretty_assertions::assert_eq;
1014 use session::context::QueryContext;
1015
1016 use super::*;
1017 use crate::batching_mode::time_window::find_time_window_expr;
1018 use crate::batching_mode::utils::sql_to_df_plan;
1019 use crate::test_utils::create_test_query_engine;
1020
1021 #[test]
1022 fn test_task_state_records_last_execution_time() {
1023 let query_ctx = QueryContext::arc();
1024 let (_tx, rx) = tokio::sync::oneshot::channel();
1025 let mut state = TaskState::new(query_ctx, rx);
1026
1027 assert_eq!(None, state.last_execution_time_millis());
1028 state.after_query_exec(std::time::Duration::from_millis(1), false);
1029 assert_eq!(None, state.last_execution_time_millis());
1030
1031 state.after_query_exec(std::time::Duration::from_millis(1), true);
1032 assert!(state.last_execution_time_millis().is_some());
1033 }
1034
1035 #[test]
1036 fn test_merge_dirty_time_windows() {
1037 let merge_dist = DirtyTimeWindows::default().time_window_merge_threshold;
1038 let testcases = vec![
1039 (
1041 vec![
1042 Timestamp::new_second(0),
1043 Timestamp::new_second((1 + merge_dist as i64) * 5 * 60),
1044 ],
1045 (chrono::Duration::seconds(5 * 60), None),
1046 BTreeMap::from([(
1047 Timestamp::new_second(0),
1048 Some(Timestamp::new_second((2 + merge_dist as i64) * 5 * 60)),
1049 )]),
1050 Some(
1051 "((ts >= CAST('1970-01-01 00:00:00' AS TIMESTAMP)) AND (ts < CAST('1970-01-01 00:25:00' AS TIMESTAMP)))",
1052 ),
1053 ),
1054 (
1056 vec![
1057 Timestamp::new_second(0),
1058 Timestamp::new_second((2 + merge_dist as i64) * 5 * 60),
1059 ],
1060 (chrono::Duration::seconds(5 * 60), None),
1061 BTreeMap::from([
1062 (
1063 Timestamp::new_second(0),
1064 Some(Timestamp::new_second(5 * 60)),
1065 ),
1066 (
1067 Timestamp::new_second((2 + merge_dist as i64) * 5 * 60),
1068 Some(Timestamp::new_second((3 + merge_dist as i64) * 5 * 60)),
1069 ),
1070 ]),
1071 Some(
1072 "(((ts >= CAST('1970-01-01 00:00:00' AS TIMESTAMP)) AND (ts < CAST('1970-01-01 00:05:00' AS TIMESTAMP))) OR ((ts >= CAST('1970-01-01 00:25:00' AS TIMESTAMP)) AND (ts < CAST('1970-01-01 00:30:00' AS TIMESTAMP))))",
1073 ),
1074 ),
1075 (
1077 vec![
1078 Timestamp::new_second(0),
1079 Timestamp::new_second((merge_dist as i64) * 5 * 60),
1080 ],
1081 (chrono::Duration::seconds(5 * 60), None),
1082 BTreeMap::from([(
1083 Timestamp::new_second(0),
1084 Some(Timestamp::new_second((1 + merge_dist as i64) * 5 * 60)),
1085 )]),
1086 Some(
1087 "((ts >= CAST('1970-01-01 00:00:00' AS TIMESTAMP)) AND (ts < CAST('1970-01-01 00:20:00' AS TIMESTAMP)))",
1088 ),
1089 ),
1090 (
1092 vec![
1093 Timestamp::new_second(0),
1094 Timestamp::new_second((merge_dist as i64) * 3),
1095 Timestamp::new_second((merge_dist as i64) * 3 * 2),
1096 ],
1097 (chrono::Duration::seconds(3), None),
1098 BTreeMap::from([(
1099 Timestamp::new_second(0),
1100 Some(Timestamp::new_second((merge_dist as i64) * 7)),
1101 )]),
1102 Some(
1103 "((ts >= CAST('1970-01-01 00:00:00' AS TIMESTAMP)) AND (ts < CAST('1970-01-01 00:00:21' AS TIMESTAMP)))",
1104 ),
1105 ),
1106 (
1108 Vec::from_iter((0..20).map(|i| Timestamp::new_second(i * 3)).chain(
1109 std::iter::once(Timestamp::new_second(
1110 60 + 3 * (DirtyTimeWindows::MERGE_DIST as i64 + 1),
1111 )),
1112 )),
1113 (chrono::Duration::seconds(3), None),
1114 BTreeMap::from([
1115 (Timestamp::new_second(0), Some(Timestamp::new_second(60))),
1116 (
1117 Timestamp::new_second(60 + 3 * (DirtyTimeWindows::MERGE_DIST as i64 + 1)),
1118 Some(Timestamp::new_second(
1119 60 + 3 * (DirtyTimeWindows::MERGE_DIST as i64 + 1) + 3,
1120 )),
1121 ),
1122 ]),
1123 Some(
1124 "((ts >= CAST('1970-01-01 00:00:00' AS TIMESTAMP)) AND (ts < CAST('1970-01-01 00:01:00' AS TIMESTAMP)))",
1125 ),
1126 ),
1127 (
1129 Vec::from_iter((0..40).map(|i| Timestamp::new_second(i * 3))),
1130 (chrono::Duration::seconds(3), None),
1131 BTreeMap::from([(
1132 Timestamp::new_second(0),
1133 Some(Timestamp::new_second(40 * 3)),
1134 )]),
1135 Some(
1136 "((ts >= CAST('1970-01-01 00:00:00' AS TIMESTAMP)) AND (ts < CAST('1970-01-01 00:01:00' AS TIMESTAMP)))",
1137 ),
1138 ),
1139 (
1141 Vec::from_iter(
1142 std::iter::once(Timestamp::new_second(0))
1143 .chain((0..40).map(|i| Timestamp::new_second(20 + i * 3))),
1144 ),
1145 (chrono::Duration::seconds(3), None),
1146 BTreeMap::from([
1147 (Timestamp::new_second(0), Some(Timestamp::new_second(3))),
1148 (Timestamp::new_second(20), Some(Timestamp::new_second(140))),
1149 ]),
1150 Some(
1151 "(((ts >= CAST('1970-01-01 00:00:00' AS TIMESTAMP)) AND (ts < CAST('1970-01-01 00:00:03' AS TIMESTAMP))) OR ((ts >= CAST('1970-01-01 00:00:20' AS TIMESTAMP)) AND (ts < CAST('1970-01-01 00:01:17' AS TIMESTAMP))))",
1152 ),
1153 ),
1154 (
1156 vec![
1157 Timestamp::new_second(0),
1158 Timestamp::new_second((merge_dist as i64) * 5 * 60),
1159 ],
1160 (
1161 chrono::Duration::seconds(5 * 60),
1162 Some(Timestamp::new_second((merge_dist as i64) * 6 * 60)),
1163 ),
1164 BTreeMap::from([]),
1165 None,
1166 ),
1167 ];
1168 for (lower_bounds, (window_size, expire_lower_bound), expected, expected_filter_expr) in
1171 testcases
1172 {
1173 let mut dirty = DirtyTimeWindows::default();
1174 dirty.add_lower_bounds(lower_bounds.into_iter());
1175 dirty
1176 .merge_dirty_time_windows(window_size, expire_lower_bound)
1177 .unwrap();
1178 assert_eq!(expected, dirty.windows);
1179 let filter_expr = dirty
1180 .gen_filter_exprs(
1181 "ts",
1182 expire_lower_bound,
1183 window_size,
1184 dirty.max_filter_num_per_query,
1185 0,
1186 None,
1187 )
1188 .unwrap()
1189 .map(|e| e.expr);
1190
1191 let unparser = datafusion::sql::unparser::Unparser::default();
1192 let to_sql = filter_expr
1193 .as_ref()
1194 .map(|e| unparser.expr_to_sql(e).unwrap().to_string());
1195 assert_eq!(expected_filter_expr, to_sql.as_deref());
1196 }
1197 }
1198
1199 #[test]
1200 fn test_merge_dirty_time_windows_with_bounded_ranges() {
1201 let window_size = chrono::Duration::seconds(5);
1202 let testcases = vec![
1203 (
1206 vec![
1207 (Timestamp::new_second(0), Some(Timestamp::new_second(15))),
1208 (Timestamp::new_second(5), Some(Timestamp::new_second(10))),
1209 ],
1210 BTreeMap::from([(Timestamp::new_second(0), Some(Timestamp::new_second(15)))]),
1211 ),
1212 (
1216 vec![
1217 (Timestamp::new_second(0), Some(Timestamp::new_second(15))),
1218 (Timestamp::new_second(3), None),
1219 ],
1220 BTreeMap::from([(Timestamp::new_second(0), Some(Timestamp::new_second(15)))]),
1221 ),
1222 (
1224 vec![
1225 (Timestamp::new_second(0), Some(Timestamp::new_second(5))),
1226 (Timestamp::new_second(100), Some(Timestamp::new_second(110))),
1227 ],
1228 BTreeMap::from([
1229 (Timestamp::new_second(0), Some(Timestamp::new_second(5))),
1230 (Timestamp::new_second(100), Some(Timestamp::new_second(110))),
1231 ]),
1232 ),
1233 (
1236 vec![
1237 (Timestamp::new_second(0), Some(Timestamp::new_second(10))),
1238 (Timestamp::new_second(5), Some(Timestamp::new_second(20))),
1239 ],
1240 BTreeMap::from([(Timestamp::new_second(0), Some(Timestamp::new_second(20)))]),
1241 ),
1242 ];
1243
1244 for (windows, expected) in testcases {
1245 let mut dirty = DirtyTimeWindows::default();
1246 for (start, end) in windows {
1247 dirty.add_window(start, end);
1248 }
1249 dirty.merge_dirty_time_windows(window_size, None).unwrap();
1250 assert_eq!(expected, dirty.windows);
1251 }
1252
1253 let expire_testcases = vec![
1255 (
1257 vec![(Timestamp::new_second(0), Some(Timestamp::new_second(10)))],
1258 BTreeMap::from([]),
1259 ),
1260 (
1262 vec![(Timestamp::new_second(0), Some(Timestamp::new_second(5)))],
1263 BTreeMap::from([]),
1264 ),
1265 (
1268 vec![(Timestamp::new_second(0), Some(Timestamp::new_second(15)))],
1269 BTreeMap::from([(Timestamp::new_second(10), Some(Timestamp::new_second(15)))]),
1270 ),
1271 (
1273 vec![(Timestamp::new_second(10), Some(Timestamp::new_second(15)))],
1274 BTreeMap::from([(Timestamp::new_second(10), Some(Timestamp::new_second(15)))]),
1275 ),
1276 (vec![(Timestamp::new_second(5), None)], BTreeMap::from([])),
1279 (
1281 vec![(Timestamp::new_second(10), None)],
1282 BTreeMap::from([(Timestamp::new_second(10), None)]),
1283 ),
1284 ];
1285
1286 for (windows, expected) in expire_testcases {
1287 let mut dirty = DirtyTimeWindows::default();
1288 for (start, end) in windows {
1289 dirty.add_window(start, end);
1290 }
1291 dirty
1292 .merge_dirty_time_windows(window_size, Some(Timestamp::new_second(10)))
1293 .unwrap();
1294 assert_eq!(expected, dirty.windows);
1295 }
1296 }
1297
1298 #[tokio::test]
1299 async fn test_align_time_window() {
1300 type TimeWindow = (Timestamp, Option<Timestamp>);
1301 struct TestCase {
1302 sql: String,
1303 aligns: Vec<(TimeWindow, TimeWindow)>,
1304 }
1305 let testcases: Vec<TestCase> = vec![TestCase{
1306 sql: "SELECT date_bin(INTERVAL '5 second', ts) AS time_window FROM numbers_with_ts GROUP BY time_window;".to_string(),
1307 aligns: vec![
1308 ((Timestamp::new_second(3), None), (Timestamp::new_second(0), None)),
1309 ((Timestamp::new_second(8), None), (Timestamp::new_second(5), None)),
1310 ((Timestamp::new_second(8), Some(Timestamp::new_second(10))), (Timestamp::new_second(5), Some(Timestamp::new_second(10)))),
1311 ((Timestamp::new_second(8), Some(Timestamp::new_second(9))), (Timestamp::new_second(5), Some(Timestamp::new_second(10)))),
1312 ],
1313 }];
1314
1315 let query_engine = create_test_query_engine();
1316 let ctx = QueryContext::arc();
1317 for TestCase { sql, aligns } in testcases {
1318 let plan = sql_to_df_plan(ctx.clone(), query_engine.clone(), &sql, true)
1319 .await
1320 .unwrap();
1321
1322 let (column_name, time_window_expr, _, df_schema) = find_time_window_expr(
1323 &plan,
1324 query_engine.engine_state().catalog_manager().clone(),
1325 ctx.clone(),
1326 )
1327 .await
1328 .unwrap();
1329
1330 let time_window_expr = time_window_expr
1331 .map(|expr| {
1332 TimeWindowExpr::from_expr(
1333 &expr,
1334 &column_name,
1335 &df_schema,
1336 &query_engine.engine_state().session_state(),
1337 )
1338 })
1339 .transpose()
1340 .unwrap()
1341 .unwrap();
1342
1343 for (before_align, expected_after_align) in aligns {
1344 let after_align = DirtyTimeWindows::align_time_window(
1345 before_align.0,
1346 before_align.1,
1347 &time_window_expr,
1348 )
1349 .unwrap();
1350 assert_eq!(expected_after_align, after_align);
1351 }
1352 }
1353 }
1354
1355 #[test]
1356 fn test_task_state_checkpoint_mode_and_advancement() {
1357 let query_ctx = QueryContext::arc();
1358 let (_tx, rx) = tokio::sync::oneshot::channel();
1359 let mut state = TaskState::new(query_ctx, rx);
1360
1361 assert_eq!(state.checkpoint_mode(), CheckpointMode::FullSnapshot);
1362 assert!(state.checkpoints().is_empty());
1363
1364 state.advance_checkpoints(HashMap::from([(1_u64, 10_u64), (2_u64, 20_u64)]));
1365 assert_eq!(state.checkpoint_mode(), CheckpointMode::Incremental);
1366 assert_eq!(
1367 state.checkpoints(),
1368 &BTreeMap::from([(1_u64, 10_u64), (2_u64, 20_u64)])
1369 );
1370
1371 state.mark_full_snapshot();
1372 assert_eq!(state.checkpoint_mode(), CheckpointMode::FullSnapshot);
1373 assert_eq!(
1374 state.checkpoints(),
1375 &BTreeMap::from([(1_u64, 10_u64), (2_u64, 20_u64)])
1376 );
1377 }
1378
1379 #[test]
1380 fn test_mark_full_snapshot_restores_pending_fenced_repair_windows() {
1381 let query_ctx = QueryContext::arc();
1382 let (_tx, rx) = tokio::sync::oneshot::channel();
1383 let mut state = TaskState::new(query_ctx, rx);
1384 state
1385 .dirty_time_windows
1386 .add_window(Timestamp::new_second(10), Some(Timestamp::new_second(15)));
1387 state
1388 .dirty_time_windows
1389 .add_window(Timestamp::new_second(100), Some(Timestamp::new_second(105)));
1390
1391 state
1392 .start_fenced_repair(BTreeMap::from([(1_u64, 10_u64)]))
1393 .unwrap();
1394 assert!(state.dirty_time_windows.is_empty());
1395 assert_eq!(
1396 state
1397 .pending_fenced_repair()
1398 .unwrap()
1399 .pending_windows()
1400 .len(),
1401 2
1402 );
1403
1404 state.mark_full_snapshot();
1405
1406 assert_eq!(state.checkpoint_mode(), CheckpointMode::FullSnapshot);
1407 assert!(state.pending_fenced_repair().is_none());
1408 assert_eq!(state.dirty_time_windows.len(), 2);
1409 }
1410
1411 #[test]
1412 fn test_explicit_fenced_repair_windows_keep_live_windows_separate() {
1413 let mut state = state_with_past_update(Duration::from_secs(1));
1414 state
1415 .dirty_time_windows
1416 .add_window(Timestamp::new_second(0), Some(Timestamp::new_second(1_000)));
1417 let high = BTreeMap::from([(1, 10)]);
1418 state.start_fenced_repair_windows(
1419 high.clone(),
1420 vec![(Timestamp::new_second(10), Timestamp::new_second(15))],
1421 );
1422
1423 assert_eq!(state.checkpoint_mode(), CheckpointMode::FullSnapshot);
1424 assert_eq!(state.pending_fenced_repair().unwrap().high(), &high);
1425 assert_eq!(
1426 state
1427 .pending_fenced_repair()
1428 .unwrap()
1429 .pending_windows()
1430 .len(),
1431 1
1432 );
1433 assert_eq!(state.dirty_time_windows.len(), 1);
1434
1435 let filter = state
1436 .gen_scoped_filter_exprs(
1437 "ts",
1438 Some(Timestamp::new_second(100)),
1439 chrono::Duration::seconds(5),
1440 1,
1441 1,
1442 None,
1443 )
1444 .unwrap()
1445 .unwrap();
1446 assert_eq!(
1447 filter.time_ranges,
1448 vec![(Timestamp::new_second(10), Timestamp::new_second(15))]
1449 );
1450 state
1451 .dirty_time_windows
1452 .add_window(Timestamp::new_second(20), Some(Timestamp::new_second(25)));
1453 state.restore_scoped_windows(&filter);
1454
1455 assert_eq!(
1456 state
1457 .pending_fenced_repair()
1458 .unwrap()
1459 .pending_windows()
1460 .len(),
1461 1
1462 );
1463 assert_eq!(state.dirty_time_windows.len(), 2);
1464 }
1465
1466 #[test]
1467 fn test_explicit_fenced_repair_keeps_empty_high_and_old_windows() {
1468 let mut state = state_with_past_update(Duration::from_secs(1));
1469 state
1470 .dirty_time_windows
1471 .add_window(Timestamp::new_second(0), Some(Timestamp::new_second(1_000)));
1472 let high = BTreeMap::from([(1, 10)]);
1473 state.start_fenced_repair_windows(high.clone(), Vec::new());
1474
1475 assert!(
1476 state
1477 .gen_scoped_filter_exprs(
1478 "ts",
1479 Some(Timestamp::new_second(100)),
1480 chrono::Duration::seconds(5),
1481 1,
1482 1,
1483 None,
1484 )
1485 .unwrap()
1486 .is_none()
1487 );
1488 assert_eq!(state.pending_fenced_repair().unwrap().high(), &high);
1489 assert_eq!(state.dirty_time_windows.len(), 1);
1490 assert_eq!(state.finish_fenced_repair(), Some(high));
1491
1492 state.start_fenced_repair_windows(
1493 BTreeMap::from([(1, 11)]),
1494 vec![(Timestamp::new_second(-20), Timestamp::new_second(-15))],
1495 );
1496 let filter = state
1497 .gen_scoped_filter_exprs(
1498 "ts",
1499 Some(Timestamp::new_second(100)),
1500 chrono::Duration::seconds(5),
1501 1,
1502 1,
1503 None,
1504 )
1505 .unwrap()
1506 .unwrap();
1507 assert_eq!(
1508 filter.time_ranges,
1509 vec![(Timestamp::new_second(-20), Timestamp::new_second(-15))]
1510 );
1511 }
1512
1513 #[test]
1514 fn test_disable_incremental_persists_full_snapshot_mode() {
1515 let query_ctx = QueryContext::arc();
1516 let (_tx, rx) = tokio::sync::oneshot::channel();
1517 let mut state = TaskState::new(query_ctx, rx);
1518
1519 assert!(!state.is_incremental_disabled());
1520
1521 state.disable_incremental();
1523 assert!(state.is_incremental_disabled());
1524 assert_eq!(state.checkpoint_mode(), CheckpointMode::FullSnapshot);
1525
1526 state.advance_checkpoints(HashMap::from([(1_u64, 10_u64), (2_u64, 20_u64)]));
1528 assert_eq!(state.checkpoint_mode(), CheckpointMode::FullSnapshot);
1529 assert_eq!(
1530 state.checkpoints(),
1531 &BTreeMap::from([(1_u64, 10_u64), (2_u64, 20_u64)])
1532 );
1533
1534 state.mark_full_snapshot();
1536 assert!(state.is_incremental_disabled());
1537 assert_eq!(state.checkpoint_mode(), CheckpointMode::FullSnapshot);
1538 }
1539
1540 #[test]
1541 fn test_full_snapshot_checkpoint_advancement_requires_participating_regions() {
1542 let query_ctx = QueryContext::arc();
1543 let (_tx, rx) = tokio::sync::oneshot::channel();
1544 let state = TaskState::new(query_ctx, rx);
1545
1546 assert!(!state.can_advance_full_snapshot_checkpoints(&BTreeSet::new(), &HashMap::new()));
1547 assert!(!state.can_advance_full_snapshot_checkpoints(
1548 &BTreeSet::from([1_u64, 2_u64]),
1549 &HashMap::from([(1_u64, 10_u64)]),
1550 ));
1551 assert!(state.can_advance_full_snapshot_checkpoints(
1552 &BTreeSet::from([1_u64, 2_u64]),
1553 &HashMap::from([(1_u64, 10_u64), (2_u64, 20_u64)]),
1554 ));
1555 }
1556
1557 #[test]
1558 fn test_incremental_checkpoint_advancement_requires_participation_alignment() {
1559 let query_ctx = QueryContext::arc();
1560 let (_tx, rx) = tokio::sync::oneshot::channel();
1561 let mut state = TaskState::new(query_ctx, rx);
1562 state.advance_checkpoints(HashMap::from([(1_u64, 10_u64), (2_u64, 20_u64)]));
1563
1564 assert!(
1565 state.can_advance_incremental_checkpoints_with_participation(
1566 &BTreeSet::from([1_u64]),
1567 &HashMap::from([(1_u64, 11_u64)]),
1568 )
1569 );
1570 assert!(
1571 !state.can_advance_incremental_checkpoints_with_participation(
1572 &BTreeSet::from([1_u64, 2_u64]),
1573 &HashMap::from([(1_u64, 11_u64)]),
1574 )
1575 );
1576 assert!(
1577 !state.can_advance_incremental_checkpoints_with_participation(
1578 &BTreeSet::from([3_u64]),
1579 &HashMap::from([(3_u64, 11_u64)]),
1580 )
1581 );
1582 assert!(
1583 !state.can_advance_incremental_checkpoints_with_participation(
1584 &BTreeSet::from([1_u64]),
1585 &HashMap::from([(1_u64, 9_u64)]),
1586 )
1587 );
1588 assert!(
1589 state.can_advance_incremental_checkpoints_with_participation(
1590 &BTreeSet::from([1_u64, 2_u64]),
1591 &HashMap::from([(1_u64, 11_u64), (2_u64, 21_u64)]),
1592 )
1593 );
1594
1595 state.disable_incremental();
1596 assert!(
1597 !state.can_advance_incremental_checkpoints_with_participation(
1598 &BTreeSet::from([1_u64, 2_u64]),
1599 &HashMap::from([(1_u64, 12_u64), (2_u64, 22_u64)]),
1600 )
1601 );
1602 }
1603
1604 #[test]
1605 fn test_incremental_checkpoint_advancement_merges_participating_subset() {
1606 let query_ctx = QueryContext::arc();
1607 let (_tx, rx) = tokio::sync::oneshot::channel();
1608 let mut state = TaskState::new(query_ctx, rx);
1609 state.advance_checkpoints(HashMap::from([
1610 (1_u64, 10_u64),
1611 (2_u64, 20_u64),
1612 (3_u64, 30_u64),
1613 ]));
1614
1615 state.advance_incremental_checkpoints_with_participation(
1616 &BTreeSet::from([1_u64, 3_u64]),
1617 HashMap::from([(1_u64, 12_u64), (3_u64, 35_u64)]),
1618 );
1619
1620 assert_eq!(state.checkpoint_mode(), CheckpointMode::Incremental);
1621 assert_eq!(
1622 state.checkpoints(),
1623 &BTreeMap::from([(1_u64, 12_u64), (2_u64, 20_u64), (3_u64, 35_u64)])
1624 );
1625 }
1626
1627 #[test]
1628 fn test_filter_expr_info_predicate_for_col_empty_ranges() {
1629 let filter = FilterExprInfo {
1630 expr: datafusion_expr::col("ts"),
1631 col_name: "ts".to_string(),
1632 time_ranges: vec![],
1633 window_size: chrono::Duration::seconds(1),
1634 };
1635
1636 assert!(filter.predicate_for_col("time_window").unwrap().is_none());
1637 }
1638
1639 #[test]
1640 fn test_filter_expr_info_predicate_for_col_single_range() {
1641 let filter = FilterExprInfo {
1642 expr: datafusion_expr::col("ts"),
1643 col_name: "ts".to_string(),
1644 time_ranges: vec![(Timestamp::new_second(0), Timestamp::new_second(1))],
1645 window_size: chrono::Duration::seconds(1),
1646 };
1647
1648 let predicate = filter.predicate_for_col("time_window").unwrap().unwrap();
1649 let unparser = datafusion::sql::unparser::Unparser::default();
1650 assert_eq!(
1651 "((time_window >= CAST('1970-01-01 00:00:00' AS TIMESTAMP)) AND (time_window < CAST('1970-01-01 00:00:01' AS TIMESTAMP)))",
1652 unparser.expr_to_sql(&predicate).unwrap().to_string()
1653 );
1654 }
1655
1656 #[test]
1657 fn test_filter_expr_info_predicate_for_col_multiple_ranges() {
1658 let filter = FilterExprInfo {
1659 expr: datafusion_expr::col("ts"),
1660 col_name: "ts".to_string(),
1661 time_ranges: vec![
1662 (Timestamp::new_second(0), Timestamp::new_second(1)),
1663 (Timestamp::new_second(10), Timestamp::new_second(11)),
1664 ],
1665 window_size: chrono::Duration::seconds(1),
1666 };
1667
1668 let predicate = filter.predicate_for_col("time_window").unwrap().unwrap();
1669 let unparser = datafusion::sql::unparser::Unparser::default();
1670 assert_eq!(
1671 "(((time_window >= CAST('1970-01-01 00:00:00' AS TIMESTAMP)) AND (time_window < CAST('1970-01-01 00:00:01' AS TIMESTAMP))) OR ((time_window >= CAST('1970-01-01 00:00:10' AS TIMESTAMP)) AND (time_window < CAST('1970-01-01 00:00:11' AS TIMESTAMP))))",
1672 unparser.expr_to_sql(&predicate).unwrap().to_string()
1673 );
1674 }
1675
1676 fn state_with_past_update(age: Duration) -> TaskState {
1678 let query_ctx = QueryContext::arc();
1679 let (_tx, rx) = tokio::sync::oneshot::channel();
1680 let mut state = TaskState::new(query_ctx, rx);
1681 state.last_update_time = Instant::now() - age;
1682 state
1683 }
1684
1685 #[test]
1686 fn test_short_incremental_cadence_uses_min_refresh() {
1687 let state = state_with_past_update(Duration::from_secs(10));
1691
1692 let time_window_size = Some(Duration::from_secs(60)); let min_refresh = Duration::from_secs(5);
1694 let flow_id = 1;
1695
1696 let result = state.get_next_start_query_time(
1697 flow_id,
1698 &time_window_size,
1699 min_refresh,
1700 None,
1701 20,
1702 true, );
1704
1705 let expected = state.last_update_time + min_refresh;
1707 assert_eq!(result, expected);
1708 }
1709
1710 #[test]
1711 fn test_short_incremental_cadence_respects_last_query_duration() {
1712 let mut state = state_with_past_update(Duration::from_secs(10));
1713 state.last_query_duration = Duration::from_secs(20);
1714
1715 let time_window_size = Some(Duration::from_secs(60));
1716 let min_refresh = Duration::from_secs(5);
1717 let flow_id = 1;
1718
1719 let result = state.get_next_start_query_time(
1720 flow_id,
1721 &time_window_size,
1722 min_refresh,
1723 None,
1724 20,
1725 true,
1726 );
1727
1728 assert_eq!(result, state.last_update_time + state.last_query_duration);
1729 }
1730
1731 #[test]
1732 fn test_short_incremental_cadence_respects_max_timeout() {
1733 let mut state = state_with_past_update(Duration::from_secs(10));
1734 state.last_query_duration = Duration::from_secs(20);
1735
1736 let time_window_size = Some(Duration::from_secs(60));
1737 let min_refresh = Duration::from_secs(30);
1738 let max_timeout = Duration::from_secs(5);
1739 let flow_id = 1;
1740
1741 let result = state.get_next_start_query_time(
1742 flow_id,
1743 &time_window_size,
1744 min_refresh,
1745 Some(max_timeout),
1746 20,
1747 true,
1748 );
1749
1750 assert_eq!(result, state.last_update_time + max_timeout);
1751 }
1752
1753 #[test]
1754 fn test_full_snapshot_ignores_short_cadence() {
1755 let mut state = state_with_past_update(Duration::from_secs(10));
1758 state.last_query_duration = Duration::from_secs(1);
1760
1761 let time_window_size = Some(Duration::from_secs(60)); let min_refresh = Duration::from_secs(5);
1763 let flow_id = 1;
1764
1765 let result = state.get_next_start_query_time(
1766 flow_id,
1767 &time_window_size,
1768 min_refresh,
1769 None,
1770 20,
1771 false, );
1773
1774 let expected = state.last_update_time + Duration::from_secs(60);
1777 assert_eq!(result, expected);
1778 }
1779
1780 #[test]
1781 fn test_dirty_window_overflow_schedules_immediately_even_with_short_cadence() {
1782 let mut state = state_with_past_update(Duration::from_secs(10));
1785 state
1787 .dirty_time_windows
1788 .add_window(Timestamp::new_second(0), Some(Timestamp::new_second(3600)));
1789
1790 let time_window_size = Some(Duration::from_secs(1)); let min_refresh = Duration::from_secs(5);
1792 let flow_id = 1;
1793
1794 let result = state.get_next_start_query_time(
1796 flow_id,
1797 &time_window_size,
1798 min_refresh,
1799 None,
1800 1, true,
1802 );
1803 assert!(
1804 result <= Instant::now(),
1805 "dirty overflow should schedule immediately"
1806 );
1807
1808 let result2 = state.get_next_start_query_time(
1810 flow_id,
1811 &time_window_size,
1812 min_refresh,
1813 None,
1814 1,
1815 false,
1816 );
1817 assert!(
1818 result2 <= Instant::now(),
1819 "dirty overflow should schedule immediately"
1820 );
1821 }
1822
1823 #[test]
1824 fn test_pending_fenced_repair_schedules_immediately() {
1825 let mut state = state_with_past_update(Duration::from_secs(10));
1826 state
1827 .dirty_time_windows
1828 .add_window(Timestamp::new_second(0), Some(Timestamp::new_second(5)));
1829 state
1830 .start_fenced_repair(BTreeMap::from([(1_u64, 10_u64)]))
1831 .unwrap();
1832 assert!(state.dirty_time_windows.is_empty());
1833 assert!(!state.fenced_repair_pending_is_empty());
1834
1835 let result = state.get_next_start_query_time(
1836 1,
1837 &Some(Duration::from_secs(60)),
1838 Duration::from_secs(5),
1839 None,
1840 20,
1841 false,
1842 );
1843
1844 assert!(
1845 result <= Instant::now(),
1846 "pending fenced repair backlog should schedule immediately"
1847 );
1848 }
1849
1850 #[test]
1851 fn test_incremental_disabled_ignores_short_cadence() {
1852 let mut state = state_with_past_update(Duration::from_secs(10));
1862 state.last_query_duration = Duration::from_secs(1);
1863
1864 let time_window_size = Some(Duration::from_secs(60));
1865 let min_refresh = Duration::from_secs(5);
1866 let flow_id = 1;
1867
1868 let result = state.get_next_start_query_time(
1869 flow_id,
1870 &time_window_size,
1871 min_refresh,
1872 None,
1873 20,
1874 false, );
1876
1877 let expected = state.last_update_time + Duration::from_secs(60);
1879 assert_eq!(result, expected);
1880 }
1881}