Skip to main content

flow/
adapter.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Flow source schema management and stateless streaming execution.
16
17use std::collections::BTreeMap;
18use std::sync::Arc;
19#[cfg(test)]
20use std::sync::atomic::AtomicUsize;
21use std::sync::atomic::{AtomicBool, Ordering};
22
23use common_base::memory_limit::MemoryLimit;
24use common_config::Configurable;
25use common_error::ext::BoxedError;
26use common_meta::distributed_time_constants::BASE_HEARTBEAT_INTERVAL;
27use common_meta::key::TableMetadataManagerRef;
28use common_options::memory::MemoryOptions;
29use common_recordbatch::map_dictionary_to_values_data_type;
30use common_stat::get_total_cpu_cores;
31use common_telemetry::logging::{LoggingOptions, TracingOptions};
32use common_telemetry::{error, info};
33use datafusion_common::TableReference;
34use datafusion_common::tree_node::TreeNode;
35use datafusion_expr::logical_plan::Distinct;
36use datafusion_expr::{Expr, LogicalPlan};
37use datatypes::schema::{ColumnSchema, SchemaRef};
38use itertools::Itertools;
39use meta_client::MetaClientOptions;
40use query::QueryEngine;
41use query::options::QueryOptions;
42use serde::{Deserialize, Serialize};
43use servers::grpc::GrpcOptions;
44use servers::http::HttpOptions;
45use snafu::{IntoError, OptionExt, ResultExt, ensure};
46use store_api::storage::{ConcreteDataType, RegionId};
47use tokio::sync::RwLock;
48use tokio::time::Instant;
49
50use crate::adapter::stateless::StatelessFlow;
51use crate::adapter::table_source::ManagedTableSource;
52use crate::adapter::util::{
53    relation_desc_to_column_schemas_with_fallback, table_info_value_to_relation_desc,
54};
55use crate::batching_mode::BatchingModeOptions;
56use crate::batching_mode::frontend_client::FrontendClient;
57use crate::batching_mode::utils::sql_to_df_plan;
58use crate::error::{
59    DatafusionSnafu, Error, ExternalSnafu, FlowNotFoundSnafu, InsertIntoFlowSnafu, InternalSnafu,
60    InvalidQuerySnafu, UnexpectedSnafu,
61};
62use crate::repr::{ColumnType, DiffRow, RelationDesc, Row};
63use crate::{CreateFlowArgs, FlowId, TableName};
64
65pub(crate) mod flownode_impl;
66pub(crate) mod stateless;
67pub(crate) mod table_source;
68#[cfg(test)]
69mod tests;
70pub(crate) mod util;
71
72/// Converts a retained plan's output to logical schemas. Only a direct source
73/// column (optionally wrapped in an alias) carries source semantics; computed
74/// expressions deliberately use the physical output field metadata instead.
75fn output_column_schemas(
76    plan: &LogicalPlan,
77    source_schema: &SchemaRef,
78) -> Result<(Vec<ColumnSchema>, Vec<Option<usize>>), Error> {
79    // Plain DISTINCT preserves the selected columns and their source semantics. Inspect its
80    // input for lineage, while leaving computed expressions without source metadata.
81    let expressions = match plan {
82        LogicalPlan::Projection(projection) => Some(&projection.expr),
83        LogicalPlan::Distinct(Distinct::All(input)) => match input.as_ref() {
84            LogicalPlan::Projection(projection) => Some(&projection.expr),
85            _ => None,
86        },
87        _ => None,
88    };
89    let is_pass_through = matches!(plan, LogicalPlan::TableScan(_) | LogicalPlan::Filter(_))
90        || matches!(plan, LogicalPlan::Distinct(Distinct::All(input)) if matches!(input.as_ref(), LogicalPlan::TableScan(_) | LogicalPlan::Filter(_)));
91    let mut lineage = Vec::new();
92    let columns = plan
93        .schema()
94        .fields()
95        .iter()
96        .enumerate()
97        .map(|(idx, field)| {
98            let mut column = ColumnSchema::try_from(field.as_ref())
99                .map_err(BoxedError::new)
100                .context(ExternalSnafu)?;
101            let source_index = expressions
102                .and_then(|exprs| {
103                    let expr = exprs.get(idx)?;
104                    let expr = match expr {
105                        Expr::Column(column) => Some(column),
106                        Expr::Alias(alias) => match alias.expr.as_ref() {
107                            Expr::Column(column) => Some(column),
108                            _ => None,
109                        },
110                        _ => None,
111                    }?;
112                    source_schema.column_index_by_name(&expr.name)
113                })
114                .or_else(|| {
115                    is_pass_through
116                        .then(|| source_schema.column_index_by_name(field.name()))
117                        .flatten()
118                });
119            if let Some(source_index) = source_index {
120                let mut source = source_schema.column_schemas()[source_index].clone();
121                source.name = field.name().clone();
122                source.data_type = map_dictionary_to_values_data_type(&source.data_type);
123                column = source;
124            } else {
125                column.data_type = map_dictionary_to_values_data_type(&column.data_type);
126                // DataFusion may propagate field metadata through a computed expression.  Such
127                // an expression has no source-column lineage and must not inherit its PK/time
128                // semantics into the sink relation.
129                if expressions.is_some() {
130                    column = column.with_time_index(false);
131                }
132            }
133            lineage.push(source_index);
134            Ok(column)
135        })
136        .collect::<Result<Vec<_>, Error>>()?;
137    Ok((columns, lineage))
138}
139
140fn relation_desc_from_output(
141    columns: &[ColumnSchema],
142    lineage: &[Option<usize>],
143    source_primary_key_indices: &[usize],
144) -> RelationDesc {
145    let keys = source_primary_key_indices
146        .iter()
147        .filter_map(|source_index| {
148            lineage
149                .iter()
150                .position(|index| index == &Some(*source_index))
151        })
152        .collect_vec();
153    let time_index = columns.iter().position(ColumnSchema::is_time_index);
154    RelationDesc {
155        typ: crate::repr::RelationType {
156            column_types: columns
157                .iter()
158                .map(|column| ColumnType::new(column.data_type.clone(), column.is_nullable()))
159                .collect(),
160            keys: if keys.is_empty() {
161                vec![]
162            } else {
163                vec![crate::repr::Key::from(keys)]
164            },
165            time_index,
166            auto_columns: vec![],
167        },
168        names: columns
169            .iter()
170            .map(|column| Some(column.name.clone()))
171            .collect(),
172    }
173}
174
175fn default_num_workers() -> usize {
176    get_total_cpu_cores().div_ceil(2)
177}
178
179/// Returns whether an existing sink uses the legacy explicit source-time-index layout.
180///
181/// This check is deliberately separate from [`resolve_sink_layout`]: ordinary auto-column
182/// resolution remains the first choice, and this compatibility path is only for a sink whose
183/// trailing column is a real, user-defined time index.
184pub(crate) fn is_explicit_source_timestamp_compatibility(
185    output_schema: &[ColumnSchema],
186    output_lineage: &[Option<usize>],
187    sink_schema: &[ColumnSchema],
188    source_schema: &SchemaRef,
189) -> bool {
190    if sink_schema.len() != output_schema.len() + 1 || output_lineage.len() != output_schema.len() {
191        return false;
192    }
193    if !output_schema
194        .iter()
195        .zip(&sink_schema[..output_schema.len()])
196        .all(|(output, sink)| output.data_type == sink.data_type)
197    {
198        return false;
199    }
200
201    let Some(source_timestamp_index) = source_schema.timestamp_index() else {
202        return false;
203    };
204    let sink_timestamp = &sink_schema[output_schema.len()];
205    sink_timestamp.data_type == source_schema.column_schemas()[source_timestamp_index].data_type
206        && sink_timestamp.data_type.is_timestamp()
207        && sink_timestamp.is_time_index()
208        && sink_timestamp.default_constraint().is_some()
209        && sink_timestamp.name != AUTO_CREATED_UPDATE_AT_TS_COL
210        && sink_timestamp.name != AUTO_CREATED_PLACEHOLDER_TS_COL
211        && !output_lineage.contains(&Some(source_timestamp_index))
212}
213
214pub const AUTO_CREATED_PLACEHOLDER_TS_COL: &str = "__ts_placeholder";
215pub const AUTO_CREATED_UPDATE_AT_TS_COL: &str = "update_at";
216
217/// Resolves the columns appended by the flow, validating the complete output/sink layout.
218///
219/// A sink with the same arity as the query is an ordinary sink, even when its last
220/// column happens to be named `update_at`. Auto columns are only inferred from the
221/// arity difference and their exact trailing layout.
222pub(crate) fn resolve_sink_layout(
223    output_schema: &[ColumnSchema],
224    sink_schema: &[ColumnSchema],
225) -> Result<Vec<ColumnSchema>, Error> {
226    ensure!(
227        sink_schema.len() >= output_schema.len() && sink_schema.len() - output_schema.len() <= 2,
228        InvalidQuerySnafu {
229            reason: format!(
230                "Flow output has {} columns, but sink has {} columns; only zero, one, or two trailing auto columns are supported",
231                output_schema.len(),
232                sink_schema.len()
233            )
234        }
235    );
236    let suffix_len = sink_schema.len() - output_schema.len();
237    for (idx, (output, sink)) in output_schema.iter().zip(sink_schema.iter()).enumerate() {
238        ensure!(
239            output.data_type == sink.data_type,
240            InvalidQuerySnafu {
241                reason: format!(
242                    "Flow output column {idx} has type {:?}, but sink column {} has type {:?}",
243                    output.data_type, sink.name, sink.data_type
244                )
245            }
246        );
247    }
248    let suffix = &sink_schema[output_schema.len()..];
249    match suffix_len {
250        0 => Ok(vec![]),
251        1 => {
252            let column = &suffix[0];
253            ensure!(
254                column.name == AUTO_CREATED_UPDATE_AT_TS_COL && column.data_type.is_timestamp(),
255                InvalidQuerySnafu {
256                    reason: format!(
257                        "The trailing sink column must be timestamp {}",
258                        AUTO_CREATED_UPDATE_AT_TS_COL
259                    )
260                }
261            );
262            Ok(suffix.to_vec())
263        }
264        2 => {
265            let update_at = &suffix[0];
266            let placeholder = &suffix[1];
267            ensure!(
268                update_at.name == AUTO_CREATED_UPDATE_AT_TS_COL
269                    && update_at.data_type.is_timestamp()
270                    && placeholder.name == AUTO_CREATED_PLACEHOLDER_TS_COL
271                    && placeholder.data_type.is_timestamp()
272                    && placeholder.is_time_index(),
273                InvalidQuerySnafu {
274                    reason: "The two trailing sink columns must be timestamp update_at followed by timestamp time-index __ts_placeholder".to_string()
275                }
276            );
277            Ok(suffix.to_vec())
278        }
279        _ => unreachable!(),
280    }
281}
282
283/// Legacy helper for callers that only have a sink schema. New stateless flows
284/// resolve the suffix against both schemas and store it in `StatelessFlow`.
285pub(crate) fn sink_output_column_count(sink_schema: &[ColumnSchema]) -> Result<usize, Error> {
286    let mut count = sink_schema.len();
287    if sink_schema
288        .last()
289        .is_some_and(|column| column.name == AUTO_CREATED_PLACEHOLDER_TS_COL)
290    {
291        let placeholder = &sink_schema[count - 1];
292        ensure!(
293            placeholder.data_type.is_timestamp() && placeholder.is_time_index(),
294            InvalidQuerySnafu {
295                reason: format!(
296                    "Auto-created sink column {} must be a timestamp time index",
297                    AUTO_CREATED_PLACEHOLDER_TS_COL
298                )
299            }
300        );
301        count -= 1;
302    }
303    if sink_schema
304        .get(count.saturating_sub(1))
305        .is_some_and(|column| column.name == AUTO_CREATED_UPDATE_AT_TS_COL)
306    {
307        ensure!(
308            sink_schema[count - 1].data_type.is_timestamp(),
309            InvalidQuerySnafu {
310                reason: format!(
311                    "Auto-created sink column {} must be a timestamp",
312                    AUTO_CREATED_UPDATE_AT_TS_COL
313                )
314            }
315        );
316        count -= 1;
317    }
318    Ok(count)
319}
320
321pub(crate) fn validate_sink_layout(
322    output_schema: &[ColumnSchema],
323    sink_schema: &[ColumnSchema],
324) -> Result<(), Error> {
325    resolve_sink_layout(output_schema, sink_schema).map(|_| ())
326}
327
328pub(crate) fn validate_auto_column_names(output_schema: &[ColumnSchema]) -> Result<(), Error> {
329    for column in output_schema {
330        ensure!(
331            column.name != AUTO_CREATED_UPDATE_AT_TS_COL
332                && column.name != AUTO_CREATED_PLACEHOLDER_TS_COL,
333            InvalidQuerySnafu {
334                reason: format!(
335                    "Flow output column {} is reserved for an auto-created sink column",
336                    column.name
337                )
338            }
339        );
340    }
341    Ok(())
342}
343
344pub(crate) fn validate_sink_layout_with_suffix(
345    output_schema: &[ColumnSchema],
346    sink_schema: &[ColumnSchema],
347    suffix: &[ColumnSchema],
348) -> Result<(), Error> {
349    ensure!(
350        sink_schema.len() == output_schema.len() + suffix.len()
351            && sink_schema[output_schema.len()..] == *suffix,
352        InvalidQuerySnafu {
353            reason: "Stored sink auto-column layout no longer matches the sink schema".to_string()
354        }
355    );
356    for (idx, (output, sink)) in output_schema
357        .iter()
358        .zip(&sink_schema[..output_schema.len()])
359        .enumerate()
360    {
361        ensure!(
362            output.data_type == sink.data_type,
363            InvalidQuerySnafu {
364                reason: format!(
365                    "Flow output column {idx} has type {:?}, but sink column {} has type {:?}",
366                    output.data_type, sink.name, sink.data_type
367                )
368            }
369        );
370    }
371    Ok(())
372}
373
374#[derive(Clone, Debug, Serialize, Deserialize, PartialEq)]
375#[serde(default)]
376pub struct FlowConfig {
377    /// Deprecated and ignored. Flow workers have been removed.
378    #[deprecated(note = "flow workers have been removed; this field is ignored")]
379    #[serde(default = "default_num_workers")]
380    pub num_workers: usize,
381    pub batching_mode: BatchingModeOptions,
382}
383
384#[allow(deprecated)]
385impl Default for FlowConfig {
386    fn default() -> Self {
387        Self {
388            num_workers: get_total_cpu_cores().div_ceil(2),
389            batching_mode: BatchingModeOptions::default(),
390        }
391    }
392}
393
394#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
395#[serde(default)]
396pub struct FlownodeOptions {
397    pub node_id: Option<u64>,
398    pub flow: FlowConfig,
399    pub grpc: GrpcOptions,
400    pub http: HttpOptions,
401    pub meta_client: Option<MetaClientOptions>,
402    pub logging: LoggingOptions,
403    pub tracing: TracingOptions,
404    pub query: QueryOptions,
405    pub memory: MemoryOptions,
406}
407
408impl Default for FlownodeOptions {
409    fn default() -> Self {
410        Self {
411            node_id: None,
412            flow: FlowConfig::default(),
413            grpc: GrpcOptions::default().with_bind_addr("127.0.0.1:3004"),
414            http: HttpOptions::default(),
415            meta_client: None,
416            logging: LoggingOptions::default(),
417            tracing: TracingOptions::default(),
418            query: QueryOptions {
419                parallelism: 1,
420                allow_query_fallback: false,
421                memory_pool_size: MemoryLimit::default(),
422                enable_per_region_metrics: false,
423                ..Default::default()
424            },
425            memory: MemoryOptions::default(),
426        }
427    }
428}
429
430impl Configurable for FlownodeOptions {
431    fn validate_sanitize(&mut self) -> common_config::error::Result<()> {
432        Ok(())
433    }
434}
435
436pub type FlowStreamingEngineRef = Arc<StreamingEngine>;
437
438#[derive(Default)]
439struct StatelessFlowRuntime {
440    flow: Option<Arc<StatelessFlow>>,
441    failed_rebuild: Option<(u32, Instant)>,
442}
443
444struct StatelessFlowSlot {
445    runtime: Arc<RwLock<StatelessFlowRuntime>>,
446    /// Cleared at registry-detach time, fencing creators that already captured this slot.
447    active: std::sync::atomic::AtomicBool,
448    #[cfg(test)]
449    rebuild_attempts: AtomicUsize,
450}
451
452fn validate_captured_slot(
453    slot: &StatelessFlowSlot,
454    current_source_table_id: Option<table::metadata::TableId>,
455    expected_table_id: table::metadata::TableId,
456    flow_id: FlowId,
457) -> Result<(), Error> {
458    ensure!(
459        slot.active.load(Ordering::Acquire),
460        FlowNotFoundSnafu { id: flow_id }
461    );
462    let current_source_table_id =
463        current_source_table_id.context(FlowNotFoundSnafu { id: flow_id })?;
464    ensure!(
465        current_source_table_id == expected_table_id,
466        InvalidQuerySnafu {
467            reason: format!("Flow {flow_id} source table changed while it was selected")
468        }
469    );
470    Ok(())
471}
472
473pub struct StreamingEngine {
474    pub query_engine: Arc<dyn QueryEngine>,
475    pub frontend_client: Arc<FrontendClient>,
476    table_info_source: ManagedTableSource,
477    stateless_flows: RwLock<BTreeMap<FlowId, Arc<StatelessFlowSlot>>>,
478    pub node_id: Option<u32>,
479}
480
481impl StreamingEngine {
482    pub fn new(
483        node_id: Option<u32>,
484        query_engine: Arc<dyn QueryEngine>,
485        table_meta: TableMetadataManagerRef,
486        frontend_client: Arc<FrontendClient>,
487    ) -> Self {
488        let table_info_source = ManagedTableSource::new(
489            table_meta.table_info_manager().clone(),
490            table_meta.table_name_manager().clone(),
491        );
492        Self {
493            query_engine,
494            frontend_client,
495            table_info_source,
496            stateless_flows: Default::default(),
497            node_id,
498        }
499    }
500
501    pub async fn handle_write_request(
502        &self,
503        region_id: RegionId,
504        rows: Vec<DiffRow>,
505        batch_datatypes: &[ConcreteDataType],
506        source_schema_version: u32,
507    ) -> Result<(), Error> {
508        let table_id = region_id.table_id();
509        let flow_ids = self.flow_ids_for_table(table_id).await;
510        let mut failed_flow_ids = Vec::new();
511        let mut first_error = None;
512        for (flow_id, slot) in flow_ids {
513            let result = self
514                .execute_flow(
515                    flow_id,
516                    slot,
517                    table_id,
518                    rows.clone(),
519                    batch_datatypes,
520                    source_schema_version,
521                )
522                .await;
523            if let Err(err) = result {
524                error!(err; "Failed to insert into flow={}, region_id={}", flow_id, region_id);
525                failed_flow_ids.push(flow_id);
526                if first_error.is_none() {
527                    first_error = Some(BoxedError::new(err));
528                }
529            }
530        }
531        match first_error {
532            Some(source) => Err(InsertIntoFlowSnafu {
533                region_id: u64::from(region_id),
534                flow_ids: failed_flow_ids,
535            }
536            .into_error(source)),
537            None => Ok(()),
538        }
539    }
540
541    async fn flow_ids_for_table(
542        &self,
543        table_id: table::metadata::TableId,
544    ) -> Vec<(FlowId, Arc<StatelessFlowSlot>)> {
545        let slots = self
546            .stateless_flows
547            .read()
548            .await
549            .iter()
550            .map(|(id, slot)| (*id, Arc::clone(slot)))
551            .collect::<Vec<_>>();
552        let mut flows = Vec::new();
553        for (id, slot) in slots {
554            if slot
555                .runtime
556                .read()
557                .await
558                .flow
559                .as_ref()
560                .is_some_and(|flow| flow.source_table_id == table_id)
561            {
562                flows.push((id, slot));
563            }
564        }
565        flows
566    }
567
568    async fn execute_flow(
569        &self,
570        flow_id: FlowId,
571        slot: Arc<StatelessFlowSlot>,
572        expected_table_id: table::metadata::TableId,
573        rows: Vec<DiffRow>,
574        batch_datatypes: &[ConcreteDataType],
575        source_schema_version: u32,
576    ) -> Result<usize, Error> {
577        // Keep the captured lifecycle slot rather than resolving the ID again after a drop.
578        let mut guard = slot.runtime.clone().read_owned().await;
579        let current = guard
580            .flow
581            .as_ref()
582            .context(FlowNotFoundSnafu { id: flow_id })?;
583        validate_captured_slot(
584            &slot,
585            Some(current.source_table_id),
586            expected_table_id,
587            flow_id,
588        )?;
589        if current.source_schema_version != source_schema_version {
590            if guard
591                .failed_rebuild
592                .as_ref()
593                .is_some_and(|(version, until)| {
594                    *version == source_schema_version && *until > Instant::now()
595                })
596            {
597                return InvalidQuerySnafu {
598                    reason: format!(
599                        "Flow {flow_id} schema rebuild for source version {source_schema_version} is cooling down"
600                    ),
601                }
602                .fail();
603            }
604            drop(guard);
605            // Serialize rebuilds with definition replacement using the existing publication
606            // lease. Another write may already have rebuilt this schema while we waited.
607            let mut published = slot.runtime.clone().write_owned().await;
608            let current = published
609                .flow
610                .as_ref()
611                .context(FlowNotFoundSnafu { id: flow_id })?;
612            validate_captured_slot(
613                &slot,
614                Some(current.source_table_id),
615                expected_table_id,
616                flow_id,
617            )?;
618            if current.source_schema_version != source_schema_version {
619                if published
620                    .failed_rebuild
621                    .as_ref()
622                    .is_some_and(|(version, until)| {
623                        *version == source_schema_version && *until > Instant::now()
624                    })
625                {
626                    return InvalidQuerySnafu {
627                        reason: format!(
628                            "Flow {flow_id} schema rebuild for source version {source_schema_version} is cooling down"
629                        ),
630                    }
631                    .fail();
632                }
633                #[cfg(test)]
634                slot.rebuild_attempts.fetch_add(1, Ordering::Relaxed);
635                let replacement = self
636                    .build_stateless_flow(&current.create_args, false)
637                    .await
638                    .and_then(|replacement| {
639                        ensure!(
640                            replacement.source_table_id == expected_table_id
641                                && replacement.source_schema_version == source_schema_version,
642                            InvalidQuerySnafu {
643                                reason: format!(
644                                    "Source schema changed while rebuilding flow {flow_id}"
645                                )
646                            }
647                        );
648                        Ok(replacement)
649                    });
650                match replacement {
651                    Ok(replacement) => {
652                        published.flow = Some(Arc::new(replacement));
653                        published.failed_rebuild = None;
654                    }
655                    Err(error) => {
656                        // Reuse the default heartbeat retry baseline; this is neither a
657                        // negotiated interval nor a cache-freshness guarantee.
658                        published.failed_rebuild = Some((
659                            source_schema_version,
660                            Instant::now() + BASE_HEARTBEAT_INTERVAL,
661                        ));
662                        return Err(error);
663                    }
664                }
665            }
666            // Do not open a replacement/drop gap between preparation and sink execution.
667            guard = tokio::sync::OwnedRwLockWriteGuard::downgrade(published);
668        }
669        let flow = guard
670            .flow
671            .as_ref()
672            .context(FlowNotFoundSnafu { id: flow_id })?;
673        // This is the last check before planning.  In particular, a request normalized against
674        // an old source schema is never allowed to reach a newly published plan.
675        let latest = self
676            .table_info_source
677            .get_table_info_value(&flow.source_table_id)
678            .await?
679            .context(UnexpectedSnafu {
680                reason: "Source table metadata is missing",
681            })?
682            .table_info
683            .meta
684            .schema
685            .version();
686        ensure!(
687            source_schema_version == latest && flow.source_schema_version == latest,
688            InvalidQuerySnafu {
689                reason: format!("Source schema version changed before flow {flow_id} execution")
690            }
691        );
692        stateless::execute(
693            flow,
694            &rows,
695            batch_datatypes,
696            &self.query_engine,
697            &self.frontend_client,
698            latest,
699        )
700        .await
701    }
702
703    pub async fn remove_flow_inner(&self, flow_id: FlowId) -> Result<(), Error> {
704        // Detach first: this is the tombstone/linearization point. A concurrent create must
705        // install a different slot and can never republish into the removed one.
706        let slot = self
707            .stateless_flows
708            .write()
709            .await
710            .remove(&flow_id)
711            .context(FlowNotFoundSnafu { id: flow_id })?;
712        slot.active.store(false, Ordering::Release);
713        let mut runtime = slot.runtime.write().await;
714        runtime.flow.take();
715        runtime.failed_rebuild = None;
716        Ok(())
717    }
718
719    async fn publish_initial_flow(
720        &self,
721        slot: Arc<StatelessFlowSlot>,
722        flow: Arc<StatelessFlow>,
723        create_if_not_exists: bool,
724        or_replace: bool,
725    ) -> Result<bool, Error> {
726        let mut runtime = slot.runtime.write().await;
727        ensure!(
728            slot.active.load(Ordering::Acquire),
729            FlowNotFoundSnafu {
730                id: flow.create_args.flow_id
731            }
732        );
733        if runtime.flow.is_some() && !or_replace {
734            if create_if_not_exists {
735                return Ok(false);
736            }
737            return crate::error::FlowAlreadyExistSnafu {
738                id: flow.create_args.flow_id,
739            }
740            .fail();
741        }
742        let latest = self
743            .table_info_source
744            .get_table_info_value(&flow.source_table_id)
745            .await?
746            .context(UnexpectedSnafu {
747                reason: "Source table metadata is missing",
748            })?
749            .table_info
750            .meta
751            .schema
752            .version();
753        ensure!(
754            latest == flow.source_schema_version,
755            InvalidQuerySnafu {
756                reason: format!(
757                    "Source schema changed while building flow: built version {}, current version {latest}",
758                    flow.source_schema_version
759                )
760            }
761        );
762        runtime.flow = Some(flow);
763        runtime.failed_rebuild = None;
764        Ok(true)
765    }
766
767    pub async fn create_flow_inner(&self, args: CreateFlowArgs) -> Result<Option<FlowId>, Error> {
768        let flow_id = args.flow_id;
769        let flow = Arc::new(self.build_stateless_flow(&args, true).await?);
770        let slot = {
771            let mut slots = self.stateless_flows.write().await;
772            if let Some(slot) = slots.get(&flow_id) {
773                Arc::clone(slot)
774            } else {
775                let slot = Arc::new(StatelessFlowSlot {
776                    runtime: Arc::new(RwLock::new(StatelessFlowRuntime::default())),
777                    active: AtomicBool::new(true),
778                    #[cfg(test)]
779                    rebuild_attempts: AtomicUsize::new(0),
780                });
781                slots.insert(flow_id, Arc::clone(&slot));
782                slot
783            }
784        };
785        let published = self
786            .publish_initial_flow(
787                slot.clone(),
788                flow,
789                args.create_if_not_exists,
790                args.or_replace,
791            )
792            .await?;
793        ensure!(
794            self.stateless_flows
795                .read()
796                .await
797                .get(&flow_id)
798                .is_some_and(|current| Arc::ptr_eq(current, &slot)),
799            FlowNotFoundSnafu { id: flow_id }
800        );
801        if published {
802            info!("Successfully create flow with id={flow_id}");
803            Ok(Some(flow_id))
804        } else {
805            Ok(None)
806        }
807    }
808
809    async fn build_stateless_flow(
810        &self,
811        args: &CreateFlowArgs,
812        create_sink: bool,
813    ) -> Result<StatelessFlow, Error> {
814        let CreateFlowArgs {
815            flow_id,
816            sink_table_name,
817            source_table_ids,
818            sql,
819            query_ctx,
820            ..
821        } = args;
822        ensure!(
823            source_table_ids.len() == 1,
824            InvalidQuerySnafu {
825                reason: "Stateless streaming flow does not support multiple source tables",
826            }
827        );
828
829        let query_ctx = query_ctx.clone().map(Arc::new).context(UnexpectedSnafu {
830            reason: "Query context is missing",
831        })?;
832        let source_table_id = source_table_ids[0];
833        let source_table_name = self
834            .table_info_source
835            .get_table_name(&source_table_id)
836            .await?;
837        let source_table_info = self
838            .table_info_source
839            .get_table_info_value(&source_table_id)
840            .await?
841            .context(UnexpectedSnafu {
842                reason: "Source table metadata is missing",
843            })?;
844        let flow_plan =
845            sql_to_df_plan(query_ctx.clone(), self.query_engine.clone(), sql, true).await?;
846        stateless::validate_plan(&flow_plan)?;
847        let source_meta = source_table_info.table_info.meta;
848        let source_schema = source_meta.schema;
849        let source_primary_key_indices = source_meta.primary_key_indices;
850        stateless::validate_source_scan(&flow_plan, source_table_id, &source_schema)?;
851        let (inferred_schema, lineage) = output_column_schemas(&flow_plan, &source_schema)?;
852        let inferred_relation =
853            relation_desc_from_output(&inferred_schema, &lineage, &source_primary_key_indices);
854        let sink_exists = self.fetch_table_pk_schema(sink_table_name).await?.is_some();
855        if !sink_exists {
856            validate_auto_column_names(&inferred_schema)?;
857        }
858        if !sink_exists
859            && create_sink
860            && !self
861                .create_table_from_relation(
862                    &format!("flow-id={flow_id}"),
863                    sink_table_name,
864                    &inferred_relation,
865                )
866                .await?
867        {
868            return UnexpectedSnafu {
869                reason: format!("Failed to create table {sink_table_name:?}"),
870            }
871            .fail();
872        }
873        ensure!(
874            sink_exists || create_sink,
875            UnexpectedSnafu {
876                reason: format!("Sink table metadata is missing: {sink_table_name:?}"),
877            }
878        );
879
880        // Fetch the metadata after validation or auto-creation. The sink's layout is the insert
881        // contract: flow output aliases must not leak into the request schema.
882        let (sink_primary_keys, _, sink_schema) = self
883            .fetch_table_pk_schema(sink_table_name)
884            .await?
885            .context(UnexpectedSnafu {
886                reason: format!("Sink table metadata is missing: {sink_table_name:?}"),
887            })?;
888        let (auto_columns, plan) = match resolve_sink_layout(&inferred_schema, &sink_schema) {
889            Ok(auto_columns) => (auto_columns, flow_plan),
890            Err(normal_error) => {
891                // Appending a hidden source timestamp changes DISTINCT's key.  It is therefore
892                // never a valid compatibility rewrite for a DISTINCT flow.
893                let has_distinct = {
894                    let mut found = false;
895                    flow_plan
896                        .apply(|node| {
897                            if matches!(node, LogicalPlan::Distinct(Distinct::All(_))) {
898                                found = true;
899                            }
900                            Ok(datafusion_common::tree_node::TreeNodeRecursion::Continue)
901                        })
902                        .context(DatafusionSnafu {
903                            context: "Failed to inspect flow plan",
904                        })?;
905                    found
906                };
907                if has_distinct
908                    || !sink_exists
909                    || !is_explicit_source_timestamp_compatibility(
910                        &inferred_schema,
911                        &lineage,
912                        &sink_schema,
913                        &source_schema,
914                    )
915                {
916                    return Err(normal_error);
917                }
918                let source_timestamp_index = source_schema.timestamp_index().unwrap();
919                let source_timestamp_name =
920                    &source_schema.column_schemas()[source_timestamp_index].name;
921                let plan = stateless::rewrite_source_timestamp(
922                    flow_plan,
923                    &TableReference::full(
924                        source_table_name[0].clone(),
925                        source_table_name[1].clone(),
926                        source_table_name[2].clone(),
927                    ),
928                    source_timestamp_name,
929                )?;
930                let (effective_output, _) = output_column_schemas(&plan, &source_schema)?;
931                ensure!(
932                    effective_output.len() == sink_schema.len()
933                        && effective_output
934                            .iter()
935                            .zip(&sink_schema)
936                            .all(|(output, sink)| output.data_type == sink.data_type),
937                    InvalidQuerySnafu {
938                        reason: "Compatibility plan output does not match the full sink schema"
939                    }
940                );
941                (vec![], plan)
942            }
943        };
944
945        Ok(StatelessFlow {
946            source_table_id,
947            source_table_name,
948            source_schema_version: source_schema.version(),
949            source_schema,
950            sink_table_name: sink_table_name.clone(),
951            sink_schema,
952            sink_primary_keys,
953            auto_columns,
954            plan,
955            query_ctx,
956            create_args: args.clone(),
957        })
958    }
959
960    pub async fn flush_flow_inner(&self, _flow_id: FlowId) -> Result<usize, Error> {
961        Ok(0)
962    }
963
964    pub(crate) async fn stateless_flow_ids(&self) -> Vec<FlowId> {
965        let slots = self.stateless_flows.read().await;
966        let mut ids = Vec::new();
967        for (id, slot) in slots.iter() {
968            if slot.runtime.read().await.flow.is_some() {
969                ids.push(*id);
970            }
971        }
972        ids
973    }
974
975    pub async fn flow_exist_inner(&self, flow_id: FlowId) -> Result<bool, Error> {
976        let slot = self.stateless_flows.read().await.get(&flow_id).cloned();
977        Ok(match slot {
978            Some(slot) => slot.runtime.read().await.flow.is_some(),
979            None => false,
980        })
981    }
982
983    async fn fetch_table_pk_schema(
984        &self,
985        table_name: &TableName,
986    ) -> Result<Option<(Vec<String>, Option<usize>, Vec<ColumnSchema>)>, Error> {
987        if let Some(table_id) = self
988            .table_info_source
989            .get_opt_table_id_from_name(table_name)
990            .await?
991        {
992            let table_info = self
993                .table_info_source
994                .get_table_info_value(&table_id)
995                .await?
996                .unwrap();
997            let meta = table_info.table_info.meta;
998            let schema = meta.schema.column_schemas().to_vec();
999            let primary_keys = meta
1000                .primary_key_indices
1001                .into_iter()
1002                .map(|i| schema[i].name.clone())
1003                .collect_vec();
1004            Ok(Some((primary_keys, meta.schema.timestamp_index(), schema)))
1005        } else {
1006            Ok(None)
1007        }
1008    }
1009
1010    async fn adjust_auto_created_table_schema(
1011        &self,
1012        schema: &RelationDesc,
1013    ) -> Result<(Vec<String>, Vec<ColumnSchema>, bool), Error> {
1014        let primary_keys = schema
1015            .typ()
1016            .keys
1017            .first()
1018            .map(|key| {
1019                key.column_indices
1020                    .iter()
1021                    .map(|i| {
1022                        schema
1023                            .get_name(*i)
1024                            .clone()
1025                            .unwrap_or_else(|| format!("col_{i}"))
1026                    })
1027                    .collect_vec()
1028            })
1029            .unwrap_or_default();
1030        let mut columns = relation_desc_to_column_schemas_with_fallback(schema);
1031        columns.push(ColumnSchema::new(
1032            AUTO_CREATED_UPDATE_AT_TS_COL,
1033            ConcreteDataType::timestamp_millisecond_datatype(),
1034            true,
1035        ));
1036        let no_time_index = schema.typ().time_index.is_none();
1037        if no_time_index {
1038            columns.push(
1039                ColumnSchema::new(
1040                    AUTO_CREATED_PLACEHOLDER_TS_COL,
1041                    ConcreteDataType::timestamp_millisecond_datatype(),
1042                    true,
1043                )
1044                .with_time_index(true),
1045            );
1046        }
1047        Ok((primary_keys, columns, no_time_index))
1048    }
1049}
1050
1051impl StreamingEngine {
1052    async fn handle_inserts_inner(
1053        &self,
1054        request: api::v1::region::InsertRequests,
1055    ) -> Result<(), Error> {
1056        // A mirrored envelope can contain several regions of one source.  Normalize each request,
1057        // then concatenate by source so DISTINCT is evaluated once over the complete envelope.
1058        let mut grouped: BTreeMap<_, (_, Vec<DiffRow>, Vec<ConcreteDataType>, u32)> =
1059            BTreeMap::new();
1060        let mut poisoned_tables = std::collections::HashSet::new();
1061        let mut first_error = None;
1062        for write_request in request.requests {
1063            let region_id = RegionId::from(write_request.region_id);
1064            let table_id = region_id.table_id();
1065            if poisoned_tables.contains(&table_id) {
1066                continue;
1067            }
1068            match self.handle_insert_request(write_request).await {
1069                Ok((rows, types, version)) => {
1070                    let entry = grouped
1071                        .entry(table_id)
1072                        .or_insert_with(|| (region_id, vec![], types.clone(), version));
1073                    if entry.2 != types || entry.3 != version {
1074                        poisoned_tables.insert(table_id);
1075                        grouped.remove(&table_id);
1076                        let err = InvalidQuerySnafu { reason: format!("Source table {table_id} metadata changed within one insert envelope") }.build();
1077                        let ids: Vec<FlowId> = self
1078                            .flow_ids_for_table(table_id)
1079                            .await
1080                            .into_iter()
1081                            .map(|(id, _)| id)
1082                            .collect();
1083                        let err = InsertIntoFlowSnafu {
1084                            region_id: u64::from(region_id),
1085                            flow_ids: ids,
1086                        }
1087                        .into_error(BoxedError::new(err));
1088                        error!(err; "Failed to normalize flow insert request for region_id={region_id}");
1089                        if first_error.is_none() {
1090                            first_error = Some(err);
1091                        }
1092                    } else {
1093                        entry.1.extend(rows);
1094                    }
1095                }
1096                Err(err) => {
1097                    let ids: Vec<FlowId> = self
1098                        .flow_ids_for_table(table_id)
1099                        .await
1100                        .into_iter()
1101                        .map(|(id, _)| id)
1102                        .collect();
1103                    poisoned_tables.insert(table_id);
1104                    grouped.remove(&table_id);
1105                    let err = InsertIntoFlowSnafu {
1106                        region_id: u64::from(region_id),
1107                        flow_ids: ids,
1108                    }
1109                    .into_error(BoxedError::new(err));
1110                    error!(err; "Failed to normalize flow insert request for region_id={region_id}");
1111                    if first_error.is_none() {
1112                        first_error = Some(err);
1113                    }
1114                }
1115            }
1116        }
1117        for (_, (region_id, rows, types, version)) in grouped {
1118            if let Err(err) = self
1119                .handle_write_request(region_id, rows, &types, version)
1120                .await
1121                && first_error.is_none()
1122            {
1123                first_error = Some(err);
1124            }
1125        }
1126        match first_error {
1127            Some(err) => Err(err),
1128            None => Ok(()),
1129        }
1130    }
1131
1132    async fn handle_insert_request(
1133        &self,
1134        write_request: api::v1::region::InsertRequest,
1135    ) -> Result<(Vec<DiffRow>, Vec<ConcreteDataType>, u32), Error> {
1136        let region_id = write_request.region_id;
1137        let table_id = RegionId::from(region_id).table_id();
1138        let (insert_schema, rows_proto) = write_request
1139            .rows
1140            .map(|r| (r.schema, r.rows))
1141            .unwrap_or_default();
1142        let now = common_time::util::current_time_millis();
1143        let (table_types, fetch_order, source_schema_version) = {
1144            // Fetch the source metadata once for both current-schema normalization and
1145            // retained-plan validation. In particular, do not execute a retained plan
1146            // against rows normalized with a newer schema.
1147            let table_info = self
1148                .table_info_source
1149                .get_table_info_value(&table_id)
1150                .await?
1151                .context(UnexpectedSnafu {
1152                    reason: format!("Table metadata is missing for table id {table_id}"),
1153                })?;
1154            let source_schema_version = table_info.table_info.meta.schema.version();
1155            let table_schema = table_info_value_to_relation_desc(table_info)?;
1156            let defaults = table_schema
1157                .default_values
1158                .iter()
1159                .zip(table_schema.relation_desc.typ().column_types.iter())
1160                .map(|(value, ty)| {
1161                    value.as_ref().and_then(|value| {
1162                        value.create_default(ty.scalar_type(), ty.nullable()).ok()
1163                    })
1164                })
1165                .collect_vec();
1166            let types = table_schema
1167                .relation_desc
1168                .typ()
1169                .column_types
1170                .iter()
1171                .map(|ty| ty.scalar_type.clone())
1172                .collect_vec();
1173            let names = table_schema
1174                .relation_desc
1175                .names
1176                .iter()
1177                .enumerate()
1178                .map(|(idx, name)| {
1179                    name.clone().context(InternalSnafu {
1180                        reason: format!("Column {idx} of table {table_id} has no name"),
1181                    })
1182                })
1183                .collect::<Result<Vec<_>, _>>()?;
1184            let input_columns = insert_schema
1185                .iter()
1186                .enumerate()
1187                .map(|(idx, column)| (&column.column_name, idx))
1188                .collect::<std::collections::HashMap<_, _>>();
1189            let order = names
1190                .iter()
1191                .zip(defaults)
1192                .map(|(name, default)| {
1193                    input_columns
1194                        .get(name)
1195                        .copied()
1196                        .map(FetchFromRow::Idx)
1197                        .or_else(|| default.map(FetchFromRow::Default))
1198                        .with_context(|| UnexpectedSnafu {
1199                            reason: format!("Column not found: {name}"),
1200                        })
1201                })
1202                .collect::<Result<Vec<_>, _>>()?;
1203            (types, order, source_schema_version)
1204        };
1205        let rows = rows_proto
1206            .into_iter()
1207            .map(|row| {
1208                let row = Row::from(row);
1209                let row = fetch_order
1210                    .iter()
1211                    .map(|item| item.fetch(&row))
1212                    .collect_vec();
1213                (Row::new(row), now, 1)
1214            })
1215            .collect_vec();
1216        Ok((rows, table_types, source_schema_version))
1217    }
1218}
1219
1220#[derive(Debug, Clone)]
1221enum FetchFromRow {
1222    Idx(usize),
1223    Default(datatypes::value::Value),
1224}
1225
1226impl FetchFromRow {
1227    fn fetch(&self, row: &Row) -> datatypes::value::Value {
1228        match self {
1229            Self::Idx(idx) => row
1230                .get(*idx)
1231                .cloned()
1232                .unwrap_or(datatypes::value::Value::Null),
1233            Self::Default(value) => value.clone(),
1234        }
1235    }
1236}