Skip to main content

flow/batching_mode/
state.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Batching mode task state, which changes frequently
16//!
17
18use 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/// The state of the [`BatchingTask`].
39#[derive(Debug)]
40pub struct TaskState {
41    /// Query context
42    pub(crate) query_ctx: QueryContextRef,
43    /// last query complete time
44    last_update_time: Instant,
45    /// last time query duration
46    last_query_duration: Duration,
47    /// Last successful execution time in unix timestamp milliseconds.
48    last_exec_time_millis: Option<i64>,
49    /// First execution time in unix timestamp milliseconds, set once.
50    start_time_millis: Option<i64>,
51    /// Dirty Time windows need to be updated
52    /// mapping of `start -> end` and non-overlapping
53    pub(crate) dirty_time_windows: DirtyTimeWindows,
54    checkpoint_mode: CheckpointMode,
55    pending_fenced_repair: Option<FencedRepair>,
56    /// Region id -> last consumed watermark sequence. Incremental scans use
57    /// this as the next lower sequence bound for each source region.
58    checkpoints: BTreeMap<u64, u64>,
59    /// Once set, the task will never attempt incremental mode again.
60    /// Set when the flow's query shape is deterministically incompatible
61    /// with incremental execution (e.g. unsupported aggregate expressions).
62    incremental_disabled: bool,
63    exec_state: ExecState,
64    /// Shutdown receiver
65    pub(crate) shutdown_rx: oneshot::Receiver<()>,
66    /// Task handle
67    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    /// Record the first-execution start time. Call this once, just before
97    /// the first frontend query is dispatched, not after it completes.
98    pub fn record_start_time_if_first(&mut self) {
99        if self.start_time_millis.is_none() {
100            // start_time is recorded just before the first frontend query is dispatched
101            // (pre-execution), so it may be marginally earlier than the streaming engine's
102            // start_time which is set post-execution. Both are valid approximations of
103            // "when this flow first ran".
104            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    /// First execution time in unix timestamp milliseconds, set once.
122    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    /// Read the live dirty-window queue without exposing TaskState internals to execution owners.
135    pub fn dirty_time_windows(&self) -> &DirtyTimeWindows {
136        &self.dirty_time_windows
137    }
138
139    /// Returns the in-progress fenced repair, if the task is repairing dirty
140    /// windows under a frozen full-snapshot high watermark.
141    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    /// Permanently disable incremental mode for this task and
150    /// immediately fall back to full snapshot for the current cycle.
151    pub fn disable_incremental(&mut self) {
152        self.incremental_disabled = true;
153        self.mark_full_snapshot();
154    }
155
156    /// Move back to top-level FullSnapshot mode. If a fenced repair is active,
157    /// restore its not-yet-in-flight pending windows to the live dirty queue so
158    /// the moved backlog is not lost.
159    pub fn mark_full_snapshot(&mut self) {
160        self.abandon_fenced_repair();
161    }
162
163    /// Replace full-snapshot checkpoints with a complete watermark proof.
164    /// Clears fenced repair state and enters Incremental unless disabled.
165    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    /// Advance only the participating regions for an incremental delta query.
174    /// This also clears any stale fenced repair sub-state.
175    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    /// Start repairing the current live dirty windows under a frozen high `H`.
192    /// The current live backlog is moved into the fenced repair so successful
193    /// chunks are consumed from that backlog. New post-`H` dirty signals can
194    /// still arrive in the live queue while the fenced repair is active.
195    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    /// Start repairing explicit bounded windows under a frozen high `H`.
212    ///
213    /// The caller supplies the complete repair scope. Live dirty windows are
214    /// left unchanged so signals received after `H` remain separate.
215    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    /// Finish the fenced repair and promote the frozen high watermark to the
231    /// checkpoint map. Incremental-disabled flows stay in FullSnapshot mode.
232    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    /// Abandon the current fenced repair and restore all not-yet-in-flight
242    /// pending windows to the live dirty queue for a fresh scoped repair.
243    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    /// Restore a scoped query's windows after a failed or unproven run. During
255    /// an active fenced repair this requeues into `pending_windows`; otherwise
256    /// it restores to the live dirty queue.
257    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    /// Generate the next scoped filter from the fenced-repair queue when active;
270    /// otherwise consume windows from the live dirty queue.
271    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            // Fenced windows are an explicit frozen repair scope. They must not
282            // be pruned by the moving live-data expiration boundary, and an
283            // empty repair must remain active until its high watermark is
284            // explicitly finished.
285            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    /// Returns true only when the query result's participating regions and
306    /// terminal watermarks exactly match the fenced repair's frozen high `H`.
307    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    /// Whether the active fenced repair has drained all pending windows.
329    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    /// Full-snapshot checkpoint advances require a watermark for every region
336    /// that participated in the query.
337    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    /// Incremental advances are limited to participating regions whose returned
350    /// watermark is not older than the stored checkpoint.
351    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    /// Compute the next query delay based on the time window size or the last query duration.
373    /// Aiming to avoid too frequent queries. But also not too long delay.
374    ///
375    /// next wait time is calculated as:
376    /// last query duration, capped by [max(min_run_interval, time_window_size), max_timeout],
377    /// note at most wait for `max_timeout`.
378    ///
379    /// if current the dirty time range is longer than one query can handle,
380    /// execute immediately to faster clean up dirty time windows.
381    /// Active fenced repairs also execute immediately while pending windows
382    /// remain: the current backlog has moved out of live dirty windows and into
383    /// `pending_fenced_repair.pending_windows`.
384    ///
385    /// If `prefer_short_incremental_cadence` is true, run incremental queries
386    /// more often when there is no large dirty backlog. This only reduces the
387    /// chance of hitting a stale cursor after flush; it is not required for
388    /// correctness.
389    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        // = last query duration, capped by [max(min_run_interval, time_window_size), max_timeout], note at most `max_timeout`
399        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        // compute how much time range can be handled in one query
421        let max_query_update_range = (*time_window_size)
422            .unwrap_or_default()
423            .mul_f64(max_filter_num_per_query as f64);
424        // if dirty time range is more than one query can handle, execute immediately
425        // to faster clean up dirty time windows
426        if cur_dirty_window_size < max_query_update_range {
427            if prefer_short_incremental_cadence {
428                // Run incremental queries sooner than the normal time-window
429                // cadence, while still backing off by at least the previous
430                // query duration and respecting the max-timeout cap.
431                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            // if dirty time windows can't be clean up in one query, execute immediately to faster
443            // clean up dirty time windows
444            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/// For keep recording of dirty time windows, which is time window that have new data inserted
456/// since last query.
457#[derive(Debug, Clone)]
458pub struct DirtyTimeWindows {
459    /// windows's `start -> end` and non-overlapping
460    /// `end` is exclusive(and optional)
461    windows: BTreeMap<Timestamp, Option<Timestamp>>,
462    /// Maximum number of filters allowed in a single query
463    max_filter_num_per_query: usize,
464    /// Time window merge distance
465    ///
466    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    /// Time window merge distance
501    ///
502    /// TODO(discord9): make those configurable
503    pub const MERGE_DIST: i32 = 3;
504
505    /// Add lower bounds to the dirty time windows. Upper bounds are ignored.
506    ///
507    /// # Arguments
508    ///
509    /// * `lower_bounds` - An iterator of lower bounds to be added.
510    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    /// Add all dirty markers from another dirty-window set.
540    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            // `None` is a dirty marker without a known upper bound.  When one
562            // side has a concrete end, keep it so merging a restored snapshot
563            // never shrinks an already-known dirty range with the same start.
564            (Some(end), None) | (None, Some(end)) => Some(end),
565            (None, None) => None,
566        }
567    }
568
569    /// Detach all dirty windows while retaining this queue's configured limits.
570    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    /// Clean all dirty time windows, useful when can't found time window expr
579    pub fn clean(&mut self) {
580        self.windows.clear();
581    }
582
583    /// Set windows to be dirty, only useful for full aggr without time window
584    /// to mark some new data is inserted
585    pub fn set_dirty(&mut self) {
586        self.add_or_merge_window(Timestamp::new_second(0), None);
587    }
588
589    /// Number of dirty windows.
590    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    /// Get the effective count of time windows, which is the number of time windows that can be
599    /// used for query, compute from total time window range divided by `window_size`.
600    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        // not sure window_size is zero have any meaning, but just in case
618        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    /// Generate all filter expressions consuming all time windows
626    ///
627    /// there is two limits:
628    /// - shouldn't return a too long time range(<=`window_size * window_cnt`), so that the query can be executed in a reasonable time
629    /// - shouldn't return too many time range exprs, so that the query can be parsed properly instead of causing parser to overflow
630    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        // get the first `window_cnt` time windows
682        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 time range is too long, stop
694            if cur_time_range >= max_time_range {
695                break;
696            }
697
698            // if we have enough time windows, stop
699            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                // too large a window, split it
712                // split at window_size * times
713                let surplus = max_time_range - cur_time_range;
714                if surplus.num_seconds() <= window_size.num_seconds() {
715                    // Skip splitting if surplus is smaller than window_size
716                    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                // remove the original window
727                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            // align using time window exprs
780            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    /// Align a time range `[start, end)` (end is optional and exclusive) to
819    /// time window boundaries defined by the time window expr.
820    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                    // if after aligned, end is the same, then use end(because it's already aligned) else use aligned end
836                    .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    /// Merge time windows that overlaps or get too close
844    ///
845    /// TODO(discord9): not merge and prefer to send smaller time windows? how?
846    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        // previous time window
865        let mut prev_tw = None;
866        for (mut lower_bound, upper_bound) in std::mem::take(&mut self.windows) {
867            // filter out expired time window
868            if let Some(expire_lower_bound) = expire_lower_bound {
869                match upper_bound {
870                    // A bounded range ending at or before the expire bound is
871                    // fully expired, drop it.
872                    Some(upper_bound) if upper_bound <= expire_lower_bound => continue,
873                    // A bounded range crossing the expire bound keeps its
874                    // still-live suffix. The expire bound is aligned to the
875                    // time window boundary by the caller, so the clipped start
876                    // stays aligned.
877                    Some(_) if lower_bound < expire_lower_bound => {
878                        lower_bound = expire_lower_bound;
879                    }
880                    // Unbounded windows keep the start-based behavior.
881                    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            // if cur.lower - prev.upper <= window_size * MERGE_DIST, merge
892            // this also deal with overlap windows because cur.lower > prev.lower is always true
893            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                // Union the two windows: the current window may be contained
910                // in the previous one, so keep the larger upper bound.
911                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/// Dirty windows that must be repaired under a frozen full-snapshot watermark.
951/// This is a FullSnapshot sub-state, not a separate checkpoint mode.
952#[derive(Debug, Clone)]
953pub struct FencedRepair {
954    high: BTreeMap<u64, u64>,
955    pending_windows: DirtyTimeWindows,
956}
957
958impl FencedRepair {
959    /// Frozen high watermark `H` used as the snapshot upper bound for chunks.
960    pub fn high(&self) -> &BTreeMap<u64, u64> {
961        &self.high
962    }
963
964    /// Dirty windows still waiting to be repaired under `high`.
965    pub fn pending_windows(&self) -> &DirtyTimeWindows {
966        &self.pending_windows
967    }
968}
969
970/// Filter Expression's information
971#[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            // just enough to merge
1040            (
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            // separate time window
1055            (
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            // overlapping
1076            (
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            // complex overlapping
1091            (
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            // split range
1107            (
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            // split 2 min into 1 min
1128            (
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            // split 3s + 1min into 3s + 57s
1140            (
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            // expired
1155            (
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        // let len = testcases.len();
1169        // let testcases = testcases[(len - 2)..(len - 1)].to_vec();
1170        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            // A contained bounded range must not shrink the containing window:
1204            // [0s, 15s) merged with nested [5s, 10s) stays [0s, 15s).
1205            (
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            // An unbounded dirty window nested in a bounded range must not
1213            // shrink the range either: [0s, 15s) merged with 3s (window end
1214            // 8s) stays [0s, 15s).
1215            (
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            // Disjoint bounded ranges far apart are kept separate.
1223            (
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            // Overlapping bounded ranges are unioned: [0s, 10s) and [5s, 20s)
1234            // become [0s, 20s).
1235            (
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        // Expire bound handling for bounded ranges vs unbounded windows.
1254        let expire_testcases = vec![
1255            // A bounded range ending at the expire bound is fully expired.
1256            (
1257                vec![(Timestamp::new_second(0), Some(Timestamp::new_second(10)))],
1258                BTreeMap::from([]),
1259            ),
1260            // A bounded range ending before the expire bound is fully expired.
1261            (
1262                vec![(Timestamp::new_second(0), Some(Timestamp::new_second(5)))],
1263                BTreeMap::from([]),
1264            ),
1265            // A bounded range crossing the expire bound keeps its live
1266            // suffix: [0s, 15s) with expire 10s becomes [10s, 15s).
1267            (
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            // A bounded range starting at the expire bound is kept intact.
1272            (
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            // An unbounded window starting before the expire bound is
1277            // dropped, preserving the existing start-based behavior.
1278            (vec![(Timestamp::new_second(5), None)], BTreeMap::from([])),
1279            // An unbounded window starting at the expire bound is kept.
1280            (
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        // After disable, mode becomes FullSnapshot and flag is set.
1522        state.disable_incremental();
1523        assert!(state.is_incremental_disabled());
1524        assert_eq!(state.checkpoint_mode(), CheckpointMode::FullSnapshot);
1525
1526        // `advance_checkpoints` will NOT transition to Incremental when disabled.
1527        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        // `mark_full_snapshot` does not re-enable incremental.
1535        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    /// Helper: create a `TaskState` whose `last_update_time` is a known duration in the past.
1677    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        // When prefer_short_incremental_cadence is true and dirty backlog is manageable,
1688        // the next start time should be last_update_time + min_refresh (short cadence),
1689        // ignoring the longer time_window_size.
1690        let state = state_with_past_update(Duration::from_secs(10));
1691
1692        let time_window_size = Some(Duration::from_secs(60)); // large window
1693        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, // prefer_short_incremental_cadence
1703        );
1704
1705        // With short cadence, result should be last_update_time + min_refresh.
1706        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        // When prefer_short_incremental_cadence is false (full snapshot mode),
1756        // the normal long-cadence based on time_window_size applies.
1757        let mut state = state_with_past_update(Duration::from_secs(10));
1758        // Make last_query_duration small so the lower bound (time_window_size) dominates.
1759        state.last_query_duration = Duration::from_secs(1);
1760
1761        let time_window_size = Some(Duration::from_secs(60)); // large window
1762        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, // prefer_short_incremental_cadence = false
1772        );
1773
1774        // With normal cadence, result should be last_update_time + time_window_size
1775        // (since last_query_duration < time_window_size).
1776        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        // Dirty-window overflow must always schedule immediately,
1783        // regardless of prefer_short_incremental_cadence.
1784        let mut state = state_with_past_update(Duration::from_secs(10));
1785        // Create a very large dirty backlog.
1786        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)); // tiny window => overflow
1791        let min_refresh = Duration::from_secs(5);
1792        let flow_id = 1;
1793
1794        // With short cadence flag.
1795        let result = state.get_next_start_query_time(
1796            flow_id,
1797            &time_window_size,
1798            min_refresh,
1799            None,
1800            1, // max 1 filter => tiny capacity
1801            true,
1802        );
1803        assert!(
1804            result <= Instant::now(),
1805            "dirty overflow should schedule immediately"
1806        );
1807
1808        // Without short cadence flag — same behavior.
1809        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        // When prefer_short_incremental_cadence is true but the dirty backlog is
1853        // manageable, the short cadence is applied. This test verifies that the
1854        // caller-side guard (checkpoint_mode + !is_incremental_disabled) controls
1855        // whether short cadence is requested at all — when incremental is disabled,
1856        // the flag is false, and the long cadence applies.
1857        //
1858        // This simulates the case where the caller computed
1859        // prefer_short_incremental_cadence = false (e.g. incremental disabled
1860        // or FullSnapshot mode), so the long cadence is used.
1861        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, // prefer_short_incremental_cadence = false
1875        );
1876
1877        // With normal cadence, result should be last_update_time + time_window_size.
1878        let expected = state.last_update_time + Duration::from_secs(60);
1879        assert_eq!(result, expected);
1880    }
1881}