Skip to main content

flow/batching_mode/
task.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
15use std::collections::{BTreeMap, BTreeSet, HashMap, HashSet};
16use std::sync::{Arc, RwLock};
17use std::time::{Duration, SystemTime, UNIX_EPOCH};
18
19use api::v1::{CreateTableExpr, TableName};
20use catalog::CatalogManagerRef;
21use catalog::kvbackend::KvBackendCatalogManager;
22use client::OutputWithMetrics;
23use common_error::ext::BoxedError;
24use common_meta::key::schema_name::SchemaNameKey;
25use common_query::OutputData;
26use common_query::logical_plan::breakup_insert_plan;
27use common_telemetry::tracing::warn;
28use common_telemetry::{debug, info};
29use common_time::{TimeToLive, Timestamp};
30use datafusion::datasource::DefaultTableSource;
31use datafusion::sql::unparser::expr_to_sql;
32use datafusion_common::tree_node::{Transformed, TreeNode};
33use datafusion_common::utils::quote_identifier;
34use datafusion_common::{DFSchemaRef, ScalarValue, TableReference};
35use datafusion_expr::{DmlStatement, LogicalPlan, Projection, WriteOp, col, lit};
36use datatypes::schema::Schema;
37use datatypes::vectors::Helper;
38use futures::TryStreamExt;
39use query::QueryEngineRef;
40use query::options::{
41    FLOW_INCREMENTAL_AFTER_SEQS, FLOW_INCREMENTAL_MODE, FLOW_INCREMENTAL_MODE_SEQUENCE_RANGE,
42    FLOW_RETURN_REGION_SEQ,
43};
44use query::query_engine::DefaultSerializer;
45use session::context::QueryContextRef;
46use snafu::{OptionExt, ResultExt};
47use sql::parsers::utils::is_tql;
48use store_api::mito_engine_options::MERGE_MODE_KEY;
49use substrait::{DFLogicalSubstraitConvertor, SubstraitPlan};
50use table::TableRef;
51use table::table::adapter::DfTableProviderAdapter;
52use tokio::sync::oneshot::error::TryRecvError;
53use tokio::sync::{Mutex, OwnedMutexGuard, oneshot};
54use tokio::time::Instant;
55
56use crate::batching_mode::BatchingModeOptions;
57use crate::batching_mode::checkpoint::checkpoint_mode_label;
58use crate::batching_mode::eval_schedule::{EvalSchedule, select_due_scheduled_times};
59use crate::batching_mode::frontend_client::{FrontendClient, PeerDesc};
60use crate::batching_mode::state::{
61    CheckpointMode, DirtyTimeWindows, FilterExprInfo, TaskState, to_df_literal,
62};
63use crate::batching_mode::table_creator::{QueryType, create_table_with_expr};
64use crate::batching_mode::time_window::TimeWindowExpr;
65use crate::batching_mode::utils::{
66    AddFilterRewriter, ColumnMatcherRewriter, analyze_incremental_aggregate_plan, df_plan_to_sql,
67    gen_plan_with_matching_schema_and_values, get_table_info_df_schema, sql_to_df_plan,
68};
69use crate::df_optimizer::apply_df_optimizer;
70use crate::error::{
71    DatafusionSnafu, ExternalSnafu, InvalidQuerySnafu, SubstraitEncodeLogicalPlanSnafu,
72    UnexpectedSnafu,
73};
74use crate::metrics::{
75    METRIC_FLOW_BATCHING_ENGINE_ERROR_CNT, METRIC_FLOW_BATCHING_ENGINE_QUERY_TIME,
76    METRIC_FLOW_BATCHING_ENGINE_SLOW_QUERY, METRIC_FLOW_BATCHING_ENGINE_START_QUERY_CNT,
77    METRIC_FLOW_ROWS,
78};
79use crate::{Error, FlowId};
80
81mod ckpt;
82mod inc;
83
84/// Returns the current wall-clock Unix timestamp in seconds.
85fn wall_clock_unix_secs() -> i64 {
86    SystemTime::now()
87        .duration_since(UNIX_EPOCH)
88        .unwrap_or_default()
89        .as_secs() as i64
90}
91
92/// Initial scheduler cursor for `start_scheduled_loop`: exactly one interval
93/// before `start_secs` so the first due scheduled time is `start_secs` itself.
94///
95/// Fallible: a `start_secs - interval_secs` difference that does not fit in
96/// `i64` is an explicit error instead of a saturated cursor that would make
97/// the first due scheduled time `start_secs + interval_secs` and silently skip
98/// the `start_secs` tick.
99fn initial_schedule_cursor(start_secs: i64, interval_secs: i64) -> Result<i64, Error> {
100    let cursor = i128::from(start_secs) - i128::from(interval_secs);
101    i64::try_from(cursor).map_err(|_| {
102        UnexpectedSnafu {
103            reason: format!(
104                "Cannot compute the initial eval schedule cursor one interval before start {start_secs} (interval={interval_secs}): {cursor} does not fit in i64"
105            ),
106        }
107        .build()
108    })
109}
110
111/// Whole seconds to sleep until the next scheduled time `next`, measured from
112/// the current wall clock `wall_now_secs`.
113///
114/// Fallible: `next` must be strictly after `wall_now_secs` and the difference
115/// must fit in `u64`. In practice `i64::MAX - i64::MIN` is exactly `u64::MAX`,
116/// so the difference always fits once `next > wall_now_secs`; the explicit
117/// error keeps the scheduled loop panic-free and wrap-free regardless.
118fn sleep_delta_secs(next: i64, wall_now_secs: i64) -> Result<u64, Error> {
119    let delta = i128::from(next) - i128::from(wall_now_secs);
120    u64::try_from(delta).map_err(|_| {
121        UnexpectedSnafu {
122            reason: format!(
123                "Cannot sleep until the next scheduled time {next}: the delta from wall clock {wall_now_secs} is {delta} seconds, which does not fit in u64"
124            ),
125        }
126        .build()
127    })
128}
129
130/// Scheduled time in seconds converted to milliseconds for the
131/// `FLOW_SCHEDULED_TIME_MILLIS` extension.
132///
133/// Fallible: a seconds value whose millisecond product does not fit in `i64`
134/// is an explicit error instead of a saturated `i64::MAX` that would silently
135/// misrepresent the logical scheduled time.
136fn scheduled_time_millis(scheduled_time_secs: i64) -> Result<i64, Error> {
137    scheduled_time_secs.checked_mul(1000).ok_or_else(|| {
138        UnexpectedSnafu {
139            reason: format!(
140                "Cannot convert scheduled time {scheduled_time_secs}s to milliseconds: the product exceeds i64"
141            ),
142        }
143        .build()
144    })
145}
146
147/// The task's config, immutable once created
148#[derive(Clone)]
149pub struct TaskConfig {
150    pub flow_id: FlowId,
151    pub query: String,
152    /// output schema of the query
153    pub output_schema: DFSchemaRef,
154    pub time_window_expr: Option<TimeWindowExpr>,
155    /// in seconds
156    pub expire_after: Option<i64>,
157    pub sink_table_name: [String; 3],
158    pub source_table_names: HashSet<[String; 3]>,
159    pub catalog_manager: CatalogManagerRef,
160    pub query_type: QueryType,
161    pub batch_opts: Arc<BatchingModeOptions>,
162    pub exact_sequence_range_required: bool,
163    pub flow_eval_interval: Option<Duration>,
164    /// Typed schedule configuration, pre-parsed at task creation time.
165    pub eval_schedule: Option<EvalSchedule>,
166}
167
168fn determine_query_type(query: &str, query_ctx: &QueryContextRef) -> Result<QueryType, Error> {
169    let is_tql = is_tql(query_ctx.sql_dialect(), query)
170        .map_err(BoxedError::new)
171        .context(ExternalSnafu)?;
172    Ok(if is_tql {
173        QueryType::Tql
174    } else {
175        QueryType::Sql
176    })
177}
178
179fn is_merge_mode_last_non_null(options: &HashMap<String, String>) -> bool {
180    options
181        .get(MERGE_MODE_KEY)
182        .map(|mode| mode.eq_ignore_ascii_case("last_non_null"))
183        .unwrap_or(false)
184}
185
186fn encode_insert_plan_request(
187    insert_to: TableName,
188    insert_input_plan: &LogicalPlan,
189) -> Result<api::v1::QueryRequest, Error> {
190    let message = DFLogicalSubstraitConvertor {}
191        .encode(insert_input_plan, DefaultSerializer)
192        .context(SubstraitEncodeLogicalPlanSnafu)?;
193    Ok(api::v1::QueryRequest {
194        query: Some(api::v1::query_request::Query::InsertIntoPlan(
195            api::v1::InsertIntoPlan {
196                table_name: Some(insert_to),
197                logical_plan: message.to_vec(),
198            },
199        )),
200    })
201}
202
203fn recovery_aggregate_input(plan: &LogicalPlan) -> Result<LogicalPlan, Error> {
204    let plan = match plan {
205        LogicalPlan::Projection(projection) => projection.input.as_ref(),
206        _ => plan,
207    };
208    let LogicalPlan::Aggregate(aggregate) = plan else {
209        return UnexpectedSnafu {
210            reason: "Recovery timestamp projection did not find an aggregate".to_string(),
211        }
212        .fail();
213    };
214    Ok(aggregate.input.as_ref().clone())
215}
216
217fn capture_recovery_batch_windows(
218    batch: &common_recordbatch::RecordBatch,
219    time_window_expr: &TimeWindowExpr,
220    windows: &mut BTreeSet<(Timestamp, Timestamp)>,
221) -> Result<(), Error> {
222    if batch.num_columns() != 1 {
223        return UnexpectedSnafu {
224            reason: format!(
225                "Recovery timestamp projection returned {} columns instead of one",
226                batch.num_columns()
227            ),
228        }
229        .fail();
230    }
231    let values = Helper::try_into_vector(batch.column(0).clone())
232        .map_err(BoxedError::new)
233        .context(ExternalSnafu)?;
234    for index in 0..values.len() {
235        let timestamp = values.get(index).as_timestamp().context(UnexpectedSnafu {
236            reason: "Recovery timestamp projection returned a null or non-timestamp value"
237                .to_string(),
238        })?;
239        let (start, end) = time_window_expr.eval(timestamp)?;
240        let window = (
241            start.context(UnexpectedSnafu {
242                reason: "Recovery time-window expression returned no lower bound".to_string(),
243            })?,
244            end.context(UnexpectedSnafu {
245                reason: "Recovery time-window expression returned no upper bound".to_string(),
246            })?,
247        );
248        windows.insert(window);
249    }
250    Ok(())
251}
252
253fn format_insert_target_columns(plan: &LogicalPlan) -> String {
254    plan.schema()
255        .fields()
256        .iter()
257        .map(|field| quote_identifier(field.name()).to_string())
258        .collect::<Vec<_>>()
259        .join(", ")
260}
261
262/// Owns a whole serialized execution round. It may be moved to an execution
263/// collaborator, so cancellation of the caller cannot release the round early.
264pub struct BatchingExecutionGuard {
265    _lock: OwnedMutexGuard<()>,
266    restore: Option<(Arc<RwLock<TaskState>>, QueryContextRef)>,
267}
268
269impl BatchingExecutionGuard {
270    fn new(lock: OwnedMutexGuard<()>) -> Self {
271        Self {
272            _lock: lock,
273            restore: None,
274        }
275    }
276
277    fn restore_query_context(&mut self, state: Arc<RwLock<TaskState>>, old_ctx: QueryContextRef) {
278        self.restore = Some((state, old_ctx));
279    }
280}
281
282impl Drop for BatchingExecutionGuard {
283    fn drop(&mut self) {
284        if let Some((state, old_ctx)) = self.restore.take()
285            && let Ok(mut state) = state.write()
286        {
287            state.query_ctx = old_ctx;
288        }
289    }
290}
291
292#[derive(Clone)]
293pub struct BatchingTask {
294    pub config: Arc<TaskConfig>,
295    pub state: Arc<RwLock<TaskState>>,
296    /// Serializes plan generation, execution, checkpoint advancement, and dirty
297    /// window restoration for this flow. Without this, a manual flush and the
298    /// background loop can process the same checkpoint range concurrently.
299    execution_lock: Arc<Mutex<()>>,
300    execution: Option<Arc<dyn crate::BatchingExecution>>,
301}
302
303/// Arguments for creating batching task
304pub struct TaskArgs<'a> {
305    pub flow_id: FlowId,
306    pub query: &'a str,
307    pub plan: LogicalPlan,
308    pub time_window_expr: Option<TimeWindowExpr>,
309    pub expire_after: Option<i64>,
310    pub sink_table_name: [String; 3],
311    pub source_table_names: Vec<[String; 3]>,
312    pub query_ctx: QueryContextRef,
313    pub catalog_manager: CatalogManagerRef,
314    pub shutdown_rx: oneshot::Receiver<()>,
315    pub batch_opts: Arc<BatchingModeOptions>,
316    pub flow_eval_interval: Option<Duration>,
317    /// Typed schedule configuration pre-parsed from `CreateFlowArgs`.
318    pub eval_schedule: Option<EvalSchedule>,
319}
320
321pub struct PlanInfo {
322    pub plan: LogicalPlan,
323    pub dirty_restore: DirtyRestore,
324    pub coverage: QueryCoverage,
325}
326
327#[derive(Clone)]
328pub enum QueryCoverage {
329    /// Explicit full-query snapshot coverage, e.g. TQL or evaluation-interval
330    /// SQL flows whose plan shape cannot be safely dirty-window pruned. This
331    /// must not be used as an implicit recovery path for scoped repair or an
332    /// unsafe incremental rewrite fallback.
333    UnfilteredFull,
334    /// Scoped full-snapshot repair over the current dirty windows. A successful
335    /// result may start a fenced repair if new dirty windows appeared meanwhile.
336    ScopedBaseRepair,
337    /// A chunk of windows being repaired under the frozen high-watermark `H`.
338    /// The `high` map is sent as snapshot read bounds and must be matched by
339    /// the returned terminal watermarks before checkpoints can advance.
340    FencedRepairChunk { high: BTreeMap<u64, u64> },
341    /// Incremental delta query over `(checkpoint, scan-open snapshot]`.
342    IncrementalDelta,
343}
344
345impl QueryCoverage {
346    /// Whether this query should use incremental scan extensions and
347    /// incremental checkpoint advancement rules.
348    fn is_incremental_delta(&self) -> bool {
349        matches!(self, Self::IncrementalDelta)
350    }
351
352    /// Snapshot upper bounds requested from the storage layer. Only fenced
353    /// repair chunks carry bounds; all other coverage relies on normal scans.
354    fn snapshot_seqs(&self) -> HashMap<u64, u64> {
355        match self {
356            Self::FencedRepairChunk { high } => high.iter().map(|(k, v)| (*k, *v)).collect(),
357            _ => HashMap::new(),
358        }
359    }
360}
361
362pub enum DirtyRestore {
363    /// The query was scoped to dirty time ranges; restore those ranges if the
364    /// run fails.
365    Scoped(FilterExprInfo),
366    /// The query could not be scoped to dirty time ranges, so the dirty-window
367    /// state is only a dirty signal. Restore the consumed signal if the full
368    /// run fails.
369    ///
370    /// TODO(discord9): Full-query runs only need a dirty bool flag. Refactor
371    /// the unscoped path to stop reusing `DirtyTimeWindows` for this signal.
372    Unscoped(DirtyTimeWindows),
373}
374
375pub struct ExecuteOnceOutcome {
376    pub new_query: Option<PlanInfo>,
377    /// Execution result of the generated insert plan.
378    ///
379    /// `Ok(Some((affected_rows, elapsed)))` means a query was executed.
380    /// `Ok(None)` means no query was generated because there was no dirty signal.
381    /// `Err(_)` means plan generation or execution failed.
382    pub result: Result<Option<(usize, Duration)>, Error>,
383}
384
385impl BatchingTask {
386    #[allow(clippy::too_many_arguments)]
387    pub fn try_new(args: TaskArgs<'_>) -> Result<Self, Error> {
388        Self::try_new_with_exact_sequence_range_required(args, false)
389    }
390
391    pub fn try_new_with_exact_sequence_range_required(
392        TaskArgs {
393            flow_id,
394            query,
395            plan,
396            time_window_expr,
397            expire_after,
398            sink_table_name,
399            source_table_names,
400            query_ctx,
401            catalog_manager,
402            shutdown_rx,
403            batch_opts,
404            flow_eval_interval,
405            eval_schedule,
406        }: TaskArgs<'_>,
407        exact_sequence_range_required: bool,
408    ) -> Result<Self, Error> {
409        let mut state = TaskState::with_dirty_time_windows(
410            query_ctx.clone(),
411            shutdown_rx,
412            DirtyTimeWindows::new(
413                batch_opts.experimental_max_filter_num_per_query,
414                batch_opts.experimental_time_window_merge_threshold,
415            ),
416        );
417        if !batch_opts.experimental_enable_incremental_read {
418            state.disable_incremental();
419        }
420
421        Ok(Self {
422            config: Arc::new(TaskConfig {
423                flow_id,
424                query: query.to_string(),
425                time_window_expr,
426                expire_after,
427                sink_table_name,
428                source_table_names: source_table_names.into_iter().collect(),
429                catalog_manager,
430                output_schema: plan.schema().clone(),
431                query_type: determine_query_type(query, &query_ctx)?,
432                exact_sequence_range_required,
433                batch_opts,
434                flow_eval_interval,
435                eval_schedule,
436            }),
437            state: Arc::new(RwLock::new(state)),
438            execution_lock: Arc::new(Mutex::new(())),
439            execution: None,
440        })
441    }
442
443    pub(crate) fn with_execution(
444        mut self,
445        execution: Option<Arc<dyn crate::BatchingExecution>>,
446    ) -> Self {
447        self.execution = execution;
448        self
449    }
450
451    pub(crate) fn stop_execution(&self) {
452        if let Some(execution) = &self.execution {
453            execution.stop();
454        }
455    }
456
457    pub fn last_execution_time_millis(&self) -> Option<i64> {
458        self.state.read().unwrap().last_execution_time_millis()
459    }
460
461    pub fn start_time_millis(&self) -> Option<i64> {
462        self.state.read().unwrap().start_time_millis()
463    }
464
465    /// Collect flow-related extensions from the task's query context that should be
466    /// forwarded to the frontend (e.g. scheduled time).
467    fn frontend_extensions(&self) -> HashMap<String, String> {
468        let ctx = self.state.read().unwrap();
469        let all = ctx.query_ctx.extensions();
470        let mut flow_exts = HashMap::new();
471        // Propagate the scheduled time extension if present so that frontend
472        // execution can use the same logical time.
473        if let Some(v) = all.get(query::options::FLOW_SCHEDULED_TIME_MILLIS) {
474            flow_exts.insert(
475                query::options::FLOW_SCHEDULED_TIME_MILLIS.to_string(),
476                v.clone(),
477            );
478        }
479        flow_exts
480    }
481
482    /// mark time window range (now - expire_after, now) as dirty (or (0, now) if expire_after not set)
483    ///
484    /// useful for flush_flow to flush dirty time windows range
485    pub fn mark_all_windows_as_dirty(&self) -> Result<(), Error> {
486        let now = SystemTime::now();
487        let now = Timestamp::new_second(
488            now.duration_since(UNIX_EPOCH)
489                .expect("Time went backwards")
490                .as_secs() as _,
491        );
492        let lower_bound = self
493            .config
494            .expire_after
495            .map(|e| now.sub_duration(Duration::from_secs(e as _)))
496            .transpose()
497            .map_err(BoxedError::new)
498            .context(ExternalSnafu)?
499            .unwrap_or(Timestamp::new_second(0));
500        debug!(
501            "Flow {} mark range ({:?}, {:?}) as dirty",
502            self.config.flow_id, lower_bound, now
503        );
504        self.state
505            .write()
506            .unwrap()
507            .dirty_time_windows
508            .add_window(lower_bound, Some(now));
509        Ok(())
510    }
511
512    /// Create the sink table if needed, then return it.
513    pub async fn check_or_create_sink_table(
514        &self,
515        engine: &QueryEngineRef,
516        frontend_client: &Arc<FrontendClient>,
517    ) -> Result<TableRef, Error> {
518        if !self.is_table_exist(&self.config.sink_table_name).await? {
519            let create_table = self.gen_create_table_expr(engine.clone()).await?;
520            info!(
521                "Try creating sink table(if not exists) with expr: {:?}",
522                create_table
523            );
524            self.create_table(frontend_client, create_table).await?;
525            info!(
526                "Sink table {}(if not exists) created",
527                self.config.sink_table_name.join(".")
528            );
529        }
530        let (table, _) = get_table_info_df_schema(
531            self.config.catalog_manager.clone(),
532            self.config.sink_table_name.clone(),
533        )
534        .await?;
535        Ok(table)
536    }
537
538    /// Returns a forked snapshot of the task query context for extension-owned planning.
539    pub fn query_context_snapshot(&self) -> QueryContextRef {
540        let query_ctx = self.state.read().unwrap().query_ctx.clone();
541        Arc::new(query_ctx.fork())
542    }
543
544    /// Returns the retention lower bound aligned to this task's time window.
545    pub fn recovery_retention_lower_bound(&self) -> Result<Option<Timestamp>, Error> {
546        let Some(expire_after) = self.config.expire_after else {
547            return Ok(None);
548        };
549        let expire_after = u64::try_from(expire_after).map_err(|_| {
550            UnexpectedSnafu {
551                reason: format!(
552                    "Flow {} has negative expire_after {expire_after}",
553                    self.config.flow_id
554                ),
555            }
556            .build()
557        })?;
558        let now = Timestamp::new_second(
559            SystemTime::now()
560                .duration_since(UNIX_EPOCH)
561                .map_err(|err| {
562                    UnexpectedSnafu {
563                        reason: format!("Failed to read recovery wall clock: {err}"),
564                    }
565                    .build()
566                })?
567                .as_secs() as i64,
568        );
569        let lower = now
570            .sub_duration(Duration::from_secs(expire_after))
571            .map_err(BoxedError::new)
572            .context(ExternalSnafu)?;
573        self.config
574            .time_window_expr
575            .as_ref()
576            .context(UnexpectedSnafu {
577                reason: "Recovery expiry requires a time-window expression".to_string(),
578            })?
579            .eval(lower)?
580            .0
581            .context(UnexpectedSnafu {
582                reason: "Recovery expiry time-window expression returned no lower bound"
583                    .to_string(),
584            })
585            .map(Some)
586    }
587
588    /// Verifies that recovery can still read all source data required by its retained scope.
589    pub async fn validate_recovery_retention(
590        &self,
591        retention_lower: Option<Timestamp>,
592        windows: &[(Timestamp, Timestamp)],
593    ) -> Result<(), Error> {
594        for name in &self.config.source_table_names {
595            let table = self
596                .config
597                .catalog_manager
598                .table(&name[0], &name[1], &name[2], None)
599                .await
600                .map_err(BoxedError::new)
601                .context(ExternalSnafu)?
602                .context(UnexpectedSnafu {
603                    reason: format!(
604                        "Flow {} source table {} is unavailable for recovery retention validation",
605                        self.config.flow_id,
606                        name.join(".")
607                    ),
608                })?;
609            let ttl = if let Some(ttl) = table.table_info().meta.options.ttl {
610                ttl
611            } else {
612                let manager = self
613                    .config
614                    .catalog_manager
615                    .as_any()
616                    .downcast_ref::<KvBackendCatalogManager>()
617                    .context(UnexpectedSnafu {
618                        reason: format!(
619                            "Flow {} cannot resolve inherited TTL for source table {} during recovery",
620                            self.config.flow_id,
621                            name.join(".")
622                        ),
623                    })?;
624                manager
625                    .table_metadata_manager_ref()
626                    .schema_manager()
627                    .get(SchemaNameKey::new(&name[0], &name[1]))
628                    .await
629                    .map_err(BoxedError::new)
630                    .context(ExternalSnafu)?
631                    .context(UnexpectedSnafu {
632                        reason: format!(
633                            "Flow {} schema {}.{} is unavailable for recovery retention validation",
634                            self.config.flow_id, name[0], name[1]
635                        ),
636                    })?
637                    .ttl
638                    .map(Into::into)
639                    .unwrap_or(TimeToLive::Forever)
640            };
641
642            match ttl {
643                TimeToLive::Forever => {}
644                TimeToLive::Instant => {
645                    return UnexpectedSnafu {
646                        reason: format!(
647                            "Flow {} source table {} has instant TTL and cannot be recovered",
648                            self.config.flow_id,
649                            name.join(".")
650                        ),
651                    }
652                    .fail();
653                }
654                TimeToLive::Duration(ttl) => {
655                    let expire_after = self.config.expire_after.context(UnexpectedSnafu {
656                        reason: format!(
657                            "Flow {} source table {} has TTL {ttl:?}, but recovery has no expire_after",
658                            self.config.flow_id,
659                            name.join(".")
660                        ),
661                    })?;
662                    let expire_after = u64::try_from(expire_after).map_err(|_| {
663                        UnexpectedSnafu {
664                            reason: format!(
665                                "Flow {} has negative expire_after {expire_after} during recovery retention validation",
666                                self.config.flow_id
667                            ),
668                        }
669                        .build()
670                    })?;
671                    if ttl <= Duration::from_secs(expire_after) {
672                        return UnexpectedSnafu {
673                            reason: format!(
674                                "Flow {} source table {} TTL {ttl:?} must exceed expire_after {expire_after}s for recovery",
675                                self.config.flow_id,
676                                name.join(".")
677                            ),
678                        }
679                        .fail();
680                    }
681                    let mut oldest = retention_lower.context(UnexpectedSnafu {
682                        reason: format!(
683                            "Flow {} source table {} has finite TTL {ttl:?}, but recovery has no retained lower bound",
684                            self.config.flow_id,
685                            name.join(".")
686                        ),
687                    })?;
688                    for (start, _) in windows {
689                        oldest = oldest.min(*start);
690                    }
691                    let cutoff = Timestamp::current_millis()
692                        .sub_duration(ttl)
693                        .map_err(BoxedError::new)
694                        .context(ExternalSnafu)?;
695                    if oldest <= cutoff {
696                        return UnexpectedSnafu {
697                            reason: format!(
698                                "Flow {} recovery requires source table {} data from {oldest:?}, older than TTL {ttl:?} cutoff {cutoff:?}",
699                                self.config.flow_id,
700                                name.join(".")
701                            ),
702                        }
703                        .fail();
704                    }
705                }
706            }
707        }
708        Ok(())
709    }
710
711    /// Discovers the source time windows touched by the exact sequence range `(C, H]`.
712    ///
713    /// This compatibility wrapper captures the retention bound before discovery.
714    pub async fn capture_recovery_windows(
715        &self,
716        engine: &QueryEngineRef,
717        frontend_client: &FrontendClient,
718        lower: &BTreeMap<u64, u64>,
719    ) -> Result<(BTreeMap<u64, u64>, Vec<(Timestamp, Timestamp)>), Error> {
720        let retention_lower = self.recovery_retention_lower_bound()?;
721        self.capture_recovery_windows_since(engine, frontend_client, lower, retention_lower)
722            .await
723    }
724
725    /// Discovers windows under `lower` using a caller-frozen retention bound.
726    ///
727    /// The caller owns the execution guard and freezes this bound with its recovery scope.
728    /// This read-only method leaves task state untouched. Terminal proof must cover every
729    /// region in `lower`; subset or pruned proofs are rejected.
730    pub async fn capture_recovery_windows_since(
731        &self,
732        engine: &QueryEngineRef,
733        frontend_client: &FrontendClient,
734        lower: &BTreeMap<u64, u64>,
735        retention_lower: Option<Timestamp>,
736    ) -> Result<(BTreeMap<u64, u64>, Vec<(Timestamp, Timestamp)>), Error> {
737        if lower.is_empty() {
738            return UnexpectedSnafu {
739                reason: format!(
740                    "Flow {} recovery window capture requires nonempty lower sequence bounds",
741                    self.config.flow_id
742                ),
743            }
744            .fail();
745        }
746        if !self.sequence_range_capable().await? {
747            return UnexpectedSnafu {
748                reason: format!(
749                    "Flow {} recovery window capture requires sequence_range-capable sources",
750                    self.config.flow_id
751                ),
752            }
753            .fail();
754        }
755
756        let query_ctx = self.query_context_snapshot();
757        let plan = sql_to_df_plan(query_ctx, engine.clone(), &self.config.query, false).await?;
758        let input = recovery_aggregate_input(&plan)?;
759        let Some(analysis) = analyze_incremental_aggregate_plan(&plan)? else {
760            return UnexpectedSnafu {
761                reason: format!(
762                    "Flow {} recovery window capture requires a supported aggregate input",
763                    self.config.flow_id
764                ),
765            }
766            .fail();
767        };
768        if !analysis.unsupported_exprs.is_empty() {
769            return UnexpectedSnafu {
770                reason: format!(
771                    "Flow {} recovery window capture has unsupported aggregate expressions: {:?}",
772                    self.config.flow_id, analysis.unsupported_exprs
773                ),
774            }
775            .fail();
776        }
777        let time_window_expr = self
778            .config
779            .time_window_expr
780            .as_ref()
781            .context(UnexpectedSnafu {
782                reason: format!(
783                    "Flow {} recovery window capture requires a time-window expression",
784                    self.config.flow_id
785                ),
786            })?;
787        let input = if let Some(retention_lower) = retention_lower {
788            let mut add_filter = AddFilterRewriter::new(
789                col(&time_window_expr.column_name).gt_eq(lit(to_df_literal(retention_lower)?)),
790            );
791            input
792                .rewrite(&mut add_filter)
793                .with_context(|_| DatafusionSnafu {
794                    context: "Failed to apply recovery expire_after filter".to_string(),
795                })?
796                .data
797        } else {
798            input
799        };
800        let timestamp_plan = LogicalPlan::Projection(
801            Projection::try_new(vec![col(&time_window_expr.column_name)], Arc::new(input))
802                .context(DatafusionSnafu {
803                    context: "Failed to project recovery source timestamps".to_string(),
804                })?,
805        );
806        let catalog = &self.config.sink_table_name[0];
807        let schema = &self.config.sink_table_name[1];
808        let timestamp_plan = timestamp_plan
809            .clone()
810            .transform_down_with_subqueries(|p| {
811                if let LogicalPlan::TableScan(mut table_scan) = p {
812                    let resolved = table_scan.table_name.resolve(catalog, schema);
813                    table_scan.table_name = resolved.into();
814                    Ok(Transformed::yes(LogicalPlan::TableScan(table_scan)))
815                } else {
816                    Ok(Transformed::no(p))
817                }
818            })
819            .with_context(|_| DatafusionSnafu {
820                context: format!(
821                    "Failed to fix table ref in recovery timestamp plan, plan={:?}",
822                    timestamp_plan
823                ),
824            })?
825            .data;
826        let message = DFLogicalSubstraitConvertor {}
827            .encode(&timestamp_plan, DefaultSerializer)
828            .context(SubstraitEncodeLogicalPlanSnafu)?;
829        let request = api::v1::QueryRequest {
830            query: Some(api::v1::query_request::Query::LogicalPlan(message.to_vec())),
831        };
832        let lower_json = serde_json::to_string(lower).map_err(|err| {
833            UnexpectedSnafu {
834                reason: format!("Failed to serialize recovery lower sequence bounds: {err}"),
835            }
836            .build()
837        })?;
838        let extensions = [
839            (FLOW_RETURN_REGION_SEQ, "true"),
840            (FLOW_INCREMENTAL_MODE, FLOW_INCREMENTAL_MODE_SEQUENCE_RANGE),
841            (FLOW_INCREMENTAL_AFTER_SEQS, lower_json.as_str()),
842        ];
843        let mut peer_desc = None;
844        let result = frontend_client
845            .query_with_terminal_metrics(
846                catalog,
847                schema,
848                request,
849                &extensions,
850                &HashMap::new(),
851                &mut peer_desc,
852            )
853            .await?;
854        let mut aligned_windows = BTreeSet::new();
855        let metrics = result.metrics.clone();
856        match result.output.data {
857            OutputData::AffectedRows(_) => {
858                return UnexpectedSnafu {
859                    reason: "Recovery timestamp projection unexpectedly returned affected rows"
860                        .to_string(),
861                }
862                .fail();
863            }
864            OutputData::RecordBatches(batches) => {
865                for batch in batches.iter() {
866                    capture_recovery_batch_windows(batch, time_window_expr, &mut aligned_windows)?;
867                }
868            }
869            OutputData::Stream(mut stream) => {
870                while let Some(batch) = stream
871                    .try_next()
872                    .await
873                    .map_err(BoxedError::new)
874                    .context(ExternalSnafu)?
875                {
876                    capture_recovery_batch_windows(&batch, time_window_expr, &mut aligned_windows)?;
877                }
878            }
879        }
880        if !metrics.is_ready() {
881            return UnexpectedSnafu {
882                reason: "Recovery timestamp projection ended without terminal metrics".to_string(),
883            }
884            .fail();
885        }
886        let participating = metrics.participating_regions().context(UnexpectedSnafu {
887            reason: "Recovery timestamp projection has no participating-region proof".to_string(),
888        })?;
889        let high = metrics.region_watermark_map().context(UnexpectedSnafu {
890            reason: "Recovery timestamp projection has no terminal watermark proof".to_string(),
891        })?;
892        if participating.len() != lower.len()
893            || high.len() != lower.len()
894            || participating
895                .iter()
896                .any(|region| !lower.contains_key(region))
897            || lower
898                .iter()
899                .any(|(region, low)| high.get(region).is_none_or(|watermark| watermark < low))
900            || high.keys().any(|region| !participating.contains(region))
901        {
902            return UnexpectedSnafu {
903                reason: format!(
904                    "Recovery timestamp projection returned incomplete or regressing terminal proof: lower={lower:?}, participating={participating:?}, high={high:?}"
905                ),
906            }
907            .fail();
908        }
909        Ok((
910            high.into_iter().collect(),
911            aligned_windows.into_iter().collect(),
912        ))
913    }
914
915    /// Validates that the sink table schema can accept this flow's ordinary output.
916    pub async fn validate_sink_table_schema(&self, engine: &QueryEngineRef) -> Result<(), Error> {
917        self.validate_sink_table_schema_with_values(engine, &BTreeMap::new())
918            .await
919    }
920
921    /// Validates the flow output plus extension-owned typed sink values without consuming dirty work.
922    pub async fn validate_sink_table_schema_with_values(
923        &self,
924        engine: &QueryEngineRef,
925        values: &BTreeMap<String, ScalarValue>,
926    ) -> Result<(), Error> {
927        let (table, _) = get_table_info_df_schema(
928            self.config.catalog_manager.clone(),
929            self.config.sink_table_name.clone(),
930        )
931        .await?;
932
933        let table_meta = &table.table_info().meta;
934        let merge_mode_last_non_null =
935            is_merge_mode_last_non_null(&table_meta.options.extra_options);
936        let primary_key_indices = table_meta.primary_key_indices.clone();
937
938        gen_plan_with_matching_schema_and_values(
939            &self.config.query,
940            self.query_context_snapshot(),
941            engine.clone(),
942            table_meta.schema.clone(),
943            &primary_key_indices,
944            merge_mode_last_non_null,
945            values,
946        )
947        .await
948        .map(|_| ())
949    }
950
951    async fn is_table_exist(&self, table_name: &[String; 3]) -> Result<bool, Error> {
952        self.config
953            .catalog_manager
954            .table_exists(&table_name[0], &table_name[1], &table_name[2], None)
955            .await
956            .map_err(BoxedError::new)
957            .context(ExternalSnafu)
958    }
959
960    pub(crate) async fn execute_once_serialized(
961        &self,
962        engine: &QueryEngineRef,
963        frontend_client: &Arc<FrontendClient>,
964        max_window_cnt: Option<usize>,
965    ) -> Result<Option<(usize, Duration)>, Error> {
966        let outcome = self
967            .execute_once_serialized_with_outcome(engine, frontend_client, max_window_cnt)
968            .await;
969        outcome.result
970    }
971
972    /// Executes one flow evaluation under `execution_lock` and keeps the
973    /// generated query context for the background loop's error logging/backoff.
974    async fn execute_once_serialized_with_outcome(
975        &self,
976        engine: &QueryEngineRef,
977        frontend_client: &Arc<FrontendClient>,
978        max_window_cnt: Option<usize>,
979    ) -> ExecuteOnceOutcome {
980        let guard = BatchingExecutionGuard::new(self.execution_lock.clone().lock_owned().await);
981        self.execute_once_with_guard(guard, engine, frontend_client, max_window_cnt)
982            .await
983    }
984
985    async fn execute_once_with_guard(
986        &self,
987        guard: BatchingExecutionGuard,
988        engine: &QueryEngineRef,
989        frontend_client: &Arc<FrontendClient>,
990        max_window_cnt: Option<usize>,
991    ) -> ExecuteOnceOutcome {
992        if let Some(execution) = &self.execution {
993            return execution
994                .clone()
995                .execute_once(guard, self, engine, frontend_client, max_window_cnt)
996                .await;
997        }
998        self.execute_once_default_unlocked(engine, frontend_client, max_window_cnt)
999            .await
1000    }
1001
1002    /// Executes the default flow evaluation. Caller owns the execution guard.
1003    async fn execute_once_default_unlocked(
1004        &self,
1005        engine: &QueryEngineRef,
1006        frontend_client: &Arc<FrontendClient>,
1007        max_window_cnt: Option<usize>,
1008    ) -> ExecuteOnceOutcome {
1009        let new_query = match self.gen_insert_plan_unlocked(engine, max_window_cnt).await {
1010            Ok(new_query) => new_query,
1011            Err(err) => {
1012                return ExecuteOnceOutcome {
1013                    new_query: None,
1014                    result: Err(err),
1015                };
1016            }
1017        };
1018
1019        if let Some(new_query) = new_query {
1020            debug!("Generate new query: {}", new_query.plan);
1021            let res = self
1022                .execute_logical_plan_unlocked(
1023                    engine,
1024                    frontend_client,
1025                    &new_query.plan,
1026                    &new_query.dirty_restore,
1027                    &new_query.coverage,
1028                )
1029                .await;
1030            if res.is_err() {
1031                self.handle_executed_query_failure(Some(&new_query));
1032            }
1033            ExecuteOnceOutcome {
1034                new_query: Some(new_query),
1035                result: res,
1036            }
1037        } else {
1038            debug!("Generate no query");
1039            ExecuteOnceOutcome {
1040                new_query: None,
1041                result: Ok(None),
1042            }
1043        }
1044    }
1045
1046    /// Generates the ordinary insert plan. Caller must reach this through the serialized path.
1047    async fn gen_insert_plan_unlocked(
1048        &self,
1049        engine: &QueryEngineRef,
1050        max_window_cnt: Option<usize>,
1051    ) -> Result<Option<PlanInfo>, Error> {
1052        self.gen_insert_plan_with_values_unlocked(engine, max_window_cnt, &BTreeMap::new(), false)
1053            .await
1054    }
1055
1056    /// Generates the insert plan with extension-owned typed sink values. Caller must own the
1057    /// serialized execution guard.
1058    pub async fn gen_insert_plan_with_values_unlocked(
1059        &self,
1060        engine: &QueryEngineRef,
1061        max_window_cnt: Option<usize>,
1062        values: &BTreeMap<String, ScalarValue>,
1063        force_full_snapshot: bool,
1064    ) -> Result<Option<PlanInfo>, Error> {
1065        let (table, df_schema) = get_table_info_df_schema(
1066            self.config.catalog_manager.clone(),
1067            self.config.sink_table_name.clone(),
1068        )
1069        .await?;
1070
1071        let table_meta = &table.table_info().meta;
1072        let merge_mode_last_non_null =
1073            is_merge_mode_last_non_null(&table_meta.options.extra_options);
1074        let primary_key_indices = table_meta.primary_key_indices.clone();
1075
1076        let new_query = self
1077            .gen_query_with_time_window_with_values(
1078                engine.clone(),
1079                &table.table_info().meta.schema,
1080                &primary_key_indices,
1081                merge_mode_last_non_null,
1082                max_window_cnt,
1083                values,
1084                force_full_snapshot,
1085            )
1086            .await?;
1087
1088        let Some(new_query) = new_query else {
1089            return Ok(None);
1090        };
1091
1092        // first check if all columns in input query exists in sink table
1093        // since insert into ref to names in record batch generate by given query
1094        let table_columns = df_schema
1095            .columns()
1096            .into_iter()
1097            .map(|c| c.name)
1098            .collect::<BTreeSet<_>>();
1099        for column in new_query.plan.schema().columns() {
1100            if !table_columns.contains(column.name()) {
1101                self.restore_dirty_windows_after_failure(&new_query);
1102                return InvalidQuerySnafu {
1103                    reason: format!(
1104                        "Column {} not found in sink table with columns {:?}",
1105                        column, table_columns
1106                    ),
1107                }
1108                .fail();
1109            }
1110        }
1111
1112        let table_provider = Arc::new(DfTableProviderAdapter::new(table));
1113        let table_source = Arc::new(DefaultTableSource::new(table_provider));
1114
1115        // update_at& time index placeholder (if exists) should have default value
1116        let plan = LogicalPlan::Dml(DmlStatement::new(
1117            datafusion_common::TableReference::Full {
1118                catalog: self.config.sink_table_name[0].clone().into(),
1119                schema: self.config.sink_table_name[1].clone().into(),
1120                table: self.config.sink_table_name[2].clone().into(),
1121            },
1122            table_source,
1123            WriteOp::Insert(datafusion_expr::dml::InsertOp::Append),
1124            Arc::new(new_query.plan.clone()),
1125        ));
1126        let insert_into_info = PlanInfo {
1127            plan,
1128            dirty_restore: new_query.dirty_restore,
1129            coverage: new_query.coverage,
1130        };
1131        let insert_into =
1132            match insert_into_info
1133                .plan
1134                .clone()
1135                .recompute_schema()
1136                .context(DatafusionSnafu {
1137                    context: "Failed to recompute schema",
1138                }) {
1139                Ok(insert_into) => insert_into,
1140                Err(err) => {
1141                    self.restore_dirty_windows_after_failure(&insert_into_info);
1142                    return Err(err);
1143                }
1144            };
1145
1146        Ok(Some(PlanInfo {
1147            plan: insert_into,
1148            dirty_restore: insert_into_info.dirty_restore,
1149            coverage: insert_into_info.coverage,
1150        }))
1151    }
1152
1153    pub async fn create_table(
1154        &self,
1155        frontend_client: &Arc<FrontendClient>,
1156        expr: CreateTableExpr,
1157    ) -> Result<(), Error> {
1158        let catalog = &self.config.sink_table_name[0];
1159        let schema = &self.config.sink_table_name[1];
1160        frontend_client
1161            .create(expr.clone(), catalog, schema)
1162            .await?;
1163        Ok(())
1164    }
1165
1166    /// Applies checkpoint state updates around the raw insert-plan execution.
1167    async fn execute_logical_plan_unlocked(
1168        &self,
1169        engine: &QueryEngineRef,
1170        frontend_client: &Arc<FrontendClient>,
1171        plan: &LogicalPlan,
1172        dirty_restore: &DirtyRestore,
1173        coverage: &QueryCoverage,
1174    ) -> Result<Option<(usize, Duration)>, Error> {
1175        let flow_id = self.config.flow_id;
1176        let Some((res, elapsed)) = self
1177            .execute_plan_unlocked(engine, frontend_client, plan, dirty_restore, coverage)
1178            .await?
1179        else {
1180            return Ok(None);
1181        };
1182
1183        if let Err(err) = &res {
1184            let decision = {
1185                let mut state = self.state.write().unwrap();
1186                let reason = Self::query_failure_reason(err, coverage);
1187                Self::apply_query_failure_to_state(&mut state, elapsed, coverage, reason)
1188            };
1189            if let Some(decision) = decision {
1190                Self::record_checkpoint_decision(flow_id, decision);
1191            }
1192        }
1193
1194        let res = res?;
1195        let (affected_rows, _) = res.output.extract_rows_and_cost();
1196        let decision = {
1197            let mut state = self.state.write().unwrap();
1198            Self::apply_query_result_to_state(&mut state, &res, elapsed, coverage)
1199        };
1200        Self::record_checkpoint_decision(flow_id, decision);
1201        Ok(Some((affected_rows, elapsed)))
1202    }
1203
1204    /// Executes a fully constructed insert plan without checkpoint-state mutation.
1205    /// Callers must own the serialized execution guard. Its unsafe-incremental safety fallback
1206    /// still restores the supplied dirty work before returning `Ok(None)`.
1207    pub async fn execute_plan_unlocked(
1208        &self,
1209        engine: &QueryEngineRef,
1210        frontend_client: &Arc<FrontendClient>,
1211        plan: &LogicalPlan,
1212        dirty_restore: &DirtyRestore,
1213        coverage: &QueryCoverage,
1214    ) -> Result<Option<(Result<OutputWithMetrics, Error>, Duration)>, Error> {
1215        let instant = Instant::now();
1216        let flow_id = self.config.flow_id;
1217
1218        debug!(
1219            "Executing flow {flow_id}(expire_after={:?} secs) with query {}",
1220            self.config.expire_after, &plan
1221        );
1222
1223        let catalog = &self.config.sink_table_name[0];
1224        let schema = &self.config.sink_table_name[1];
1225
1226        // fix all table ref by make it fully qualified, i.e. "table_name" => "catalog_name.schema_name.table_name"
1227        let plan = plan
1228            .clone()
1229            .transform_down_with_subqueries(|p| {
1230                if let LogicalPlan::TableScan(mut table_scan) = p {
1231                    let resolved = table_scan.table_name.resolve(catalog, schema);
1232                    table_scan.table_name = resolved.into();
1233                    Ok(Transformed::yes(LogicalPlan::TableScan(table_scan)))
1234                } else {
1235                    Ok(Transformed::no(p))
1236                }
1237            })
1238            .with_context(|_| DatafusionSnafu {
1239                context: format!("Failed to fix table ref in logical plan, plan={:?}", plan),
1240            })?
1241            .data;
1242
1243        // For incremental-mode SQL queries, attempt to rewrite the delta aggregate
1244        // plan into a safe delta-LEFT-JOIN-sink form before deciding on extensions.
1245        let incremental_plan = if coverage.is_incremental_delta() {
1246            self.prepare_plan_for_incremental(engine, &plan).await?
1247        } else {
1248            None
1249        };
1250        let incremental_safe = incremental_plan.is_some();
1251        if coverage.is_incremental_delta() && !incremental_safe {
1252            debug!(
1253                "Flow {flow_id} skipped unsafe incremental delta fallback; \
1254                 restored dirty signal instead of executing an unfiltered full snapshot"
1255            );
1256            self.restore_dirty_windows(dirty_restore);
1257            return Ok(None);
1258        }
1259        let plan = incremental_plan.unwrap_or_else(|| plan.clone());
1260        let plan = match &self.execution {
1261            Some(execution) => execution.rewrite_plan(self, plan)?,
1262            None => plan,
1263        };
1264
1265        let extensions = self
1266            .build_flow_query_extensions(incremental_safe, coverage.is_incremental_delta())
1267            .await?;
1268        let frontend_extensions = self.frontend_extensions();
1269        let extension_refs = extensions
1270            .iter()
1271            .map(|(key, value)| (*key, value.as_str()))
1272            .chain(
1273                frontend_extensions
1274                    .iter()
1275                    .map(|(key, value)| (key.as_str(), value.as_str())),
1276            )
1277            .collect::<Vec<_>>();
1278        let query_mode = if extensions
1279            .iter()
1280            .any(|(key, _)| *key == FLOW_INCREMENTAL_MODE)
1281        {
1282            CheckpointMode::Incremental
1283        } else {
1284            CheckpointMode::FullSnapshot
1285        };
1286        Self::record_query_mode(flow_id, query_mode);
1287        debug!(
1288            "Flow {flow_id} executing batching query with checkpoint_mode={}, extension_count={}",
1289            checkpoint_mode_label(query_mode),
1290            extensions.len()
1291        );
1292
1293        let mut peer_desc = None;
1294        let res = {
1295            let _timer = METRIC_FLOW_BATCHING_ENGINE_QUERY_TIME
1296                .with_label_values(&[flow_id.to_string().as_str()])
1297                .start_timer();
1298
1299            let req = if let Some((insert_to, insert_input_plan)) =
1300                breakup_insert_plan(&plan, catalog, schema)
1301            {
1302                if query_mode == CheckpointMode::FullSnapshot
1303                    && matches!(self.config.query_type, QueryType::Sql)
1304                    && self.config.flow_eval_interval.is_some()
1305                    && self.config.time_window_expr.is_none()
1306                {
1307                    // Evaluation-interval SQL flows without a time-window
1308                    // expression execute as full-query snapshots. Send these
1309                    // as SQL text instead of Substrait to avoid logical-plan
1310                    // round-trip issues around complex joins/unions/CTEs and
1311                    // duplicate field aliases. Keep ordinary SQL full snapshots
1312                    // on the existing InsertIntoPlan path because SQL unparsing
1313                    // is not valid for every planned aggregate shape yet.
1314                    // If the local SQL unparser does not support this plan,
1315                    // keep the previous InsertIntoPlan transport as a fallback.
1316                    match df_plan_to_sql(&insert_input_plan) {
1317                        Ok(select_sql) => {
1318                            let target_columns = format_insert_target_columns(&insert_input_plan);
1319                            let sql = format!(
1320                                "INSERT INTO {} ({}) {}",
1321                                TableReference::full(
1322                                    insert_to.catalog_name.as_str(),
1323                                    insert_to.schema_name.as_str(),
1324                                    insert_to.table_name.as_str(),
1325                                )
1326                                .to_quoted_string(),
1327                                target_columns,
1328                                select_sql
1329                            );
1330                            api::v1::QueryRequest {
1331                                query: Some(api::v1::query_request::Query::Sql(sql)),
1332                            }
1333                        }
1334                        Err(err) => {
1335                            debug!(
1336                                "Failed to unparse full-snapshot SQL flow {} plan; \
1337                                 falling back to InsertIntoPlan: {:?}",
1338                                flow_id, err
1339                            );
1340                            encode_insert_plan_request(insert_to, &insert_input_plan)?
1341                        }
1342                    }
1343                } else {
1344                    encode_insert_plan_request(insert_to, &insert_input_plan)?
1345                }
1346            } else {
1347                let message = DFLogicalSubstraitConvertor {}
1348                    .encode(&plan, DefaultSerializer)
1349                    .context(SubstraitEncodeLogicalPlanSnafu)?;
1350
1351                api::v1::QueryRequest {
1352                    query: Some(api::v1::query_request::Query::LogicalPlan(message.to_vec())),
1353                }
1354            };
1355
1356            let snapshot_seqs = coverage.snapshot_seqs();
1357            {
1358                let mut state = self.state.write().unwrap();
1359                state.record_start_time_if_first();
1360            }
1361            frontend_client
1362                .query_with_terminal_metrics(
1363                    catalog,
1364                    schema,
1365                    req,
1366                    &extension_refs,
1367                    &snapshot_seqs,
1368                    &mut peer_desc,
1369                )
1370                .await
1371        };
1372
1373        let elapsed = instant.elapsed();
1374        let peer_label = peer_desc
1375            .as_ref()
1376            .map(ToString::to_string)
1377            .unwrap_or_else(|| PeerDesc::default().to_string());
1378        if let Err(err) = &res {
1379            // This raw helper deliberately leaves checkpoint transitions to its caller. The
1380            // default wrapper performs the single OSS transition; extension owners can persist
1381            // their durable marker before updating in-memory state.
1382            warn!(
1383                "Failed to execute Flow {flow_id} on frontend {peer_label}, result: {err:?}, elapsed: {:?} with query: {}",
1384                elapsed, &plan
1385            );
1386        }
1387
1388        // record slow query
1389        if elapsed >= self.config.batch_opts.slow_query_threshold {
1390            warn!(
1391                "Flow {flow_id} on frontend {peer_label} executed for {:?} before complete, query: {}",
1392                elapsed, &plan
1393            );
1394            let flow_id = flow_id.to_string();
1395            METRIC_FLOW_BATCHING_ENGINE_SLOW_QUERY
1396                .with_label_values(&[flow_id.as_str(), peer_label.as_str()])
1397                .observe(elapsed.as_secs_f64());
1398        }
1399
1400        match res {
1401            Ok(res) => {
1402                let (affected_rows, _) = res.output.extract_rows_and_cost();
1403                if matches!(&res.output.data, common_query::OutputData::AffectedRows(_)) {
1404                    res.metrics.wait_ready().await;
1405                }
1406                if let Some(error) = res.metrics.completion_error() {
1407                    warn!("Flow {flow_id} completed with terminal metrics error: {error}");
1408                }
1409                debug!(
1410                    "Flow {flow_id} executed, affected_rows: {affected_rows:?}, elapsed: {:?}, watermark: {:?}",
1411                    elapsed,
1412                    res.region_watermark_map()
1413                );
1414                METRIC_FLOW_ROWS
1415                    .with_label_values(&[format!("{}-out-batching", flow_id).as_str()])
1416                    .inc_by(affected_rows as _);
1417                Ok(Some((Ok(res), elapsed)))
1418            }
1419            Err(err) => Ok(Some((Err(err), elapsed))),
1420        }
1421    }
1422
1423    /// Restore dirty windows consumed by a failed query so they are retried on
1424    /// the next execution.
1425    ///
1426    pub fn restore_dirty_windows(&self, dirty_restore: &DirtyRestore) {
1427        match dirty_restore {
1428            DirtyRestore::Scoped(filter) => self.restore_scoped_dirty_windows(filter),
1429            DirtyRestore::Unscoped(dirty_windows) => self
1430                .state
1431                .write()
1432                .unwrap()
1433                .dirty_time_windows
1434                .add_dirty_windows(dirty_windows),
1435        }
1436    }
1437
1438    /// Restore the dirty signal for a plan that was generated but failed before
1439    /// it could prove any checkpoint advancement.
1440    fn restore_dirty_windows_after_failure(&self, query: &PlanInfo) {
1441        self.restore_dirty_windows(&query.dirty_restore);
1442    }
1443
1444    /// Restore scoped windows through `TaskState` so fenced repair can decide
1445    /// whether they go back to pending repair or live dirty state.
1446    fn restore_scoped_dirty_windows(&self, filter: &FilterExprInfo) {
1447        self.state.write().unwrap().restore_scoped_windows(filter);
1448    }
1449
1450    /// Run a fallible scoped operation and restore its consumed windows if plan
1451    /// generation/rewrite fails before execution.
1452    fn restore_scoped_dirty_windows_on_err<T>(
1453        &self,
1454        filter: &FilterExprInfo,
1455        result: Result<T, Error>,
1456    ) -> Result<T, Error> {
1457        result.inspect_err(|_| {
1458            self.restore_scoped_dirty_windows(filter);
1459        })
1460    }
1461
1462    /// Restore an unscoped dirty signal consumed by an explicit full-query or
1463    /// incremental-delta plan.
1464    fn restore_unscoped_dirty_windows(&self, dirty_windows: &DirtyTimeWindows) {
1465        self.state
1466            .write()
1467            .unwrap()
1468            .dirty_time_windows
1469            .add_dirty_windows(dirty_windows);
1470    }
1471
1472    /// Run a fallible unscoped operation and restore the dirty signal if it
1473    /// fails before a query is executed.
1474    fn restore_unscoped_dirty_windows_on_err<T>(
1475        &self,
1476        dirty_windows: &DirtyTimeWindows,
1477        result: Result<T, Error>,
1478    ) -> Result<T, Error> {
1479        result.inspect_err(|_| {
1480            self.restore_unscoped_dirty_windows(dirty_windows);
1481        })
1482    }
1483
1484    /// Consume the live dirty signal for an unscoped query while keeping a copy
1485    /// that can be restored if planning or execution fails.
1486    fn drain_dirty_windows_signal(&self) -> (bool, DirtyTimeWindows) {
1487        let mut state = self.state.write().unwrap();
1488        let dirty_windows_to_restore = state.dirty_time_windows.clone();
1489        let is_dirty = !dirty_windows_to_restore.is_empty();
1490        state.dirty_time_windows.clean();
1491        (is_dirty, dirty_windows_to_restore)
1492    }
1493
1494    #[allow(clippy::too_many_arguments)]
1495    /// Build an unfiltered plan for explicit full-query or incremental-delta
1496    /// coverage. Callers pass the consumed dirty signal for failure restoration.
1497    async fn gen_unfiltered_plan_info(
1498        &self,
1499        engine: QueryEngineRef,
1500        query_ctx: QueryContextRef,
1501        sink_table_schema: Arc<Schema>,
1502        primary_key_indices: &[usize],
1503        allow_partial: bool,
1504        values: &BTreeMap<String, ScalarValue>,
1505        dirty_windows_to_restore: DirtyTimeWindows,
1506        retention_filter: Option<(&str, Timestamp, &'static str)>,
1507        coverage: QueryCoverage,
1508    ) -> Result<PlanInfo, Error> {
1509        let mut plan = self.restore_unscoped_dirty_windows_on_err(
1510            &dirty_windows_to_restore,
1511            gen_plan_with_matching_schema_and_values(
1512                &self.config.query,
1513                query_ctx,
1514                engine,
1515                sink_table_schema,
1516                primary_key_indices,
1517                allow_partial,
1518                values,
1519            )
1520            .await,
1521        )?;
1522
1523        if let Some((col_name, lower_bound, context)) = retention_filter {
1524            let lower = self.restore_unscoped_dirty_windows_on_err(
1525                &dirty_windows_to_restore,
1526                to_df_literal(lower_bound),
1527            )?;
1528            let retention_filter = col(col_name).gt_eq(lit(lower));
1529            let mut add_filter = AddFilterRewriter::new(retention_filter);
1530            plan = self.restore_unscoped_dirty_windows_on_err(
1531                &dirty_windows_to_restore,
1532                plan.clone()
1533                    .rewrite(&mut add_filter)
1534                    .with_context(|_| DatafusionSnafu {
1535                        context: format!(
1536                            "Failed to apply {context} expire_after filter to plan:\n {}\n",
1537                            plan
1538                        ),
1539                    })
1540                    .map(|rewrite| rewrite.data),
1541            )?;
1542        }
1543
1544        Ok(PlanInfo {
1545            plan,
1546            dirty_restore: DirtyRestore::Unscoped(dirty_windows_to_restore),
1547            coverage,
1548        })
1549    }
1550
1551    #[allow(clippy::too_many_arguments)]
1552    /// Build an unfiltered plan only when the live dirty signal was present;
1553    /// otherwise skip this round without querying.
1554    async fn gen_unfiltered_plan_info_if_dirty(
1555        &self,
1556        engine: QueryEngineRef,
1557        query_ctx: QueryContextRef,
1558        sink_table_schema: Arc<Schema>,
1559        primary_key_indices: &[usize],
1560        allow_partial: bool,
1561        values: &BTreeMap<String, ScalarValue>,
1562        retention_filter: Option<(&str, Timestamp, &'static str)>,
1563        coverage: QueryCoverage,
1564    ) -> Result<Option<PlanInfo>, Error> {
1565        let (is_dirty, dirty_windows_to_restore) = self.drain_dirty_windows_signal();
1566        if !is_dirty {
1567            debug!("Flow id={:?}, no new data, not update", self.config.flow_id);
1568            return Ok(None);
1569        }
1570
1571        self.gen_unfiltered_plan_info(
1572            engine,
1573            query_ctx,
1574            sink_table_schema,
1575            primary_key_indices,
1576            allow_partial,
1577            values,
1578            dirty_windows_to_restore,
1579            retention_filter,
1580            coverage,
1581        )
1582        .await
1583        .map(Some)
1584    }
1585
1586    fn handle_executed_query_failure(&self, query: Option<&PlanInfo>) {
1587        if let Some(query) = query {
1588            self.restore_dirty_windows_after_failure(query);
1589        }
1590    }
1591
1592    /// start executing query in a loop, break when receive shutdown signal
1593    ///
1594    /// any error will be logged when executing query.
1595    ///
1596    /// Dispatches to:
1597    /// - scheduled loop when `flow_eval_interval.is_some()`
1598    /// - adaptive dirty-window loop otherwise
1599    pub async fn start_executing_loop(
1600        &self,
1601        engine: QueryEngineRef,
1602        frontend_client: Arc<FrontendClient>,
1603    ) {
1604        if self.config.flow_eval_interval.is_some() {
1605            self.start_scheduled_loop(engine, frontend_client).await;
1606        } else {
1607            self.start_adaptive_loop(engine, frontend_client).await;
1608        }
1609    }
1610
1611    /// Scheduled batching loop for flows with `EVAL INTERVAL`.
1612    ///
1613    /// Uses the pre-parsed `EvalSchedule` from `TaskConfig` and selects due
1614    /// scheduled times using bounded catch-up semantics. Each scheduled time is the
1615    /// scheduled evaluation time used as logical `now()` for that attempt.
1616    /// Each attempt temporarily sets `flow.scheduled_time_millis` on the
1617    /// task's `QueryContext` and executes under the existing `execution_lock`.
1618    /// After every attempt (success, no-op, or failure) the in-memory
1619    /// cursor advances.
1620    async fn start_scheduled_loop(
1621        &self,
1622        engine: QueryEngineRef,
1623        frontend_client: Arc<FrontendClient>,
1624    ) {
1625        let flow_id_str = self.config.flow_id.to_string();
1626
1627        let schedule = match &self.config.eval_schedule {
1628            Some(s) => s.clone(),
1629            None => {
1630                let eval_interval_secs = self
1631                    .config
1632                    .flow_eval_interval
1633                    .map(|d| d.as_secs() as i64)
1634                    .expect("checked by caller");
1635
1636                // Fallback: no typed config provided. Compute defaults
1637                // anchored at epoch/start=0.
1638                match EvalSchedule::from_config(Some(eval_interval_secs), None) {
1639                    Ok(Some(s)) => s,
1640                    Ok(None) => {
1641                        warn!(
1642                            "Flow {}: EVAL INTERVAL set but no schedule parsed; exiting loop",
1643                            flow_id_str
1644                        );
1645                        return;
1646                    }
1647                    Err(e) => {
1648                        warn!(
1649                            "Flow {}: Failed to parse eval schedule: {}; exiting loop",
1650                            flow_id_str, e
1651                        );
1652                        return;
1653                    }
1654                }
1655            }
1656        };
1657
1658        // Initial cursor is one interval before start so the first due
1659        // scheduled time is `start_secs`. An unrepresentable difference is an
1660        // explicit error, never a saturated cursor that would silently skip
1661        // the first scheduled tick.
1662        let mut cursor_secs =
1663            match initial_schedule_cursor(schedule.start_secs, schedule.interval_secs) {
1664                Ok(cursor) => cursor,
1665                Err(e) => {
1666                    warn!(
1667                        "Flow {}: invalid eval schedule, exiting loop: {}",
1668                        flow_id_str, e
1669                    );
1670                    return;
1671                }
1672            };
1673
1674        info!(
1675            "Flow {}: entering scheduled loop, interval={}s, start={}, anchor={}, policy={:?}, max_runs={}, max_lag={}s",
1676            flow_id_str,
1677            schedule.interval_secs,
1678            schedule.start_secs,
1679            schedule.anchor_secs,
1680            schedule.missed_tick_policy,
1681            schedule.max_runs,
1682            schedule.max_lag_secs,
1683        );
1684
1685        loop {
1686            if self.is_shutdown_signaled() {
1687                break;
1688            }
1689
1690            let wall_now_secs = wall_clock_unix_secs();
1691
1692            let due = match select_due_scheduled_times(&schedule, cursor_secs, wall_now_secs) {
1693                Ok(d) => d,
1694                Err(e) => {
1695                    warn!(
1696                        "Flow {}: invalid eval schedule, exiting loop: {}",
1697                        flow_id_str, e
1698                    );
1699                    return;
1700                }
1701            };
1702
1703            if due.scheduled_times_secs.is_empty() {
1704                if due.skipped > 0 {
1705                    warn!(
1706                        "Flow {}: all {} due scheduled times skipped by max-lag, advancing cursor to wall-clock ({wall_now_secs}) to avoid re-skipping",
1707                        flow_id_str, due.skipped
1708                    );
1709                    cursor_secs = wall_now_secs;
1710                    continue;
1711                }
1712
1713                // No due yet — sleep until the next scheduled time.
1714                let next = match schedule.next_scheduled_time_after(cursor_secs) {
1715                    Ok(next) => next,
1716                    Err(e) => {
1717                        warn!(
1718                            "Flow {}: cannot advance eval schedule past cursor {cursor_secs}: {e}; exiting loop",
1719                            flow_id_str
1720                        );
1721                        return;
1722                    }
1723                };
1724                if next <= wall_now_secs {
1725                    // Shouldn't happen given select_due_scheduled_times returned empty,
1726                    // but guard against clock skew / logic error.
1727                    cursor_secs = wall_now_secs;
1728                    continue;
1729                }
1730                let wait_secs = match sleep_delta_secs(next, wall_now_secs) {
1731                    Ok(wait_secs) => wait_secs,
1732                    Err(e) => {
1733                        warn!(
1734                            "Flow {}: cannot sleep until next scheduled time {}: {e}; exiting loop",
1735                            flow_id_str, next
1736                        );
1737                        return;
1738                    }
1739                };
1740                let wait_dur = Duration::from_secs(wait_secs);
1741                debug!(
1742                    "Flow {}: no due scheduled times, sleeping for {}s until next scheduled time at {}",
1743                    flow_id_str, wait_secs, next
1744                );
1745                tokio::time::sleep(wait_dur).await;
1746                continue;
1747            }
1748
1749            if due.skipped > 0 {
1750                info!(
1751                    "Flow {}: {} due scheduled times, {} skipped (catch-up)",
1752                    flow_id_str,
1753                    due.scheduled_times_secs.len(),
1754                    due.skipped
1755                );
1756            }
1757
1758            // Execute scheduled times oldest → newest.
1759            for scheduled_time_secs in &due.scheduled_times_secs {
1760                if self.is_shutdown_signaled() {
1761                    break;
1762                }
1763
1764                METRIC_FLOW_BATCHING_ENGINE_START_QUERY_CNT
1765                    .with_label_values(&[&flow_id_str])
1766                    .inc();
1767
1768                let outcome = self
1769                    .execute_once_serialized_at_scheduled_time(
1770                        &engine,
1771                        &frontend_client,
1772                        *scheduled_time_secs,
1773                    )
1774                    .await;
1775
1776                // Advance cursor regardless of outcome.
1777                cursor_secs = *scheduled_time_secs;
1778
1779                match outcome.result {
1780                    Ok(Some((rows, elapsed))) => {
1781                        debug!(
1782                            "Flow {}: scheduled time {} completed, rows={}, elapsed={:?}",
1783                            flow_id_str, scheduled_time_secs, rows, elapsed
1784                        );
1785                    }
1786                    Ok(None) => {
1787                        debug!(
1788                            "Flow {}: scheduled time {} produced no query (no dirty signal or no-op)",
1789                            flow_id_str, scheduled_time_secs
1790                        );
1791                    }
1792                    Err(err) => {
1793                        warn!(
1794                            "Flow {}: scheduled time {} failed: {:?}",
1795                            flow_id_str, scheduled_time_secs, err
1796                        );
1797                        METRIC_FLOW_BATCHING_ENGINE_ERROR_CNT
1798                            .with_label_values(&[&flow_id_str])
1799                            .inc();
1800                        // Dirty-window restoration is handled by the
1801                        // existing `handle_executed_query_failure` inside
1802                        // default execution.
1803                    }
1804                }
1805            }
1806        }
1807    }
1808
1809    /// Existing adaptive dirty-window loop for flows without `EVAL INTERVAL`.
1810    async fn start_adaptive_loop(
1811        &self,
1812        engine: QueryEngineRef,
1813        frontend_client: Arc<FrontendClient>,
1814    ) {
1815        let flow_id_str = self.config.flow_id.to_string();
1816        let mut max_window_cnt = None;
1817        loop {
1818            if self.is_shutdown_signaled() {
1819                break;
1820            }
1821            METRIC_FLOW_BATCHING_ENGINE_START_QUERY_CNT
1822                .with_label_values(&[&flow_id_str])
1823                .inc();
1824
1825            let min_refresh = self.config.batch_opts.experimental_min_refresh_duration;
1826
1827            let outcome = self
1828                .execute_once_serialized_with_outcome(&engine, &frontend_client, max_window_cnt)
1829                .await;
1830
1831            match outcome.result {
1832                Ok(Some(_)) => {
1833                    max_window_cnt = max_window_cnt.map(|cnt| {
1834                        (cnt + 1).min(self.config.batch_opts.experimental_max_filter_num_per_query)
1835                    });
1836
1837                    let sleep_until = {
1838                        let state = self.state.write().unwrap();
1839
1840                        let time_window_size = self
1841                            .config
1842                            .time_window_expr
1843                            .as_ref()
1844                            .and_then(|t| *t.time_window_size());
1845
1846                        let prefer_short_incremental_cadence = state.checkpoint_mode()
1847                            == CheckpointMode::Incremental
1848                            && !state.is_incremental_disabled();
1849
1850                        state.get_next_start_query_time(
1851                            self.config.flow_id,
1852                            &time_window_size,
1853                            min_refresh,
1854                            Some(self.config.batch_opts.query_timeout),
1855                            self.config.batch_opts.experimental_max_filter_num_per_query,
1856                            prefer_short_incremental_cadence,
1857                        )
1858                    };
1859
1860                    tokio::time::sleep_until(sleep_until).await;
1861                }
1862                Ok(None) => {
1863                    debug!(
1864                        "Flow id = {:?} found no new data, sleep for {:?} then continue",
1865                        self.config.flow_id, min_refresh
1866                    );
1867                    tokio::time::sleep(min_refresh).await;
1868                    continue;
1869                }
1870                Err(err) => {
1871                    METRIC_FLOW_BATCHING_ENGINE_ERROR_CNT
1872                        .with_label_values(&[&flow_id_str])
1873                        .inc();
1874                    match outcome.new_query {
1875                        Some(query) => {
1876                            common_telemetry::error!(err; "Failed to execute query for flow={} with query: {}", self.config.flow_id, query.plan);
1877                            max_window_cnt = Some(1);
1878                        }
1879                        None => {
1880                            common_telemetry::error!(err; "Failed to generate query for flow={}", self.config.flow_id)
1881                        }
1882                    }
1883                    tokio::time::sleep(min_refresh).await;
1884                }
1885            }
1886        }
1887    }
1888
1889    /// Check whether the shutdown signal has been received.
1890    fn is_shutdown_signaled(&self) -> bool {
1891        let mut state = self.state.write().unwrap();
1892        match state.shutdown_rx.try_recv() {
1893            Ok(()) | Err(TryRecvError::Closed) => true,
1894            Err(TryRecvError::Empty) => false,
1895        }
1896    }
1897
1898    /// Execute one scheduled attempt, temporarily setting
1899    /// `flow.scheduled_time_millis` on the task's QueryContext so
1900    /// SQL/TQL `now()` resolves to the logical scheduled time.
1901    ///
1902    /// The extension is removed after the attempt so a later manual
1903    /// `flush_flow` does not reuse a stale scheduled time.
1904    async fn execute_once_serialized_at_scheduled_time(
1905        &self,
1906        engine: &QueryEngineRef,
1907        frontend_client: &Arc<FrontendClient>,
1908        scheduled_time_secs: i64,
1909    ) -> ExecuteOnceOutcome {
1910        let mut guard = BatchingExecutionGuard::new(self.execution_lock.clone().lock_owned().await);
1911
1912        // Convert to milliseconds before touching the task state so an
1913        // unrepresentable scheduled time fails as an explicit error without
1914        // ever installing a saturated (off-phase) extension value.
1915        let scheduled_time_millis = match scheduled_time_millis(scheduled_time_secs) {
1916            Ok(millis) => millis,
1917            Err(e) => {
1918                return ExecuteOnceOutcome {
1919                    new_query: None,
1920                    result: Err(e),
1921                };
1922            }
1923        };
1924
1925        // Clone the current QueryContext and add the scheduled time
1926        // extension, then swap it into the task state for this attempt.
1927        let old_ctx = {
1928            let mut state = self.state.write().unwrap();
1929            let old = state.query_ctx.clone();
1930            let mut new_ctx = (*old).clone();
1931            new_ctx.set_extension(
1932                query::options::FLOW_SCHEDULED_TIME_MILLIS,
1933                scheduled_time_millis.to_string(),
1934            );
1935            state.query_ctx = Arc::new(new_ctx);
1936            old
1937        };
1938        guard.restore_query_context(self.state.clone(), old_ctx);
1939
1940        // A collaborator may retain the guard in an owned child task. Its Drop
1941        // restores the scheduled context before releasing serialization.
1942        self.execute_once_with_guard(guard, engine, frontend_client, None)
1943            .await
1944    }
1945
1946    /// Generate the create table SQL
1947    ///
1948    /// the auto created table will automatically added a `update_at` Milliseconds DEFAULT now() column in the end
1949    /// (for compatibility with flow streaming mode)
1950    ///
1951    /// and it will use first timestamp column as time index, all other columns will be added as normal columns and nullable
1952    async fn gen_create_table_expr(
1953        &self,
1954        engine: QueryEngineRef,
1955    ) -> Result<CreateTableExpr, Error> {
1956        let query_ctx = self.state.read().unwrap().query_ctx.clone();
1957        let plan =
1958            sql_to_df_plan(query_ctx.clone(), engine.clone(), &self.config.query, true).await?;
1959        create_table_with_expr(&plan, &self.config.sink_table_name, &self.config.query_type)
1960    }
1961
1962    /// Incremental delta scans are unfiltered by dirty windows; the sequence
1963    /// range, not a time predicate, defines source correctness.
1964    fn should_use_unfiltered_incremental_delta(&self) -> bool {
1965        let state = self.state.read().unwrap();
1966        state.checkpoint_mode() == CheckpointMode::Incremental
1967            && !state.is_incremental_disabled()
1968            && matches!(self.config.query_type, QueryType::Sql)
1969    }
1970
1971    /// Generate the next ordinary plan and classify its coverage.
1972    async fn gen_query_with_time_window(
1973        &self,
1974        engine: QueryEngineRef,
1975        sink_table_schema: &Arc<Schema>,
1976        primary_key_indices: &[usize],
1977        allow_partial: bool,
1978        max_window_cnt: Option<usize>,
1979    ) -> Result<Option<PlanInfo>, Error> {
1980        self.gen_query_with_time_window_with_values(
1981            engine,
1982            sink_table_schema,
1983            primary_key_indices,
1984            allow_partial,
1985            max_window_cnt,
1986            &BTreeMap::new(),
1987            false,
1988        )
1989        .await
1990    }
1991
1992    /// Generate the next plan with extension-owned sink values and an optional forced full scan.
1993    #[allow(clippy::too_many_arguments)]
1994    async fn gen_query_with_time_window_with_values(
1995        &self,
1996        engine: QueryEngineRef,
1997        sink_table_schema: &Arc<Schema>,
1998        primary_key_indices: &[usize],
1999        allow_partial: bool,
2000        max_window_cnt: Option<usize>,
2001        values: &BTreeMap<String, ScalarValue>,
2002        force_full_snapshot: bool,
2003    ) -> Result<Option<PlanInfo>, Error> {
2004        let query_ctx = self.state.read().unwrap().query_ctx.clone();
2005        let start = SystemTime::now();
2006        let since_the_epoch = start
2007            .duration_since(UNIX_EPOCH)
2008            .expect("Time went backwards");
2009        let low_bound = self
2010            .config
2011            .expire_after
2012            .map(|e| since_the_epoch.as_secs() - e as u64)
2013            .unwrap_or(u64::MIN);
2014
2015        let low_bound = Timestamp::new_second(low_bound as i64);
2016
2017        let expire_time_window_bound = self
2018            .config
2019            .time_window_expr
2020            .as_ref()
2021            .map(|expr| expr.eval(low_bound))
2022            .transpose()?;
2023
2024        let (expire_lower_bound, expire_upper_bound) = match (
2025            expire_time_window_bound,
2026            &self.config.query_type,
2027        ) {
2028            (Some((Some(l), Some(u))), QueryType::Sql) => (l, u),
2029            (None, QueryType::Sql) if self.config.flow_eval_interval.is_none() => {
2030                return UnexpectedSnafu {
2031                    reason: format!(
2032                        "Flow id={} reached execution without a time-window expression or EVAL INTERVAL; create-flow validation should have rejected it",
2033                        self.config.flow_id
2034                    ),
2035                }
2036                .fail();
2037            }
2038            _ => {
2039                // Explicit full-query flows (TQL and evaluation-interval SQL
2040                // plans whose shape cannot be safely dirty-window pruned) are
2041                // allowed to run as unfiltered full snapshots. This is distinct
2042                // from using unfiltered full as a fallback after scoped repair or
2043                // incremental rewrite failed.
2044                let (_, dirty_windows_to_restore) = self.drain_dirty_windows_signal();
2045
2046                let plan_info = self
2047                    .gen_unfiltered_plan_info(
2048                        engine,
2049                        query_ctx,
2050                        sink_table_schema.clone(),
2051                        primary_key_indices,
2052                        allow_partial,
2053                        values,
2054                        dirty_windows_to_restore,
2055                        None,
2056                        QueryCoverage::UnfilteredFull,
2057                    )
2058                    .await?;
2059
2060                return Ok(Some(plan_info));
2061            }
2062        };
2063
2064        debug!(
2065            "Flow id = {:?}, found time window: precise_lower_bound={:?}, precise_upper_bound={:?} with dirty time windows: {:?}",
2066            self.config.flow_id,
2067            expire_lower_bound,
2068            expire_upper_bound,
2069            self.state.read().unwrap().dirty_time_windows
2070        );
2071        let window_size = expire_upper_bound
2072            .sub(&expire_lower_bound)
2073            .with_context(|| UnexpectedSnafu {
2074                reason: format!(
2075                    "Can't get window size from {expire_upper_bound:?} - {expire_lower_bound:?}"
2076                ),
2077            })?;
2078        let col_name = self
2079            .config
2080            .time_window_expr
2081            .as_ref()
2082            .map(|expr| expr.column_name.clone())
2083            .with_context(|| UnexpectedSnafu {
2084                reason: format!(
2085                    "Flow id={:?}, Failed to get column name from time window expr",
2086                    self.config.flow_id
2087                ),
2088            })?;
2089
2090        if force_full_snapshot {
2091            // Durable execution cannot trust an uncertain prior checkpoint. Re-run an
2092            // unfiltered full snapshot, but retain the normal expiry predicate and the consumed
2093            // dirty signal so every planning/execution failure remains reversible.
2094            let retention_filter = self.config.expire_after.map(|_| {
2095                (
2096                    col_name.as_str(),
2097                    expire_lower_bound,
2098                    "forced full snapshot",
2099                )
2100            });
2101            let (_, dirty_windows_to_restore) = self.drain_dirty_windows_signal();
2102            return self
2103                .gen_unfiltered_plan_info(
2104                    engine,
2105                    query_ctx,
2106                    sink_table_schema.clone(),
2107                    primary_key_indices,
2108                    allow_partial,
2109                    values,
2110                    dirty_windows_to_restore,
2111                    retention_filter,
2112                    QueryCoverage::UnfilteredFull,
2113                )
2114                .await
2115                .map(Some);
2116        }
2117
2118        if self.should_use_unfiltered_incremental_delta() {
2119            // In incremental mode, source correctness is defined by the
2120            // per-region sequence range `(checkpoint, scan-open snapshot]`, not
2121            // by dirty-window predicates. Dirty windows are only a scheduling
2122            // signal here. Applying a stale dirty-window filter to the source can
2123            // exclude rows that are inside the returned watermark and make a
2124            // checkpoint advance skip them forever. The sink side is also left
2125            // unfiltered by dirty windows; the incremental rewrite joins the
2126            // delta groups with the full sink state for correctness. Future
2127            // dynamic filters can prune sink reads as a pure optimization.
2128            let retention_filter = self
2129                .config
2130                .expire_after
2131                .map(|_| (col_name.as_str(), expire_lower_bound, "incremental"));
2132            return self
2133                .gen_unfiltered_plan_info_if_dirty(
2134                    engine,
2135                    query_ctx,
2136                    sink_table_schema.clone(),
2137                    primary_key_indices,
2138                    allow_partial,
2139                    values,
2140                    retention_filter,
2141                    QueryCoverage::IncrementalDelta,
2142                )
2143                .await;
2144        }
2145
2146        let (expr, coverage) = {
2147            let mut state = self.state.write().unwrap();
2148            let window_cnt = max_window_cnt
2149                .unwrap_or(self.config.batch_opts.experimental_max_filter_num_per_query);
2150            let expr = state.gen_scoped_filter_exprs(
2151                &col_name,
2152                Some(expire_lower_bound),
2153                window_size,
2154                window_cnt,
2155                self.config.flow_id,
2156                Some(self),
2157            )?;
2158            let repair_high = state
2159                .pending_fenced_repair()
2160                .map(|repair| repair.high().clone());
2161            let coverage = if let Some(high) = repair_high {
2162                QueryCoverage::FencedRepairChunk { high }
2163            } else {
2164                QueryCoverage::ScopedBaseRepair
2165            };
2166            (expr, coverage)
2167        };
2168
2169        let Some(expr) = expr else {
2170            // no new data, hence no need to update
2171            debug!("Flow id={:?}, no new data, not update", self.config.flow_id);
2172            return Ok(None);
2173        };
2174
2175        let filter_sql = expr_to_sql(&expr.expr)
2176            .map(|sql| sql.to_string())
2177            .unwrap_or_else(|err| format!("<failed to format filter expr: {err}>"));
2178
2179        debug!(
2180            "Flow id={:?}, Generated filter expr: {:?}",
2181            self.config.flow_id, filter_sql
2182        );
2183
2184        let mut add_filter = AddFilterRewriter::new(expr.expr.clone());
2185        let mut add_auto_column = ColumnMatcherRewriter::new_with_values(
2186            sink_table_schema.clone(),
2187            primary_key_indices.to_vec(),
2188            allow_partial,
2189            values.clone(),
2190        );
2191
2192        let plan = self.restore_scoped_dirty_windows_on_err(
2193            &expr,
2194            sql_to_df_plan(query_ctx.clone(), engine.clone(), &self.config.query, false).await,
2195        )?;
2196        let rewrite = self.restore_scoped_dirty_windows_on_err(
2197            &expr,
2198            plan.clone()
2199                .rewrite(&mut add_filter)
2200                .and_then(|p| p.data.rewrite(&mut add_auto_column))
2201                .with_context(|_| DatafusionSnafu {
2202                    context: format!("Failed to rewrite plan:\n {}\n", plan),
2203                })
2204                .map(|rewrite| rewrite.data),
2205        )?;
2206        // only apply optimize after complex rewrite is done
2207        let new_plan = self.restore_scoped_dirty_windows_on_err(
2208            &expr,
2209            apply_df_optimizer(rewrite, &query_ctx).await,
2210        )?;
2211
2212        let info = PlanInfo {
2213            plan: new_plan.clone(),
2214            dirty_restore: DirtyRestore::Scoped(expr),
2215            coverage,
2216        };
2217
2218        Ok(Some(info))
2219    }
2220}
2221
2222#[cfg(test)]
2223mod test;