Skip to main content

operator/
insert.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;
16use std::sync::atomic::{AtomicU64, Ordering};
17use std::sync::{Arc, LazyLock};
18use std::time::{Duration, Instant};
19
20use ahash::{HashMap, HashMapExt, HashSet, HashSetExt};
21use api::v1::alter_table_expr::Kind;
22use api::v1::column_def::{options_from_skipping, try_as_column_def};
23use api::v1::region::{
24    InsertRequest as RegionInsertRequest, InsertRequests as RegionInsertRequests,
25    RegionRequestHeader,
26};
27use api::v1::value::ValueData;
28use api::v1::{
29    AlterTableExpr, ColumnDataType, ColumnSchema, CreateTableExpr, InsertRequests,
30    RowInsertRequest, RowInsertRequests, Rows, SemanticType,
31};
32use arrow::datatypes::{DataType as ArrowDataType, Schema as ArrowSchema};
33use catalog::CatalogManagerRef;
34use client::{OutputData, OutputMeta};
35use common_catalog::consts::{
36    DEFAULT_PRIVATE_SCHEMA_NAME, PARENT_SPAN_ID_COLUMN, SERVICE_NAME_COLUMN, TRACE_ID_COLUMN,
37    TRACE_TABLE_NAME, TRACE_TABLE_NAME_SESSION_KEY, default_engine, is_ddl_reserved_table,
38    trace_operations_table_name, trace_services_table_name,
39};
40use common_event_recorder::DEFAULT_EVENTS_TABLE_NAME;
41use common_frontend::slow_query_event::SLOW_QUERY_TABLE_NAME;
42use common_grpc_expr::util::ColumnExpr;
43use common_meta::cache::TableFlownodeSetCacheRef;
44use common_meta::datanode::REGION_STATS_HISTORY_TABLE_NAME;
45use common_meta::node_manager::{AffectedRows, NodeManagerRef};
46use common_meta::peer::Peer;
47use common_meta::rpc::ddl::TriggerReason;
48use common_query::Output;
49use common_query::native_histogram::{is_native_histogram_value_type, native_histogram_value_type};
50use common_query::prelude::{greptime_timestamp, greptime_value};
51use common_telemetry::tracing_context::TracingContext;
52use common_telemetry::{debug, error, warn};
53use common_time::Timestamp;
54use common_time::timestamp::TimeUnit;
55use datatypes::schema::SkippingIndexOptions;
56use futures_util::future;
57use meter_core::data::MeterRecord;
58use meter_macros::write_meter;
59use partition::manager::PartitionRuleManagerRef;
60use session::context::QueryContextRef;
61use snafu::ResultExt;
62use snafu::prelude::*;
63use sql::partition::partition_rule_for_hexstring;
64use sql::statements::create::Partitions;
65use sql::statements::insert::Insert;
66use store_api::metric_engine_consts::{
67    LOGICAL_TABLE_METADATA_KEY, METRIC_ENGINE_NAME, PHYSICAL_TABLE_METADATA_KEY,
68};
69use store_api::mito_engine_options::{
70    APPEND_MODE_KEY, COMPACTION_TYPE, COMPACTION_TYPE_TWCS, MERGE_MODE_KEY, TTL_KEY,
71    TWCS_TIME_WINDOW,
72};
73use store_api::storage::{RegionId, TableId};
74use table::TableRef;
75use table::metadata::{TableInfo, TableInfoRef};
76use table::requests::{
77    AUTO_CREATE_TABLE_KEY, InsertRequest as TableInsertRequest, SEMANTIC_PER_TABLE_INDEX_KEY,
78    SEMANTIC_PIPELINE, TABLE_DATA_MODEL, TABLE_DATA_MODEL_TRACE_V1, TABLE_DATA_MODEL_TRACE_V2,
79    TRACE_TABLE_PARTITIONS_HINT_KEY, VALID_TABLE_OPTION_KEYS, is_semantic_option_key,
80    validate_semantic_option,
81};
82use table::table_reference::TableReference;
83
84use crate::batcher::PendingRowsBatcher;
85use crate::error::{
86    CatalogSnafu, ColumnOptionsSnafu, CreatePartitionRulesSnafu, FindRegionLeaderSnafu,
87    InvalidInsertRequestSnafu, JoinTaskSnafu, RequestInsertsSnafu, Result, TableNotFoundSnafu,
88    WriteRejectedSnafu,
89};
90use crate::expr_helper;
91use crate::region_req_factory::RegionRequestFactory;
92use crate::req_convert::common::preprocess_row_insert_requests;
93use crate::req_convert::insert::{
94    ColumnToRow, ImpureDefaultFiller, RowToRegion, StatementToRegion, TableToRegion,
95    fill_reqs_with_impure_default, rows_to_record_batch,
96};
97use crate::statement::StatementExecutor;
98
99pub struct Inserter {
100    catalog_manager: CatalogManagerRef,
101    pub(crate) partition_manager: PartitionRuleManagerRef,
102    pub(crate) node_manager: NodeManagerRef,
103    pub(crate) table_flownode_set_cache: TableFlownodeSetCacheRef,
104    /// Server-side upper bound for auto table creation on write.
105    /// When `false`, missing tables are never auto-created regardless of the
106    /// per-request `auto_create_table` hint. When `true`, the hint still applies.
107    auto_create_table: bool,
108    pending_rows_batcher: Option<Arc<dyn PendingRowsBatcher>>,
109    /// Rows handed to detached flow mirror tasks that have not finished yet.
110    /// Bounds how much row data in-flight mirror writes may keep alive, see
111    /// [`MAX_MIRROR_PENDING_ROWS`].
112    mirror_pending_rows: Arc<AtomicU64>,
113}
114
115pub type InserterRef = Arc<Inserter>;
116
117/// Hint for the table type to create automatically.
118#[derive(Clone)]
119pub enum AutoCreateTableType {
120    /// A logical table with the physical table name.
121    Logical(String),
122    /// A physical table.
123    Physical,
124    /// A log table which is append-only.
125    Log,
126    /// A table that merges rows by `last_non_null` strategy.
127    LastNonNull,
128    /// Create table that build index and default partition rules on trace_id
129    Trace { alter_existing: bool },
130}
131
132impl AutoCreateTableType {
133    pub fn as_str(&self) -> &'static str {
134        match self {
135            AutoCreateTableType::Logical(_) => "logical",
136            AutoCreateTableType::Physical => "physical",
137            AutoCreateTableType::Log => "log",
138            AutoCreateTableType::LastNonNull => "last_non_null",
139            AutoCreateTableType::Trace { .. } => "trace",
140        }
141    }
142
143    fn alter_existing(&self) -> bool {
144        !matches!(
145            self,
146            Self::Trace {
147                alter_existing: false
148            }
149        )
150    }
151}
152
153/// Split insert requests into normal and instant requests.
154///
155/// Where instant requests are requests with ttl=instant,
156/// and normal requests are requests with ttl set to other values.
157///
158/// This is used to split requests for different processing.
159#[derive(Clone)]
160pub struct InstantAndNormalInsertRequests {
161    /// Requests with normal ttl.
162    pub normal_requests: RegionInsertRequests,
163    /// Requests with ttl=instant.
164    /// Will be discarded immediately at frontend, wouldn't even insert into memtable, and only sent to flow node if needed.
165    pub instant_requests: RegionInsertRequests,
166}
167
168impl Inserter {
169    /// Checks the assumptions of the logical bulk path without changing tables.
170    /// Unsupported requests retain ordinary insertion, including schema policy,
171    /// defaults, instant TTL and row-based Flow delivery.
172    pub async fn can_batch_metric_rows(
173        &self,
174        requests: &RowInsertRequests,
175        ctx: &QueryContextRef,
176        physical_table: &str,
177    ) -> Result<bool> {
178        if self.auto_create_disabled_reason(ctx)?.is_some() || ctx.extension(TTL_KEY).is_some() {
179            return Ok(false);
180        }
181        for request in &requests.inserts {
182            // The logical bulk encoder only supports scalar metric schemas;
183            // any time index unit is accepted — requests are converted to the
184            // destination table's unit during batch alignment.
185            // Check new tables too, before catalog lookup or schema changes.
186            if request.rows.as_ref().is_some_and(|rows| {
187                rows.schema.iter().any(|column| {
188                    column.datatype_extension.is_some()
189                        || !matches!(
190                            ColumnDataType::try_from(column.datatype),
191                            Ok(ColumnDataType::TimestampSecond
192                                | ColumnDataType::TimestampMillisecond
193                                | ColumnDataType::TimestampMicrosecond
194                                | ColumnDataType::TimestampNanosecond
195                                | ColumnDataType::Float64
196                                | ColumnDataType::String)
197                        )
198                })
199            }) {
200                return Ok(false);
201            }
202
203            let Some(table) = self
204                .get_table(
205                    ctx.current_catalog(),
206                    &ctx.current_schema(),
207                    &request.table_name,
208                )
209                .await?
210            else {
211                continue;
212            };
213            let info = table.table_info();
214            if info.meta.engine != METRIC_ENGINE_NAME
215                || info.is_ttl_instant_table()
216                || info
217                    .meta
218                    .options
219                    .extra_options
220                    .get(LOGICAL_TABLE_METADATA_KEY)
221                    .map(String::as_str)
222                    != Some(physical_table)
223                || info
224                    .meta
225                    .schema
226                    .column_schemas()
227                    .iter()
228                    .any(|column| column.default_constraint().is_some())
229            {
230                return Ok(false);
231            }
232            // Physical metric tags are nullable even when their logical schema
233            // is not. Arrow alignment is stricter than ordinary metric insertion.
234            if info
235                .meta
236                .primary_key_indices
237                .iter()
238                .any(|&index| !info.meta.schema.column_schemas()[index].is_nullable())
239            {
240                return Ok(false);
241            }
242            // The current Flow cache does not distinguish streaming and batch
243            // flows. Keep all Flow sources on the row-based delivery path.
244            match self.table_flownode_set_cache.get(info.table_id()).await {
245                Ok(None) => {}
246                Ok(Some(flows)) if flows.is_empty() => {}
247                _ => return Ok(false),
248            }
249        }
250        Ok(true)
251    }
252
253    /// Meters an original logical-table request before bulk routing, without
254    /// cloning its rows or changing the request boundary used for accounting.
255    pub async fn meter_row_inserts(
256        requests: &mut RowInsertRequests,
257        ctx: &QueryContextRef,
258    ) -> Result<u64> {
259        let metered = InstantAndNormalInsertRequests {
260            normal_requests: RegionInsertRequests {
261                requests: requests
262                    .inserts
263                    .iter_mut()
264                    .map(|request| RegionInsertRequest {
265                        rows: request.rows.take(),
266                        ..Default::default()
267                    })
268                    .collect(),
269            },
270            instant_requests: RegionInsertRequests::default(),
271        };
272        let cost = write_meter!(
273            ctx.current_catalog(),
274            ctx.current_schema(),
275            metered,
276            ctx.write_rows_to_admit(
277                ctx.current_catalog(),
278                &ctx.current_schema(),
279                count_insert_rows(&metered)?
280            ),
281            ctx.channel() as u8
282        )
283        .await
284        .context(WriteRejectedSnafu);
285        for (request, region) in requests
286            .inserts
287            .iter_mut()
288            .zip(metered.normal_requests.requests)
289        {
290            request.rows = region.rows;
291        }
292        cost
293    }
294
295    pub fn new(
296        catalog_manager: CatalogManagerRef,
297        partition_manager: PartitionRuleManagerRef,
298        node_manager: NodeManagerRef,
299        table_flownode_set_cache: TableFlownodeSetCacheRef,
300        auto_create_table: bool,
301    ) -> Self {
302        Self {
303            catalog_manager,
304            partition_manager,
305            node_manager,
306            table_flownode_set_cache,
307            auto_create_table,
308            pending_rows_batcher: None,
309            mirror_pending_rows: Arc::new(AtomicU64::new(0)),
310        }
311    }
312
313    /// Installs the shared batcher; callers explicitly select its ingestion entry point.
314    pub fn with_pending_rows_batcher(
315        mut self,
316        batcher: Option<Arc<dyn PendingRowsBatcher>>,
317    ) -> Self {
318        self.pending_rows_batcher = batcher;
319        self
320    }
321
322    pub async fn handle_column_inserts(
323        &self,
324        requests: InsertRequests,
325        ctx: QueryContextRef,
326        statement_executor: &StatementExecutor,
327    ) -> Result<Output> {
328        let row_inserts = ColumnToRow::convert(requests)?;
329        self.handle_row_inserts(row_inserts, ctx, statement_executor, false, false)
330            .await
331    }
332
333    /// Handles row inserts request and creates a physical table on demand.
334    pub async fn handle_row_inserts(
335        &self,
336        mut requests: RowInsertRequests,
337        ctx: QueryContextRef,
338        statement_executor: &StatementExecutor,
339        accommodate_existing_schema: bool,
340        is_single_value: bool,
341    ) -> Result<Output> {
342        preprocess_row_insert_requests(&mut requests.inserts)?;
343        self.handle_row_inserts_with_create_type(
344            requests,
345            ctx,
346            statement_executor,
347            AutoCreateTableType::Physical,
348            accommodate_existing_schema,
349            is_single_value,
350        )
351        .await
352    }
353
354    /// Handles row inserts request and creates a log table on demand.
355    pub async fn handle_log_inserts(
356        &self,
357        requests: RowInsertRequests,
358        ctx: QueryContextRef,
359        statement_executor: &StatementExecutor,
360    ) -> Result<Output> {
361        self.handle_row_inserts_with_create_type(
362            requests,
363            ctx,
364            statement_executor,
365            AutoCreateTableType::Log,
366            false,
367            false,
368        )
369        .await
370    }
371
372    pub async fn handle_trace_inserts(
373        &self,
374        requests: RowInsertRequests,
375        ctx: QueryContextRef,
376        statement_executor: &StatementExecutor,
377    ) -> Result<Output> {
378        self.handle_row_inserts_with_create_type(
379            requests,
380            ctx,
381            statement_executor,
382            AutoCreateTableType::Trace {
383                alter_existing: true,
384            },
385            false,
386            false,
387        )
388        .await
389    }
390
391    /// Handles row inserts request and creates a table with `last_non_null` merge mode on demand.
392    pub async fn handle_last_non_null_inserts(
393        &self,
394        requests: RowInsertRequests,
395        ctx: QueryContextRef,
396        statement_executor: &StatementExecutor,
397        accommodate_existing_schema: bool,
398        is_single_value: bool,
399    ) -> Result<Output> {
400        self.handle_row_inserts_with_create_type(
401            requests,
402            ctx,
403            statement_executor,
404            AutoCreateTableType::LastNonNull,
405            accommodate_existing_schema,
406            is_single_value,
407        )
408        .await
409    }
410
411    /// Handles row inserts request with specified [AutoCreateTableType].
412    async fn handle_row_inserts_with_create_type(
413        &self,
414        mut requests: RowInsertRequests,
415        ctx: QueryContextRef,
416        statement_executor: &StatementExecutor,
417        create_type: AutoCreateTableType,
418        accommodate_existing_schema: bool,
419        is_single_value: bool,
420    ) -> Result<Output> {
421        let skip_wal = ctx.skip_wal();
422
423        let batcher = self
424            .pending_rows_batcher
425            .as_ref()
426            .filter(|_| ctx.batching_enabled());
427
428        // remove empty requests
429        requests.inserts.retain(|req| {
430            req.rows
431                .as_ref()
432                .map(|r| !r.rows.is_empty())
433                .unwrap_or_default()
434        });
435        validate_column_count_match(&requests)?;
436
437        let CreateAlterTableResult {
438            instant_table_ids,
439            table_infos,
440        } = self
441            .create_or_alter_tables_on_demand(
442                &mut requests,
443                &ctx,
444                create_type,
445                statement_executor,
446                accommodate_existing_schema,
447                is_single_value,
448                None,
449            )
450            .await?;
451
452        // Instant tables have no persisted data for dirty-window Flow to read.
453        // Metric tables keep their existing dedicated ingestion path.
454        if let Some(batcher) = batcher
455            && instant_table_ids.is_empty()
456            && table_infos
457                .values()
458                .all(|info| info.meta.engine == default_engine())
459        {
460            return self
461                .submit_pending_rows(requests, table_infos, ctx, batcher)
462                .await;
463        }
464
465        let name_to_info = table_infos
466            .values()
467            .map(|info| (info.name.clone(), info.clone()))
468            .collect::<HashMap<_, _>>();
469        let inserts = RowToRegion::new(
470            name_to_info,
471            instant_table_ids,
472            self.partition_manager.as_ref(),
473        )
474        .convert(requests, skip_wal)
475        .await?;
476
477        self.do_request(inserts, &table_infos, &ctx).await
478    }
479
480    async fn submit_pending_rows(
481        &self,
482        mut requests: RowInsertRequests,
483        table_infos: HashMap<TableId, Arc<TableInfo>>,
484        ctx: QueryContextRef,
485        batcher: &Arc<dyn PendingRowsBatcher>,
486    ) -> Result<Output> {
487        // All entry points, including single-table and SQL writes, skip empty input
488        // before evaluating defaults or converting prepared rows.
489        requests.inserts.retain(|request| {
490            request
491                .rows
492                .as_ref()
493                .is_some_and(|rows| !rows.rows.is_empty())
494        });
495        let by_name = table_infos
496            .values()
497            .map(|info| (info.name.as_str(), info))
498            .collect::<HashMap<_, _>>();
499        let mut prepared = Vec::with_capacity(requests.inserts.len());
500        for request in &mut requests.inserts {
501            let table_info =
502                by_name
503                    .get(request.table_name.as_str())
504                    .context(TableNotFoundSnafu {
505                        table_name: &request.table_name,
506                    })?;
507            let Some(rows) = &mut request.rows else {
508                continue;
509            };
510            ImpureDefaultFiller::new((*table_info).clone())?.fill_rows(rows);
511            let batch = rows_to_record_batch(rows, table_info)?;
512            prepared.push(((*table_info).clone(), batch));
513        }
514        // Preserve the existing meter input and original request boundary. These
515        // envelopes are only for accounting; routing happens after batching.
516        let metered = InstantAndNormalInsertRequests {
517            normal_requests: RegionInsertRequests {
518                requests: requests
519                    .inserts
520                    .into_iter()
521                    .map(|request| RegionInsertRequest {
522                        rows: request.rows,
523                        ..Default::default()
524                    })
525                    .collect(),
526            },
527            instant_requests: RegionInsertRequests::default(),
528        };
529        let table_info = table_infos.values().next();
530        let catalog = table_info.map_or(ctx.current_catalog(), |info| info.catalog_name.as_str());
531        let schema =
532            table_info.map_or_else(|| ctx.current_schema(), |info| info.schema_name.clone());
533        let write_cost = write_meter!(
534            catalog,
535            &schema,
536            metered,
537            ctx.write_rows_to_admit(catalog, &schema, count_insert_rows(&metered)?),
538            ctx.channel() as u8
539        )
540        .await
541        .context(WriteRejectedSnafu)?;
542        prepared.retain(|(_, batch)| batch.num_rows() != 0);
543        let results = if prepared.is_empty() {
544            Vec::new()
545        } else {
546            // One original request shares admission across all table submissions.
547            let permit = batcher.acquire().await?;
548            let submissions = prepared.into_iter().map(|(info, batch)| {
549                // Route to the same target database used for admission above.
550                let mut target_ctx = ctx.fork();
551                target_ctx.set_current_catalog(&info.catalog_name);
552                target_ctx.set_current_schema(&info.schema_name);
553                batcher.submit(info, batch, Arc::new(target_ctx), permit.clone())
554            });
555            // Observe every table completion even when another table fails.
556            future::join_all(submissions).await
557        };
558        let mut affected_rows = 0;
559        let mut meta = OutputMeta::new_with_cost(write_cost as _);
560        for result in results {
561            let output = result?;
562            affected_rows += output.extract_rows_and_cost().0;
563            meta.write_completions.extend(output.meta.write_completions);
564        }
565        Ok(Output::new(OutputData::AffectedRows(affected_rows), meta))
566    }
567
568    /// Handles row inserts request with metric engine.
569    pub async fn handle_metric_row_inserts(
570        &self,
571        mut requests: RowInsertRequests,
572        ctx: QueryContextRef,
573        statement_executor: &StatementExecutor,
574        physical_table: String,
575    ) -> Result<Output> {
576        let skip_wal = ctx.skip_wal();
577
578        // remove empty requests
579        requests.inserts.retain(|req| {
580            req.rows
581                .as_ref()
582                .map(|r| !r.rows.is_empty())
583                .unwrap_or_default()
584        });
585        validate_column_count_match(&requests)?;
586
587        // check and create physical table
588        let physical_table_ref = self
589            .create_physical_table_on_demand(&ctx, physical_table.clone(), statement_executor)
590            .await?;
591
592        // check and create logical tables; `create_or_alter_tables_on_demand`
593        // aligns each request's time index unit with the unit of the table it
594        // targets, inside its existing table lookups: existing tables keep
595        // their own unit (which matches the physical table they are bound
596        // to), and new tables use the selected physical table's unit (from
597        // `physical_table_ref`). Ingestion endpoints encode timestamps in a
598        // fixed unit (prometheus remote write always uses millisecond; OTLP
599        // keeps nanosecond precision on the metric engine path), and the
600        // metric engine requires each logical table's requests to match its
601        // time index unit. Narrowing conversions truncate the sub-unit part
602        // (floor), following `Timestamp::convert_to`.
603        let CreateAlterTableResult {
604            instant_table_ids,
605            table_infos,
606        } = self
607            .create_or_alter_tables_on_demand(
608                &mut requests,
609                &ctx,
610                AutoCreateTableType::Logical(physical_table.clone()),
611                statement_executor,
612                true,
613                true,
614                table_time_index_unit(&physical_table_ref),
615            )
616            .await?;
617        let name_to_info = table_infos
618            .values()
619            .map(|info| (info.name.clone(), info.clone()))
620            .collect::<HashMap<_, _>>();
621        let inserts = RowToRegion::new(name_to_info, instant_table_ids, &self.partition_manager)
622            .convert(requests, skip_wal)
623            .await?;
624
625        self.do_request(inserts, &table_infos, &ctx).await
626    }
627
628    fn table_batcher(
629        &self,
630        table_info: &TableInfoRef,
631        ctx: &QueryContextRef,
632    ) -> Option<&Arc<dyn PendingRowsBatcher>> {
633        self.pending_rows_batcher.as_ref().filter(|_| {
634            ctx.batching_enabled()
635                && !table_info.is_ttl_instant_table()
636                && table_info.meta.engine == default_engine()
637        })
638    }
639
640    async fn submit_table_rows(
641        &self,
642        rows: Rows,
643        table_info: TableInfoRef,
644        ctx: QueryContextRef,
645        batcher: &Arc<dyn PendingRowsBatcher>,
646    ) -> Result<Output> {
647        let requests = RowInsertRequests {
648            inserts: vec![RowInsertRequest {
649                table_name: table_info.name.clone(),
650                rows: Some(rows),
651            }],
652        };
653        let table_infos = HashMap::from_iter([(table_info.table_id(), table_info)]);
654        self.submit_pending_rows(requests, table_infos, ctx, batcher)
655            .await
656    }
657
658    pub async fn handle_table_insert(
659        &self,
660        request: TableInsertRequest,
661        ctx: QueryContextRef,
662    ) -> Result<Output> {
663        let catalog = request.catalog_name.as_str();
664        let schema = request.schema_name.as_str();
665        let table_name = request.table_name.as_str();
666        let table = self.get_table(catalog, schema, table_name).await?;
667        let table = table.with_context(|| TableNotFoundSnafu {
668            table_name: common_catalog::format_full_table_name(catalog, schema, table_name),
669        })?;
670        let table_info = table.table_info();
671
672        let converter = TableToRegion::new(&table_info, &self.partition_manager);
673        let skip_wal = request.skip_wal;
674        let rows = converter.prepare(request)?;
675        if let Some(batcher) = self.table_batcher(&table_info, &ctx) {
676            return self.submit_table_rows(rows, table_info, ctx, batcher).await;
677        }
678        let inserts = converter.partition(rows, skip_wal).await?;
679
680        let table_infos = HashMap::from_iter([(table_info.table_id(), table_info.clone())]);
681
682        self.do_request(inserts, &table_infos, &ctx).await
683    }
684
685    pub async fn handle_statement_insert(
686        &self,
687        insert: &Insert,
688        ctx: &QueryContextRef,
689    ) -> Result<Output> {
690        let converter =
691            StatementToRegion::new(self.catalog_manager.as_ref(), &self.partition_manager, ctx);
692        let (rows, table_info) = converter.prepare(insert, ctx).await?;
693        if let Some(batcher) = self.table_batcher(&table_info, ctx) {
694            return self
695                .submit_table_rows(rows, table_info, ctx.clone(), batcher)
696                .await;
697        }
698        let inserts = converter.partition(rows, table_info.clone(), ctx).await?;
699
700        let table_infos = HashMap::from_iter([(table_info.table_id(), table_info.clone())]);
701
702        self.do_request(inserts, &table_infos, ctx).await
703    }
704}
705
706/// Admits a finite request before it is split into internal writes.
707/// The returned context preserves accounting while preventing a second row debit.
708pub async fn admit_write(rows: u64, ctx: &QueryContextRef) -> Result<QueryContextRef> {
709    // The zero value is WCU: this record only admits rows. Actual inserts retain
710    // their existing WCU accounting, so charging here would count it twice.
711    write_meter!(MeterRecord::new(
712        ctx.current_catalog().to_string(),
713        ctx.current_schema(),
714        0,
715        ctx.write_rows_to_admit(ctx.current_catalog(), &ctx.current_schema(), rows),
716        ctx.channel() as u8,
717    ))
718    .await
719    .context(WriteRejectedSnafu)?;
720    Ok(Arc::new(ctx.with_write_admission()))
721}
722
723/// Admits all database totals before dispatching any batch of a finite request.
724/// Each batch keeps its own protocol options and target database.
725pub async fn admit_row_insert_batches(
726    batches: &mut [(QueryContextRef, RowInsertRequests)],
727) -> Result<()> {
728    let mut totals = BTreeMap::<_, (QueryContextRef, u64)>::new();
729    for (ctx, requests) in batches.iter() {
730        let catalog = ctx.current_catalog();
731        let schema = ctx.current_schema();
732        if ctx.write_rows_to_admit(catalog, &schema, 1) == 0 {
733            continue;
734        }
735        let (_, total) = totals
736            .entry((catalog.to_string(), schema.clone()))
737            .or_insert_with(|| (ctx.clone(), 0));
738        for rows in requests.inserts.iter().filter_map(|r| r.rows.as_ref()) {
739            *total =
740                total
741                    .checked_add(rows.rows.len() as u64)
742                    .context(InvalidInsertRequestSnafu {
743                        reason: "Insert row count exceeds u64::MAX",
744                    })?;
745        }
746    }
747    for (ctx, rows) in totals.values() {
748        admit_write(*rows, ctx).await?;
749    }
750    for (ctx, _) in batches {
751        *ctx = Arc::new(ctx.with_write_admission());
752    }
753    Ok(())
754}
755
756fn count_insert_rows(requests: &InstantAndNormalInsertRequests) -> Result<u64> {
757    requests
758        .normal_requests
759        .requests
760        .iter()
761        .chain(&requests.instant_requests.requests)
762        .filter_map(|request| request.rows.as_ref())
763        .try_fold(0u64, |total, rows| {
764            total
765                .checked_add(rows.rows.len() as u64)
766                .context(InvalidInsertRequestSnafu {
767                    reason: "Insert row count exceeds u64::MAX",
768                })
769        })
770}
771
772impl Inserter {
773    async fn do_request(
774        &self,
775        requests: InstantAndNormalInsertRequests,
776        table_infos: &HashMap<TableId, Arc<TableInfo>>,
777        ctx: &QueryContextRef,
778    ) -> Result<Output> {
779        // Fill impure default values in the request
780        let requests = fill_reqs_with_impure_default(table_infos, requests)?;
781
782        // All tables in a batch resolve to the same database. Qualified SQL
783        // inserts may target a different database than the session's current one.
784        let table_info = table_infos.values().next();
785        let catalog = table_info.map_or(ctx.current_catalog(), |info| info.catalog_name.as_str());
786        let schema =
787            table_info.map_or_else(|| ctx.current_schema(), |info| info.schema_name.clone());
788        let write_cost = write_meter!(
789            catalog,
790            schema.clone(),
791            requests,
792            ctx.write_rows_to_admit(catalog, &schema, count_insert_rows(&requests)?),
793            ctx.channel() as u8
794        )
795        .await
796        .context(WriteRejectedSnafu)?;
797        let request_factory = RegionRequestFactory::new(RegionRequestHeader {
798            tracing_context: TracingContext::from_current_span().to_w3c(),
799            dbname: ctx.get_db_string(),
800            ..Default::default()
801        });
802
803        let InstantAndNormalInsertRequests {
804            normal_requests,
805            instant_requests,
806        } = requests;
807
808        // Mirror requests for source table to flownode asynchronously
809        let flow_mirror_task = FlowMirrorTask::new(
810            &self.table_flownode_set_cache,
811            normal_requests
812                .requests
813                .iter()
814                .chain(instant_requests.requests.iter()),
815        )
816        .await?;
817        let has_instant_rows = instant_requests.requests.iter().any(|request| {
818            request
819                .rows
820                .as_ref()
821                .is_some_and(|rows| !rows.rows.is_empty())
822        });
823        flow_mirror_task.detach(
824            self.node_manager.clone(),
825            self.mirror_pending_rows.clone(),
826            has_instant_rows,
827        )?;
828
829        // Write requests to datanode and wait for response
830        let write_tasks = self
831            .group_requests_by_peer(normal_requests)
832            .await?
833            .into_iter()
834            .map(|(peer, inserts)| {
835                let node_manager = self.node_manager.clone();
836                let request = request_factory.build_insert(inserts);
837                common_runtime::spawn_global(async move {
838                    node_manager
839                        .datanode(&peer)
840                        .await
841                        .handle(request)
842                        .await
843                        .context(RequestInsertsSnafu)
844                })
845            });
846        let results = future::try_join_all(write_tasks)
847            .await
848            .context(JoinTaskSnafu)?;
849        let affected_rows = results
850            .into_iter()
851            .map(|resp| resp.map(|r| r.affected_rows))
852            .sum::<Result<AffectedRows>>()?;
853        crate::metrics::DIST_INGEST_ROW_COUNT
854            .with_label_values(&[ctx.get_db_string().as_str()])
855            .inc_by(affected_rows as u64);
856        Ok(Output::new(
857            OutputData::AffectedRows(affected_rows),
858            OutputMeta::new_with_cost(write_cost as _),
859        ))
860    }
861
862    async fn group_requests_by_peer(
863        &self,
864        requests: RegionInsertRequests,
865    ) -> Result<HashMap<Peer, RegionInsertRequests>> {
866        // group by region ids first to reduce repeatedly call `find_region_leader`
867        // TODO(discord9): determine if a addition clone is worth it
868        let mut requests_per_region: HashMap<RegionId, RegionInsertRequests> = HashMap::new();
869        for req in requests.requests {
870            let region_id = RegionId::from_u64(req.region_id);
871            requests_per_region
872                .entry(region_id)
873                .or_default()
874                .requests
875                .push(req);
876        }
877
878        let mut inserts: HashMap<Peer, RegionInsertRequests> = HashMap::new();
879
880        for (region_id, reqs) in requests_per_region {
881            let peer = self
882                .partition_manager
883                .find_region_leader(region_id)
884                .await
885                .context(FindRegionLeaderSnafu)?;
886            inserts
887                .entry(peer)
888                .or_default()
889                .requests
890                .extend(reqs.requests);
891        }
892
893        Ok(inserts)
894    }
895
896    /// Returns `Some(reason)` if the config or request hint disables automatic
897    /// table creation. Exempt private system tables are handled by
898    /// [`Self::is_auto_create_exempt_private_table`].
899    fn auto_create_disabled_reason(&self, ctx: &QueryContextRef) -> Result<Option<&'static str>> {
900        let auto_create_table_hint = ctx
901            .extension(AUTO_CREATE_TABLE_KEY)
902            .map(|v| v.parse::<bool>())
903            .transpose()
904            .map_err(|_| {
905                InvalidInsertRequestSnafu {
906                    reason: "`auto_create_table` hint must be a boolean",
907                }
908                .build()
909            })?
910            .unwrap_or(true);
911        Ok(if !self.auto_create_table {
912            Some("auto-create table is disabled by frontend config")
913        } else if !auto_create_table_hint {
914            Some("`auto_create_table` hint is disabled")
915        } else {
916            None
917        })
918    }
919
920    /// Returns whether a private system table may infer and reconcile its schema
921    /// even when automatic table creation is disabled.
922    fn is_auto_create_exempt_private_table(schema: &str, table: &str) -> bool {
923        schema == DEFAULT_PRIVATE_SCHEMA_NAME
924            && matches!(
925                table,
926                DEFAULT_EVENTS_TABLE_NAME | SLOW_QUERY_TABLE_NAME | REGION_STATS_HISTORY_TABLE_NAME
927            )
928    }
929
930    /// Adds missing columns from a bulk stream's schema and returns the refreshed table.
931    /// Call once when initializing the stream, before writing its first batch.
932    /// Does not infer new nested or dictionary columns.
933    pub async fn ensure_bulk_insert_schema(
934        &self,
935        table: TableRef,
936        request_schema: &ArrowSchema,
937        ctx: &QueryContextRef,
938        statement_executor: &StatementExecutor,
939    ) -> Result<TableRef> {
940        let table_info = table.table_info();
941        if self.auto_create_disabled_reason(ctx)?.is_some()
942            && !Self::is_auto_create_exempt_private_table(&table_info.schema_name, &table_info.name)
943        {
944            return Ok(table);
945        }
946
947        let table_schema = table.schema();
948        let schema = request_schema
949            .fields()
950            .iter()
951            .filter(|field| table_schema.column_schema_by_name(field.name()).is_none())
952            .map(|field| {
953                let data_type = field.data_type();
954                // Dictionary values can reach the same infallible child-type conversion
955                // as nested types, even when Arrow's is_nested() returns false.
956                ensure!(
957                    !data_type.is_nested() && !matches!(data_type, ArrowDataType::Dictionary(..)),
958                    crate::error::NotSupportedSnafu {
959                        feat: format!(
960                            "automatically adding bulk insert column '{}' with type {:?}",
961                            field.name(),
962                            data_type
963                        ),
964                    }
965                );
966                let column = datatypes::schema::ColumnSchema::try_from(field.as_ref())
967                    .context(crate::error::ConvertSchemaSnafu)?;
968                // Arrow fields do not carry primary-key semantics. New columns are
969                // fields, unless explicitly marked as a time index.
970                let column_def =
971                    try_as_column_def(&column, false).context(crate::error::ColumnDataTypeSnafu)?;
972                Ok(ColumnSchema {
973                    column_name: column_def.name,
974                    datatype: column_def.data_type,
975                    semantic_type: column_def.semantic_type,
976                    datatype_extension: column_def.datatype_extension,
977                    options: column_def.options,
978                })
979            })
980            .collect::<Result<Vec<_>>>()?;
981        let mut request = RowInsertRequest {
982            table_name: table_info.name.clone(),
983            rows: Some(Rows {
984                schema,
985                rows: Vec::new(),
986            }),
987        };
988        let Some(alter_expr) =
989            self.get_alter_table_expr_on_demand(&mut request, &table, ctx, false, false, true)?
990        else {
991            return Ok(table);
992        };
993
994        statement_executor
995            .alter_table_inner(alter_expr, ctx.clone(), TriggerReason::AutoAlter)
996            .await?;
997        self.get_table(
998            &table_info.catalog_name,
999            &table_info.schema_name,
1000            &table_info.name,
1001        )
1002        .await?
1003        .with_context(|| TableNotFoundSnafu {
1004            table_name: table_info.full_table_name(),
1005        })
1006    }
1007
1008    /// Ensures a trace table has the request-global schema without requiring a
1009    /// padded data row to drive on-demand creation or alteration. When
1010    /// `alter_existing` is false, a table created after planning is left for the
1011    /// caller to re-plan.
1012    pub async fn ensure_trace_table_on_demand(
1013        &self,
1014        table_name: &str,
1015        request_schema: Vec<ColumnSchema>,
1016        alter_existing: bool,
1017        ctx: &QueryContextRef,
1018        statement_executor: &StatementExecutor,
1019    ) -> Result<()> {
1020        let mut requests = RowInsertRequests {
1021            inserts: vec![RowInsertRequest {
1022                table_name: table_name.to_string(),
1023                rows: Some(api::v1::Rows {
1024                    schema: request_schema,
1025                    rows: Vec::new(),
1026                }),
1027            }],
1028        };
1029        self.create_or_alter_tables_on_demand(
1030            &mut requests,
1031            ctx,
1032            AutoCreateTableType::Trace { alter_existing },
1033            statement_executor,
1034            false,
1035            false,
1036            None,
1037        )
1038        .await?;
1039        Ok(())
1040    }
1041
1042    /// Creates or alter tables on demand:
1043    /// - if table does not exist, create table by inferred CreateExpr
1044    /// - if table exist, check if schema matches. If any new column found, alter table by inferred `AlterExpr`
1045    ///
1046    /// Returns a mapping from table name to table id, where table name is the table name involved in the requests.
1047    /// This mapping is used in the conversion of RowToRegion.
1048    ///
1049    /// `accommodate_existing_schema` is used to determine if the existing schema should override the new schema.
1050    /// It only works for TIME_INDEX and single VALUE columns. This is for the case where the user creates a table with
1051    /// custom schema, and then inserts data with endpoints that have default schema setting, like prometheus
1052    /// remote write. This will modify the `RowInsertRequests` in place.
1053    /// `is_single_value` indicates whether the default schema only contains single value column so we can accommodate it.
1054    ///
1055    /// `align_time_index_unit` is the selected physical metric table's time
1056    /// index unit; passing `Some` (metric engine path only) rewrites each
1057    /// request's time index column, inside this function's existing table
1058    /// lookups (no extra catalog access): existing destination tables are
1059    /// converted to their own unit, and new tables to the given unit.
1060    #[allow(clippy::too_many_arguments)]
1061    async fn create_or_alter_tables_on_demand(
1062        &self,
1063        requests: &mut RowInsertRequests,
1064        ctx: &QueryContextRef,
1065        auto_create_table_type: AutoCreateTableType,
1066        statement_executor: &StatementExecutor,
1067        accommodate_existing_schema: bool,
1068        is_single_value: bool,
1069        align_time_index_unit: Option<TimeUnit>,
1070    ) -> Result<CreateAlterTableResult> {
1071        let _timer = crate::metrics::CREATE_ALTER_ON_DEMAND
1072            .with_label_values(&[auto_create_table_type.as_str()])
1073            .start_timer();
1074        let catalog = ctx.current_catalog();
1075        let schema = ctx.current_schema();
1076
1077        let auto_create_disabled_reason = self.auto_create_disabled_reason(ctx)?;
1078        // Enabled batches permit every table, so only disabled batches need a whitelist scan.
1079        let has_auto_create_exempt_table = auto_create_disabled_reason.is_some()
1080            && requests
1081                .inserts
1082                .iter()
1083                .any(|req| Self::is_auto_create_exempt_private_table(&schema, &req.table_name));
1084        let mut table_infos = HashMap::new();
1085        // Without exempt tables, verify existing tables and reject missing ones without inferring schemas.
1086        if let Some(disabled_reason) = auto_create_disabled_reason
1087            && !has_auto_create_exempt_table
1088        {
1089            let mut instant_table_ids = HashSet::new();
1090            for req in &mut requests.inserts {
1091                let table = match self.get_table(catalog, &schema, &req.table_name).await? {
1092                    Some(table) => table,
1093                    // System-defined table: created canonically by the system,
1094                    // so the auto-create config/hint does not apply.
1095                    None if is_ddl_reserved_table(&schema, &req.table_name) => {
1096                        statement_executor
1097                            .create_declared_relationships_table(catalog, ctx.clone())
1098                            .await?
1099                    }
1100                    None => {
1101                        return InvalidInsertRequestSnafu {
1102                            reason: format!(
1103                                "Table `{}` does not exist, and {}",
1104                                req.table_name, disabled_reason
1105                            ),
1106                        }
1107                        .fail();
1108                    }
1109                };
1110                // Metric path: an existing destination table keeps its own
1111                // time index unit (it may be bound to another physical
1112                // table than the one selected by this request).
1113                if align_time_index_unit.is_some()
1114                    && let Some(rows) = req.rows.as_mut()
1115                    && let Some(target_unit) = table_time_index_unit(&table)
1116                {
1117                    convert_rows_time_unit(rows, target_unit)?;
1118                }
1119                let table_info = table.table_info();
1120                if matches!(auto_create_table_type, AutoCreateTableType::Trace { .. }) {
1121                    validate_trace_table_model(&table_info, ctx)?;
1122                }
1123                if table_info.is_ttl_instant_table() {
1124                    instant_table_ids.insert(table_info.table_id());
1125                }
1126                table_infos.insert(table_info.table_id(), table.table_info());
1127            }
1128            let ret = CreateAlterTableResult {
1129                instant_table_ids,
1130                table_infos,
1131            };
1132            return Ok(ret);
1133        }
1134
1135        let mut create_tables = vec![];
1136        let mut alter_tables = vec![];
1137        let mut need_refresh_table_infos = HashSet::new();
1138        let mut instant_table_ids = HashSet::new();
1139        let mut per_table_semantics: Option<Option<PerTableSemanticIndex>> = None;
1140
1141        for req in &mut requests.inserts {
1142            // Mixed batches need a per-table decision so an exempt table cannot authorize others.
1143            let auto_create_allowed = auto_create_disabled_reason.is_none()
1144                || Self::is_auto_create_exempt_private_table(&schema, &req.table_name);
1145            match self.get_table(catalog, &schema, &req.table_name).await? {
1146                Some(table) => {
1147                    let table_info = table.table_info();
1148                    if matches!(auto_create_table_type, AutoCreateTableType::Trace { .. }) {
1149                        validate_trace_table_model(&table_info, ctx)?;
1150                    }
1151                    if table_info.is_ttl_instant_table() {
1152                        instant_table_ids.insert(table_info.table_id());
1153                    }
1154                    // Metric path: an existing destination table keeps its
1155                    // own time index unit (it may be bound to another
1156                    // physical table than the one selected by this request).
1157                    if align_time_index_unit.is_some()
1158                        && let Some(rows) = req.rows.as_mut()
1159                        && let Some(target_unit) = table_time_index_unit(&table)
1160                    {
1161                        convert_rows_time_unit(rows, target_unit)?;
1162                    }
1163                    if auto_create_allowed
1164                        && let Some(alter_expr) = self.get_alter_table_expr_on_demand(
1165                            req,
1166                            &table,
1167                            ctx,
1168                            accommodate_existing_schema,
1169                            is_single_value,
1170                            auto_create_table_type.alter_existing(),
1171                        )?
1172                    {
1173                        alter_tables.push(alter_expr);
1174                        need_refresh_table_infos.insert((
1175                            catalog.to_string(),
1176                            schema.clone(),
1177                            req.table_name.clone(),
1178                        ));
1179                    } else {
1180                        table_infos.insert(table_info.table_id(), table.table_info());
1181                    }
1182                }
1183                // A DDL-reserved table's definition never derives from the
1184                // write request; the system creates it canonically, below the
1185                // user-DDL guard that rejects the generic create path.
1186                None if is_ddl_reserved_table(&schema, &req.table_name) => {
1187                    let table = statement_executor
1188                        .create_declared_relationships_table(catalog, ctx.clone())
1189                        .await?;
1190                    let table_info = table.table_info();
1191                    if table_info.is_ttl_instant_table() {
1192                        instant_table_ids.insert(table_info.table_id());
1193                    }
1194                    table_infos.insert(table_info.table_id(), table_info);
1195                }
1196                None if !auto_create_allowed
1197                    && let Some(disabled_reason) = auto_create_disabled_reason =>
1198                {
1199                    return InvalidInsertRequestSnafu {
1200                        reason: format!(
1201                            "Table `{}` does not exist, and {}",
1202                            req.table_name, disabled_reason,
1203                        ),
1204                    }
1205                    .fail();
1206                }
1207                None => {
1208                    // Metric path: a new table uses the selected physical
1209                    // table's unit; convert before the create expression is
1210                    // derived from the request schema.
1211                    if let Some(physical_unit) = align_time_index_unit
1212                        && let Some(rows) = req.rows.as_mut()
1213                    {
1214                        convert_rows_time_unit(rows, physical_unit)?;
1215                    }
1216                    let semantic_index = per_table_semantics
1217                        .get_or_insert_with(|| parse_per_table_semantic_index(ctx))
1218                        .as_ref();
1219                    let create_expr = self.get_create_table_expr_on_demand(
1220                        req,
1221                        &auto_create_table_type,
1222                        ctx,
1223                        semantic_index,
1224                    )?;
1225                    create_tables.push(create_expr);
1226                }
1227            }
1228        }
1229
1230        match auto_create_table_type {
1231            AutoCreateTableType::Logical(_) => {
1232                if !create_tables.is_empty() {
1233                    // Creates logical tables in batch.
1234                    let tables = self
1235                        .create_logical_tables(create_tables, ctx, statement_executor)
1236                        .await?;
1237
1238                    for table in tables {
1239                        let table_info = table.table_info();
1240                        if table_info.is_ttl_instant_table() {
1241                            instant_table_ids.insert(table_info.table_id());
1242                        }
1243                        table_infos.insert(table_info.table_id(), table.table_info());
1244                    }
1245                }
1246                if !alter_tables.is_empty() {
1247                    // Alter logical tables in batch.
1248                    statement_executor
1249                        .alter_logical_tables(alter_tables, ctx.clone(), TriggerReason::AutoAlter)
1250                        .await?;
1251                }
1252            }
1253            AutoCreateTableType::Physical
1254            | AutoCreateTableType::Log
1255            | AutoCreateTableType::LastNonNull => {
1256                // note that auto create table shouldn't be ttl instant table
1257                // for it's a very unexpected behavior and should be set by user explicitly
1258                for create_table in create_tables {
1259                    let table = self
1260                        .create_physical_table(create_table, None, ctx, statement_executor)
1261                        .await?;
1262                    let table_info = table.table_info();
1263                    if table_info.is_ttl_instant_table() {
1264                        instant_table_ids.insert(table_info.table_id());
1265                    }
1266                    table_infos.insert(table_info.table_id(), table.table_info());
1267                }
1268                for alter_expr in alter_tables.into_iter() {
1269                    statement_executor
1270                        .alter_table_inner(alter_expr, ctx.clone(), TriggerReason::AutoAlter)
1271                        .await?;
1272                }
1273            }
1274
1275            AutoCreateTableType::Trace { .. } => {
1276                let trace_table_name = ctx
1277                    .extension(TRACE_TABLE_NAME_SESSION_KEY)
1278                    .unwrap_or(TRACE_TABLE_NAME);
1279
1280                let trace_table_partitions = if let Some(trace_table_partitions) =
1281                    ctx.extension(TRACE_TABLE_PARTITIONS_HINT_KEY)
1282                {
1283                    let p = trace_table_partitions.parse::<u32>().map_err(|_| {
1284                        InvalidInsertRequestSnafu {
1285                            reason: format!(
1286                                "Failed to parse trace_table_partitions: {}",
1287                                trace_table_partitions
1288                            ),
1289                        }
1290                        .build()
1291                    })?;
1292                    Some(p)
1293                } else {
1294                    None
1295                };
1296
1297                // note that auto create table shouldn't be ttl instant table
1298                // for it's a very unexpected behavior and should be set by user explicitly
1299                for mut create_table in create_tables {
1300                    if create_table.table_name == trace_services_table_name(trace_table_name)
1301                        || create_table.table_name == trace_operations_table_name(trace_table_name)
1302                    {
1303                        // Disable append mode for auxiliary tables (services/operations) since they require upsert behavior.
1304                        create_table
1305                            .table_options
1306                            .insert(APPEND_MODE_KEY.to_string(), "false".to_string());
1307                        // Remove `ttl` key from table options if it exists
1308                        create_table.table_options.remove(TTL_KEY);
1309
1310                        let table = self
1311                            .create_physical_table(create_table, None, ctx, statement_executor)
1312                            .await?;
1313                        let table_info = table.table_info();
1314                        if table_info.is_ttl_instant_table() {
1315                            instant_table_ids.insert(table_info.table_id());
1316                        }
1317                        table_infos.insert(table_info.table_id(), table.table_info());
1318                    } else {
1319                        // prebuilt partition rules for uuid data: see the function
1320                        // for more information
1321                        let partitions = if matches!(trace_table_partitions, Some(0) | Some(1)) {
1322                            // disable partitions
1323                            None
1324                        } else {
1325                            let p = partition_rule_for_hexstring(
1326                                TRACE_ID_COLUMN,
1327                                trace_table_partitions,
1328                            )
1329                            .context(CreatePartitionRulesSnafu)?;
1330                            Some(p)
1331                        };
1332
1333                        // add skip index to
1334                        // - trace_id: when searching by trace id
1335                        // - parent_span_id: when searching root span
1336                        // - span_name: when searching certain types of span
1337                        let index_columns =
1338                            [TRACE_ID_COLUMN, PARENT_SPAN_ID_COLUMN, SERVICE_NAME_COLUMN];
1339                        for index_column in index_columns {
1340                            if let Some(col) = create_table
1341                                .column_defs
1342                                .iter_mut()
1343                                .find(|c| c.name == index_column)
1344                            {
1345                                col.options =
1346                                    options_from_skipping(&SkippingIndexOptions::default())
1347                                        .context(ColumnOptionsSnafu)?;
1348                            } else {
1349                                warn!(
1350                                    "Column {} not found when creating index for trace table: {}.",
1351                                    index_column, create_table.table_name
1352                                );
1353                            }
1354                        }
1355
1356                        // use table_options to mark table model version
1357                        create_table.table_options.insert(
1358                            TABLE_DATA_MODEL.to_string(),
1359                            ctx.extension(SEMANTIC_PIPELINE)
1360                                .unwrap_or(TABLE_DATA_MODEL_TRACE_V1)
1361                                .to_string(),
1362                        );
1363
1364                        let table = self
1365                            .create_physical_table(
1366                                create_table,
1367                                partitions,
1368                                ctx,
1369                                statement_executor,
1370                            )
1371                            .await?;
1372                        let table_info = table.table_info();
1373                        if table_info.is_ttl_instant_table() {
1374                            instant_table_ids.insert(table_info.table_id());
1375                        }
1376                        table_infos.insert(table_info.table_id(), table.table_info());
1377                    }
1378                }
1379                for alter_expr in alter_tables.into_iter() {
1380                    statement_executor
1381                        .alter_table_inner(alter_expr, ctx.clone(), TriggerReason::AutoAlter)
1382                        .await?;
1383                }
1384            }
1385        }
1386
1387        // refresh table infos for altered tables
1388        for (catalog, schema, table_name) in need_refresh_table_infos {
1389            let table = self
1390                .get_table(&catalog, &schema, &table_name)
1391                .await?
1392                .context(TableNotFoundSnafu {
1393                    table_name: common_catalog::format_full_table_name(
1394                        &catalog,
1395                        &schema,
1396                        &table_name,
1397                    ),
1398                })?;
1399            let table_info = table.table_info();
1400            table_infos.insert(table_info.table_id(), table.table_info());
1401        }
1402
1403        Ok(CreateAlterTableResult {
1404            instant_table_ids,
1405            table_infos,
1406        })
1407    }
1408
1409    async fn create_physical_table_on_demand(
1410        &self,
1411        ctx: &QueryContextRef,
1412        physical_table: String,
1413        statement_executor: &StatementExecutor,
1414    ) -> Result<TableRef> {
1415        let catalog_name = ctx.current_catalog();
1416        let schema_name = ctx.current_schema();
1417
1418        // check if exist
1419        if let Some(table) = self
1420            .get_table(catalog_name, &schema_name, &physical_table)
1421            .await?
1422        {
1423            return Ok(table);
1424        }
1425
1426        // Gate here too, otherwise a disabled switch would still leak the physical table.
1427        if let Some(disabled_reason) = self.auto_create_disabled_reason(ctx)? {
1428            return InvalidInsertRequestSnafu {
1429                reason: format!(
1430                    "Physical table `{physical_table}` does not exist, and {disabled_reason}"
1431                ),
1432            }
1433            .fail();
1434        }
1435
1436        let table_reference = TableReference::full(catalog_name, &schema_name, &physical_table);
1437        debug!("Ensuring physical metric table `{table_reference}` exists for insert");
1438
1439        // schema with timestamp and field column
1440        let default_schema = vec![
1441            ColumnSchema {
1442                column_name: greptime_timestamp().to_string(),
1443                datatype: ColumnDataType::TimestampMillisecond as _,
1444                semantic_type: SemanticType::Timestamp as _,
1445                datatype_extension: None,
1446                options: None,
1447            },
1448            ColumnSchema {
1449                column_name: greptime_value().to_string(),
1450                datatype: ColumnDataType::Float64 as _,
1451                semantic_type: SemanticType::Field as _,
1452                datatype_extension: None,
1453                options: None,
1454            },
1455        ];
1456        let create_table_expr =
1457            &mut build_create_table_expr(&table_reference, &default_schema, default_engine())?;
1458
1459        create_table_expr.engine = METRIC_ENGINE_NAME.to_string();
1460        create_table_expr
1461            .table_options
1462            .insert(PHYSICAL_TABLE_METADATA_KEY.to_string(), "true".to_string());
1463
1464        // create physical table
1465        let res = statement_executor
1466            .create_table_inner(
1467                create_table_expr,
1468                None,
1469                ctx.clone(),
1470                TriggerReason::AutoCreate,
1471            )
1472            .await;
1473
1474        match res {
1475            Ok(table) => Ok(table),
1476            Err(err) => {
1477                error!(err; "Failed to create table {table_reference}");
1478                Err(err)
1479            }
1480        }
1481    }
1482
1483    async fn get_table(
1484        &self,
1485        catalog: &str,
1486        schema: &str,
1487        table: &str,
1488    ) -> Result<Option<TableRef>> {
1489        self.catalog_manager
1490            .table(catalog, schema, table, None)
1491            .await
1492            .context(CatalogSnafu)
1493    }
1494
1495    fn get_create_table_expr_on_demand(
1496        &self,
1497        req: &RowInsertRequest,
1498        create_type: &AutoCreateTableType,
1499        ctx: &QueryContextRef,
1500        semantic_index: Option<&PerTableSemanticIndex>,
1501    ) -> Result<CreateTableExpr> {
1502        let schema = ctx.current_schema();
1503        let mut table_options = std::collections::HashMap::with_capacity(4);
1504        fill_table_options_for_create(&mut table_options, create_type, ctx);
1505        apply_per_table_semantic_options(
1506            &mut table_options,
1507            semantic_index,
1508            ctx.current_schema().as_str(),
1509            &req.table_name,
1510        );
1511
1512        let engine_name = if let AutoCreateTableType::Logical(_) = create_type {
1513            // engine should be metric engine when creating logical tables.
1514            METRIC_ENGINE_NAME
1515        } else {
1516            default_engine()
1517        };
1518
1519        let table_ref = TableReference::full(ctx.current_catalog(), &schema, &req.table_name);
1520        // SAFETY: `req.rows` is guaranteed to be `Some` by `handle_row_inserts_with_create_type()`.
1521        let request_schema = req.rows.as_ref().unwrap().schema.as_slice();
1522        let mut create_table_expr =
1523            build_create_table_expr(&table_ref, request_schema, engine_name)?;
1524
1525        // extension set by the Splunk HEC handler for identity path
1526        if ctx.extension(SPLUNK_PK_METADATA_ORDER_KEY).is_some() {
1527            reorder_splunk_primary_keys(&mut create_table_expr.primary_keys);
1528        }
1529
1530        debug!("Ensuring table `{table_ref}` exists for insert");
1531        create_table_expr.table_options.extend(table_options);
1532        Ok(create_table_expr)
1533    }
1534
1535    /// Returns an alter table expression if it finds new columns in the request.
1536    /// When `accommodate_existing_schema` is false, it always adds columns if not exist.
1537    /// When `accommodate_existing_schema` is true, it may modify the input `req` to
1538    /// accommodate it with existing schema. See [`create_or_alter_tables_on_demand`](Self::create_or_alter_tables_on_demand)
1539    /// for more details.
1540    /// When `is_single_value` is true, it also rejects native-histogram/float kind changes.
1541    /// When both options are true, it considers fields when modifying the input `req`.
1542    fn get_alter_table_expr_on_demand(
1543        &self,
1544        req: &mut RowInsertRequest,
1545        table: &TableRef,
1546        ctx: &QueryContextRef,
1547        accommodate_existing_schema: bool,
1548        is_single_value: bool,
1549        alter_existing: bool,
1550    ) -> Result<Option<AlterTableExpr>> {
1551        if !alter_existing {
1552            return Ok(None);
1553        }
1554
1555        let catalog_name = ctx.current_catalog();
1556        let schema_name = ctx.current_schema();
1557        let table_name = table.table_info().name.clone();
1558
1559        // Never auto-alter a system-defined table to fit a write; a request
1560        // with unknown columns fails instead.
1561        if is_ddl_reserved_table(&schema_name, &table_name) {
1562            return Ok(None);
1563        }
1564
1565        let request_schema = req.rows.as_ref().unwrap().schema.as_slice();
1566        let request_field_count = request_schema
1567            .iter()
1568            .filter(|col| col.semantic_type == SemanticType::Field as i32)
1569            .count();
1570        let column_exprs = ColumnExpr::from_column_schemas(request_schema);
1571        let add_columns = expr_helper::extract_add_columns_expr(&table.schema(), column_exprs)?;
1572        let Some(mut add_columns) = add_columns else {
1573            return Ok(None);
1574        };
1575
1576        if is_single_value {
1577            let request_is_native_histogram = request_is_native_histogram(request_schema);
1578            let table_is_native_histogram = table_is_native_histogram(table);
1579            ensure!(
1580                request_is_native_histogram == table_is_native_histogram,
1581                InvalidInsertRequestSnafu {
1582                    reason: format!(
1583                        "Table `{table_name}` cannot mix native histogram and float sample fields"
1584                    ),
1585                }
1586            );
1587        }
1588
1589        // If accommodate_existing_schema is true, update request schema for Timestamp/Field columns
1590        if accommodate_existing_schema {
1591            let table_schema = table.schema();
1592            // Find timestamp column name
1593            let ts_col_name = table_schema.timestamp_column().map(|c| c.name.clone());
1594            // Find field column name if there is only one and `is_single_value` is true.
1595            let mut field_col_name = None;
1596            if is_single_value && request_field_count <= 1 {
1597                let mut multiple_field_cols = false;
1598                table.field_columns().for_each(|col| {
1599                    if field_col_name.is_none() {
1600                        field_col_name = Some(col.name.clone());
1601                    } else {
1602                        multiple_field_cols = true;
1603                    }
1604                });
1605                if multiple_field_cols {
1606                    field_col_name = None;
1607                }
1608            }
1609
1610            // Update column name in request schema for Timestamp/Field columns
1611            if let Some(rows) = req.rows.as_mut() {
1612                for col in &mut rows.schema {
1613                    match col.semantic_type {
1614                        x if x == SemanticType::Timestamp as i32 => {
1615                            if let Some(ref ts_name) = ts_col_name
1616                                && col.column_name != *ts_name
1617                            {
1618                                col.column_name = ts_name.clone();
1619                            }
1620                        }
1621                        x if x == SemanticType::Field as i32 => {
1622                            if let Some(ref field_name) = field_col_name
1623                                && col.column_name != *field_name
1624                            {
1625                                col.column_name = field_name.clone();
1626                            }
1627                        }
1628                        _ => {}
1629                    }
1630                }
1631            }
1632
1633            // Only keep columns that are tags or non-single field.
1634            add_columns.add_columns.retain(|col| {
1635                let def = col.column_def.as_ref().unwrap();
1636                def.semantic_type == SemanticType::Tag as i32
1637                    || (def.semantic_type == SemanticType::Field as i32 && field_col_name.is_none())
1638            });
1639
1640            if add_columns.add_columns.is_empty() {
1641                return Ok(None);
1642            }
1643        }
1644
1645        Ok(Some(AlterTableExpr {
1646            catalog_name: catalog_name.to_string(),
1647            schema_name: schema_name.clone(),
1648            table_name: table_name.clone(),
1649            kind: Some(Kind::AddColumns(add_columns)),
1650        }))
1651    }
1652
1653    /// Creates a table with options.
1654    async fn create_physical_table(
1655        &self,
1656        mut create_table_expr: CreateTableExpr,
1657        partitions: Option<Partitions>,
1658        ctx: &QueryContextRef,
1659        statement_executor: &StatementExecutor,
1660    ) -> Result<TableRef> {
1661        let res = statement_executor
1662            .create_table_inner(
1663                &mut create_table_expr,
1664                partitions,
1665                ctx.clone(),
1666                TriggerReason::AutoCreate,
1667            )
1668            .await;
1669
1670        let table_ref = TableReference::full(
1671            &create_table_expr.catalog_name,
1672            &create_table_expr.schema_name,
1673            &create_table_expr.table_name,
1674        );
1675
1676        match res {
1677            Ok(table) => {
1678                validate_trace_table_model(&table.table_info(), ctx)?;
1679                Ok(table)
1680            }
1681            Err(err) => {
1682                error!(err; "Failed to create table {}", table_ref);
1683                Err(err)
1684            }
1685        }
1686    }
1687
1688    async fn create_logical_tables(
1689        &self,
1690        create_table_exprs: Vec<CreateTableExpr>,
1691        ctx: &QueryContextRef,
1692        statement_executor: &StatementExecutor,
1693    ) -> Result<Vec<TableRef>> {
1694        let res = statement_executor
1695            .create_logical_tables(&create_table_exprs, ctx.clone(), TriggerReason::AutoCreate)
1696            .await;
1697
1698        match res {
1699            Ok(res) => Ok(res),
1700            Err(err) => {
1701                let failed_tables = create_table_exprs
1702                    .into_iter()
1703                    .map(|expr| {
1704                        format!(
1705                            "{}.{}.{}",
1706                            expr.catalog_name, expr.schema_name, expr.table_name
1707                        )
1708                    })
1709                    .collect::<Vec<_>>();
1710                error!(
1711                    err;
1712                    "Failed to create logical tables {:?}",
1713                    failed_tables
1714                );
1715                Err(err)
1716            }
1717        }
1718    }
1719
1720    pub fn node_manager(&self) -> &NodeManagerRef {
1721        &self.node_manager
1722    }
1723
1724    pub fn partition_manager(&self) -> &PartitionRuleManagerRef {
1725        &self.partition_manager
1726    }
1727
1728    pub fn table_flownode_set_cache(&self) -> &TableFlownodeSetCacheRef {
1729        &self.table_flownode_set_cache
1730    }
1731}
1732
1733fn request_is_native_histogram(request_schema: &[ColumnSchema]) -> bool {
1734    let mut fields = request_schema
1735        .iter()
1736        .filter(|col| col.semantic_type == SemanticType::Field as i32);
1737    let Some(col) = fields.next() else {
1738        return false;
1739    };
1740
1741    fields.next().is_none()
1742        && api::helper::is_column_type_value_eq(
1743            col.datatype,
1744            col.datatype_extension.clone(),
1745            native_histogram_value_type(),
1746        )
1747}
1748
1749/// Returns the table's time index unit, if any. A metric table without a
1750/// timestamp column is left alone by the unit alignment: the metric engine
1751/// rejects it anyway.
1752fn table_time_index_unit(table: &TableRef) -> Option<TimeUnit> {
1753    table
1754        .table_info()
1755        .meta
1756        .schema
1757        .timestamp_column()
1758        .and_then(|col| col.data_type.as_timestamp().map(|ts| ts.unit()))
1759}
1760
1761fn convert_rows_time_unit(rows: &mut Rows, target_unit: TimeUnit) -> Result<()> {
1762    let Some(ts_index) = rows
1763        .schema
1764        .iter()
1765        .position(|col| col.semantic_type == SemanticType::Timestamp as i32)
1766    else {
1767        return Ok(());
1768    };
1769    let Some(source_unit) = ColumnDataType::try_from(rows.schema[ts_index].datatype)
1770        .ok()
1771        .and_then(api::helper::timestamp_unit)
1772    else {
1773        return Ok(());
1774    };
1775    if source_unit == target_unit {
1776        return Ok(());
1777    }
1778
1779    rows.schema[ts_index].datatype = api::helper::timestamp_datatype(target_unit) as i32;
1780    // Timestamp columns never carry a datatype extension.
1781    rows.schema[ts_index].datatype_extension = None;
1782
1783    // Note: the schema is rewritten before the rows are converted, so an
1784    // overflow error mid-batch leaves this request half-converted. That is
1785    // harmless: the error aborts the whole insert request.
1786    //
1787    // `validate_column_count_match` guarantees every row carries exactly one
1788    // value per schema column, so the time index position is directly in
1789    // bounds; no per-value search is needed.
1790    for row in &mut rows.rows {
1791        debug_assert_eq!(row.values.len(), rows.schema.len());
1792        let value = &mut row.values[ts_index];
1793        let Some(value_data) = value.value_data.take() else {
1794            continue;
1795        };
1796        value.value_data =
1797            convert_timestamp_value_data(value_data, source_unit, target_unit, ts_index)?;
1798    }
1799    Ok(())
1800}
1801
1802fn convert_timestamp_value_data(
1803    value_data: ValueData,
1804    source_unit: TimeUnit,
1805    target_unit: TimeUnit,
1806    column_index: usize,
1807) -> Result<Option<ValueData>> {
1808    let timestamp = match value_data {
1809        ValueData::TimestampSecondValue(v) => Timestamp::new_second(v),
1810        ValueData::TimestampMillisecondValue(v) => Timestamp::new_millisecond(v),
1811        ValueData::TimestampMicrosecondValue(v) => Timestamp::new_microsecond(v),
1812        ValueData::TimestampNanosecondValue(v) => Timestamp::new_nanosecond(v),
1813        // Null or non-timestamp value; nothing to convert.
1814        other => return Ok(Some(other)),
1815    };
1816    let converted = timestamp
1817        .convert_to(target_unit)
1818        .with_context(|| InvalidInsertRequestSnafu {
1819            reason: format!(
1820                "timestamp column {column_index} value {} in unit {source_unit:?} overflows when converting to unit {target_unit:?}",
1821                timestamp.value()
1822            ),
1823        })?;
1824    Ok(api::helper::to_grpc_value(datatypes::value::Value::Timestamp(converted)).value_data)
1825}
1826
1827fn table_is_native_histogram(table: &TableRef) -> bool {
1828    let mut fields = table.field_columns();
1829    let Some(col) = fields.next() else {
1830        return false;
1831    };
1832
1833    fields.next().is_none() && is_native_histogram_value_type(&col.data_type)
1834}
1835
1836fn validate_column_count_match(requests: &RowInsertRequests) -> Result<()> {
1837    for request in &requests.inserts {
1838        let rows = request.rows.as_ref().unwrap();
1839        let column_count = rows.schema.len();
1840        rows.rows.iter().try_for_each(|r| {
1841            ensure!(
1842                r.values.len() == column_count,
1843                InvalidInsertRequestSnafu {
1844                    reason: format!(
1845                        "column count mismatch, columns: {}, values: {}",
1846                        column_count,
1847                        r.values.len()
1848                    )
1849                }
1850            );
1851            Ok(())
1852        })?;
1853    }
1854    Ok(())
1855}
1856
1857/// Rejects writes from a different built-in trace model before schema mutation.
1858/// Unstamped, explicitly created tables remain subject to normal schema validation.
1859pub fn validate_trace_table_model(table_info: &TableInfo, ctx: &QueryContextRef) -> Result<()> {
1860    let Some(expected @ (TABLE_DATA_MODEL_TRACE_V1 | TABLE_DATA_MODEL_TRACE_V2)) =
1861        ctx.extension(SEMANTIC_PIPELINE)
1862    else {
1863        return Ok(());
1864    };
1865    if let Some(actual) = table_info.meta.options.data_model() {
1866        ensure!(
1867            actual == expected,
1868            InvalidInsertRequestSnafu {
1869                reason: format!(
1870                    "Trace table `{}` uses {actual}, but the request uses {expected}",
1871                    table_info.name,
1872                ),
1873            }
1874        );
1875    }
1876    Ok(())
1877}
1878
1879/// Fill table options for a new table by create type.
1880pub fn fill_table_options_for_create(
1881    table_options: &mut std::collections::HashMap<String, String>,
1882    create_type: &AutoCreateTableType,
1883    ctx: &QueryContextRef,
1884) {
1885    for key in VALID_TABLE_OPTION_KEYS {
1886        if let Some(value) = ctx.extension(key) {
1887            table_options.insert(key.to_string(), value.to_string());
1888        }
1889    }
1890
1891    // Semantic keys use their own vocabulary instead of the fixed option list.
1892    for (key, value) in ctx.extensions() {
1893        if is_semantic_option_key(&key) && validate_semantic_option(&key, &value) {
1894            table_options.insert(key, value);
1895        }
1896    }
1897
1898    match create_type {
1899        AutoCreateTableType::Logical(physical_table) => {
1900            table_options.insert(
1901                LOGICAL_TABLE_METADATA_KEY.to_string(),
1902                physical_table.clone(),
1903            );
1904        }
1905        AutoCreateTableType::Physical => {
1906            if let Some(append_mode) = ctx.extension(APPEND_MODE_KEY) {
1907                table_options.insert(APPEND_MODE_KEY.to_string(), append_mode.to_string());
1908            }
1909            if let Some(merge_mode) = ctx.extension(MERGE_MODE_KEY) {
1910                table_options.insert(MERGE_MODE_KEY.to_string(), merge_mode.to_string());
1911            }
1912            if let Some(time_window) = ctx.extension(TWCS_TIME_WINDOW) {
1913                table_options.insert(TWCS_TIME_WINDOW.to_string(), time_window.to_string());
1914                // We need to set the compaction type explicitly.
1915                table_options.insert(
1916                    COMPACTION_TYPE.to_string(),
1917                    COMPACTION_TYPE_TWCS.to_string(),
1918                );
1919            }
1920        }
1921        // Set append_mode to true for log table.
1922        // because log tables should keep rows with the same ts and tags.
1923        AutoCreateTableType::Log => {
1924            table_options.insert(APPEND_MODE_KEY.to_string(), "true".to_string());
1925        }
1926        AutoCreateTableType::LastNonNull => {
1927            if ctx
1928                .extension(APPEND_MODE_KEY)
1929                .is_some_and(|value| value.eq_ignore_ascii_case("true"))
1930            {
1931                table_options.insert(APPEND_MODE_KEY.to_string(), "true".to_string());
1932                table_options.insert(MERGE_MODE_KEY.to_string(), "last_row".to_string());
1933            } else if let Some(merge_mode) = ctx.extension(MERGE_MODE_KEY) {
1934                table_options.insert(MERGE_MODE_KEY.to_string(), merge_mode.to_string());
1935            } else {
1936                table_options.insert(MERGE_MODE_KEY.to_string(), "last_non_null".to_string());
1937            }
1938        }
1939        AutoCreateTableType::Trace { .. } => {
1940            table_options.insert(APPEND_MODE_KEY.to_string(), "true".to_string());
1941        }
1942    }
1943}
1944
1945/// The parsed per-table semantic index: `{schema -> {table -> {key -> value}}}`,
1946/// produced by the OTLP metrics encode path (where one metric can fan out into
1947/// several tables with distinct keys) and the Prometheus remote write v2 path
1948/// (where per-series metadata declares type/unit, and a series may override its
1949/// target schema).
1950pub type PerTableSemanticIndex = BTreeMap<String, BTreeMap<String, BTreeMap<String, String>>>;
1951
1952/// Parses the per-table semantic index off the context extension. Call once per
1953/// create-planning round: a first write creating N tables would otherwise
1954/// re-parse the whole index N times. `None` when the request carries no index
1955/// (logs, traces, Prom RW v1) or it fails to parse.
1956pub fn parse_per_table_semantic_index(ctx: &QueryContextRef) -> Option<PerTableSemanticIndex> {
1957    let raw = ctx.extension(SEMANTIC_PER_TABLE_INDEX_KEY)?;
1958    match serde_json::from_str(raw) {
1959        Ok(index) => Some(index),
1960        Err(_) => {
1961            warn!("failed to parse semantic per-table index, skipping per-table options");
1962            None
1963        }
1964    }
1965}
1966
1967/// Folds the semantic keys of the table being created into `table_options`.
1968///
1969/// Common keys shared by every table in a request travel as plain semantic
1970/// extensions and are handled by [`fill_table_options_for_create`]; this
1971/// carries only the per-table tail and is applied after it, so a per-table
1972/// value (e.g. `declared` quality) wins. Keys are re-checked against the
1973/// vocabulary defensively.
1974pub fn apply_per_table_semantic_options(
1975    table_options: &mut std::collections::HashMap<String, String>,
1976    index: Option<&PerTableSemanticIndex>,
1977    schema: &str,
1978    table_name: &str,
1979) {
1980    let Some(entry) = index
1981        .and_then(|index| index.get(schema))
1982        .and_then(|tables| tables.get(table_name))
1983    else {
1984        return;
1985    };
1986    for (key, value) in entry {
1987        if is_semantic_option_key(key) && validate_semantic_option(key, value) {
1988            table_options.insert(key.clone(), value.clone());
1989        }
1990    }
1991}
1992
1993pub fn build_create_table_expr(
1994    table: &TableReference,
1995    request_schema: &[ColumnSchema],
1996    engine: &str,
1997) -> Result<CreateTableExpr> {
1998    expr_helper::create_table_expr_by_column_schemas(table, request_schema, engine, None)
1999}
2000
2001/// `QueryContext` extension key the Splunk HEC handler sets (to `"true"`) on its identity
2002/// path to request metadata-first primary-key ordering at table creation. It is absent for
2003/// user-supplied pipelines, so their primary-key order is left untouched.
2004pub const SPLUNK_PK_METADATA_ORDER_KEY: &str = "splunk_pk_metadata_order";
2005
2006/// Moves Splunk's metadata tags (`host`, `source`, `sourcetype`) to the front of the
2007/// primary key, keeping the relative order of the remaining tags.
2008fn reorder_splunk_primary_keys(primary_keys: &mut [String]) {
2009    const LEAD: [&str; 3] = ["host", "source", "sourcetype"];
2010    // Stable sort: `LEAD` columns move to the front in `host`/`source`/`sourcetype` order;
2011    // every other column keeps its existing relative position.
2012    primary_keys.sort_by_key(|name| {
2013        LEAD.iter()
2014            .position(|&lead| lead == name.as_str())
2015            .unwrap_or(LEAD.len())
2016    });
2017}
2018
2019/// Result of `create_or_alter_tables_on_demand`.
2020struct CreateAlterTableResult {
2021    /// table ids of ttl=instant tables.
2022    instant_table_ids: HashSet<TableId>,
2023    /// Table Info of the created tables.
2024    table_infos: HashMap<TableId, Arc<TableInfo>>,
2025}
2026
2027/// Upper bound of rows buffered by detached flow mirror tasks on this frontend.
2028///
2029/// Rows retained by in-flight mirror RPCs count against this bound. If ordinary
2030/// persisted source-table mirroring exceeds it, that mirror batch is best-effort
2031/// dropped while its datanode write continues. Instant-TTL batches fail retryably
2032/// on temporary saturation; oversized batches fail non-retryably so the client
2033/// can reduce the batch size.
2034const MAX_MIRROR_PENDING_ROWS: usize = 1_000_000;
2035
2036/// How often at most each mirror event is reported.
2037const MIRROR_LOG_INTERVAL: Duration = Duration::from_secs(10);
2038const MIRROR_LOG_NEVER_REPORTED: u64 = u64::MAX;
2039
2040struct MirrorLog {
2041    last_log_millis: AtomicU64,
2042    events: AtomicU64,
2043}
2044
2045impl MirrorLog {
2046    const fn new() -> Self {
2047        Self {
2048            last_log_millis: AtomicU64::new(MIRROR_LOG_NEVER_REPORTED),
2049            events: AtomicU64::new(0),
2050        }
2051    }
2052
2053    /// Returns the count since the previous report, including this event, to
2054    /// the single caller that claims the interval.
2055    fn claim_report(&self, now_millis: u64) -> Option<u64> {
2056        self.events.fetch_add(1, Ordering::Relaxed);
2057        let last = self.last_log_millis.load(Ordering::Relaxed);
2058        if last != MIRROR_LOG_NEVER_REPORTED
2059            && now_millis.saturating_sub(last) < MIRROR_LOG_INTERVAL.as_millis() as u64
2060        {
2061            return None;
2062        }
2063        if self
2064            .last_log_millis
2065            .compare_exchange(last, now_millis, Ordering::Relaxed, Ordering::Relaxed)
2066            .is_err()
2067        {
2068            return None;
2069        }
2070        Some(self.events.swap(0, Ordering::Relaxed))
2071    }
2072}
2073
2074static MIRROR_DROP_LOG: MirrorLog = MirrorLog::new();
2075static MIRROR_FAILURE_LOG: MirrorLog = MirrorLog::new();
2076
2077/// Milliseconds since the first call, so the drop log is rate limited on a
2078/// monotonic clock instead of the wall clock, which can jump.
2079fn mirror_drop_log_millis() -> u64 {
2080    static START: LazyLock<Instant> = LazyLock::new(Instant::now);
2081    START.elapsed().as_millis() as u64
2082}
2083
2084struct FlowMirrorTask {
2085    requests: HashMap<Peer, RegionInsertRequests>,
2086}
2087
2088impl FlowMirrorTask {
2089    async fn new(
2090        cache: &TableFlownodeSetCacheRef,
2091        requests: impl Iterator<Item = &RegionInsertRequest>,
2092    ) -> Result<Self> {
2093        let mut src_table_reqs: HashMap<TableId, Option<(Vec<Peer>, RegionInsertRequests)>> =
2094            HashMap::new();
2095
2096        for req in requests {
2097            let table_id = RegionId::from_u64(req.region_id).table_id();
2098            match src_table_reqs.get_mut(&table_id) {
2099                Some(Some((_peers, reqs))) => reqs.requests.push(req.clone()),
2100                // already know this is not source table
2101                Some(None) => continue,
2102                _ => {
2103                    // dedup peers
2104                    let peers = cache
2105                        .get(table_id)
2106                        .await
2107                        .context(RequestInsertsSnafu)?
2108                        .unwrap_or_default()
2109                        .values()
2110                        .cloned()
2111                        .collect::<HashSet<_>>()
2112                        .into_iter()
2113                        .collect::<Vec<_>>();
2114
2115                    if !peers.is_empty() {
2116                        let mut reqs = RegionInsertRequests::default();
2117                        reqs.requests.push(req.clone());
2118                        src_table_reqs.insert(table_id, Some((peers, reqs)));
2119                    } else {
2120                        // insert a empty entry to avoid repeat query
2121                        src_table_reqs.insert(table_id, None);
2122                    }
2123                }
2124            }
2125        }
2126
2127        let mut inserts: HashMap<Peer, RegionInsertRequests> = HashMap::new();
2128
2129        for (_table_id, (peers, reqs)) in src_table_reqs
2130            .into_iter()
2131            .filter_map(|(k, v)| v.map(|v| (k, v)))
2132        {
2133            if peers.len() == 1 {
2134                // fast path, zero copy
2135                inserts
2136                    .entry(peers[0].clone())
2137                    .or_default()
2138                    .requests
2139                    .extend(reqs.requests);
2140                continue;
2141            } else {
2142                // TODO(discord9): need to split requests to multiple flownodes
2143                for flownode in peers {
2144                    inserts
2145                        .entry(flownode.clone())
2146                        .or_default()
2147                        .requests
2148                        .extend(reqs.requests.clone());
2149                }
2150            }
2151        }
2152
2153        Ok(Self { requests: inserts })
2154    }
2155
2156    /// Rows of row data this task keeps alive once detached: the payload it
2157    /// spawns per peer, not the source requests it was built from.
2158    ///
2159    /// A source table mapped to several flownodes clones its requests per peer
2160    /// and every clone owns its own copy of the rows, so reserving from
2161    /// `self.requests` is what the per-peer shares released by [`Self::detach`]
2162    /// sum to. Counting the source requests instead under-reserves whenever a
2163    /// batch holds several requests for the same source table.
2164    fn pending_rows(&self) -> u64 {
2165        self.requests.values().map(region_inserts_rows).sum()
2166    }
2167
2168    fn detach(
2169        self,
2170        node_manager: NodeManagerRef,
2171        mirror_pending_rows: Arc<AtomicU64>,
2172        has_instant_rows: bool,
2173    ) -> Result<()> {
2174        // Reserve the cloned per-peer payload before spawning. A full budget
2175        // drops normal mirrors or rejects instant writes before datanode dispatch.
2176        let num_rows = self.pending_rows();
2177        if num_rows == 0 {
2178            return Ok(());
2179        }
2180        if has_instant_rows && num_rows > MAX_MIRROR_PENDING_ROWS as u64 {
2181            return InvalidInsertRequestSnafu {
2182                reason: format!(
2183                    "flow mirror batch has {num_rows} peer-fan-out rows, exceeding the limit of {}; reduce the batch size",
2184                    MAX_MIRROR_PENDING_ROWS
2185                ),
2186            }
2187            .fail();
2188        }
2189
2190        let pending = reserve_mirror_pending_rows(&mirror_pending_rows, num_rows);
2191        if pending > MAX_MIRROR_PENDING_ROWS as u64 {
2192            release_mirror_pending_rows(&mirror_pending_rows, num_rows);
2193            if has_instant_rows {
2194                return Err(meter_core::collect::WriteRejected::new(
2195                    "flow mirror pending rows limit exceeded for instant-TTL insert",
2196                ))
2197                .context(WriteRejectedSnafu);
2198            }
2199            crate::metrics::DIST_MIRROR_DROPPED_ROW_COUNT.inc_by(num_rows);
2200            if let Some(events) = MIRROR_DROP_LOG.claim_report(mirror_drop_log_millis()) {
2201                warn!(
2202                    "Flow mirror write dropped: {} pending rows exceeds limit {}, {} drop events since last report",
2203                    pending, MAX_MIRROR_PENDING_ROWS, events
2204                );
2205            }
2206            return Ok(());
2207        }
2208
2209        for (peer, inserts) in self.requests {
2210            // Each spawned task releases exactly its own share of the
2211            // reservation. The shares are computed over the same per-peer
2212            // payload as `pending_rows`, so they sum to the reservation.
2213            let peer_rows = region_inserts_rows(&inserts);
2214            let node_manager = node_manager.clone();
2215            let mirror_pending_rows = mirror_pending_rows.clone();
2216            common_runtime::spawn_global(async move {
2217                let result = node_manager
2218                    .flownode(&peer)
2219                    .await
2220                    .handle_inserts(inserts)
2221                    .await
2222                    .context(RequestInsertsSnafu);
2223
2224                match result {
2225                    Ok(resp) => {
2226                        let affected_rows = resp.affected_rows;
2227                        crate::metrics::DIST_MIRROR_ROW_COUNT.inc_by(affected_rows);
2228                    }
2229                    Err(err) => {
2230                        if let Some(events) =
2231                            MIRROR_FAILURE_LOG.claim_report(mirror_drop_log_millis())
2232                        {
2233                            error!(err; flownode_id = peer.id, flownode_addr = %peer.addr, "Failed to insert data into flownode ({} total mirror failures across peers since last report)", events);
2234                        }
2235                    }
2236                }
2237                // Release the reservation on both the success and the failure
2238                // path, otherwise a broken flownode would fill the budget forever.
2239                release_mirror_pending_rows(&mirror_pending_rows, peer_rows);
2240            });
2241        }
2242
2243        Ok(())
2244    }
2245}
2246
2247/// Rows of row data carried by `inserts`. A request without a `rows` payload
2248/// carries nothing to keep buffered and contributes zero.
2249fn region_inserts_rows(inserts: &RegionInsertRequests) -> u64 {
2250    inserts
2251        .requests
2252        .iter()
2253        .filter_map(|req| req.rows.as_ref())
2254        .map(|rows| rows.rows.len() as u64)
2255        .sum()
2256}
2257
2258/// Reserves `rows` of the pending mirror budget and returns the resulting
2259/// pending row count. The caller must [`release_mirror_pending_rows`] exactly
2260/// these rows once they are no longer buffered, whether it spawns the batch or
2261/// drops it.
2262fn reserve_mirror_pending_rows(mirror_pending_rows: &AtomicU64, rows: u64) -> u64 {
2263    let pending = mirror_pending_rows.fetch_add(rows, Ordering::Relaxed) + rows;
2264    // Keep the gauge in step by applying the same increment, never by setting an
2265    // absolute value: another task may move the budget in between.
2266    crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.add(rows as i64);
2267    pending
2268}
2269
2270/// Releases `rows` of the pending mirror budget.
2271///
2272/// Every release mirrors the reservation it hands back, so a well-behaved
2273/// release never subtracts more than is pending; the saturating update only
2274/// keeps a duplicated release from wrapping the shared budget around.
2275fn release_mirror_pending_rows(mirror_pending_rows: &AtomicU64, rows: u64) {
2276    let prev = mirror_pending_rows
2277        .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |pending| {
2278            Some(pending.saturating_sub(rows))
2279        })
2280        .unwrap_or_else(|prev| prev);
2281    // Subtract from the gauge with the opposite sign of
2282    // [`reserve_mirror_pending_rows`]; reading the budget here and `set`-ting it
2283    // would publish a stale value whenever another task releases concurrently.
2284    let released = rows.min(prev);
2285    if released > 0 {
2286        crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.sub(released as i64);
2287    }
2288}
2289
2290#[cfg(test)]
2291mod tests {
2292    use std::sync::Arc;
2293    use std::sync::atomic::{AtomicU64, Ordering};
2294
2295    use api::helper::ColumnDataTypeWrapper;
2296    use api::v1::flow::FlowResponse;
2297    use api::v1::helper::{field_column_schema, time_index_column_schema};
2298    use api::v1::{Row, RowInsertRequest, Rows, Value};
2299    use common_catalog::consts::{DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME};
2300    use common_meta::cache::new_table_flownode_set_cache;
2301    use common_meta::ddl::test_util::datanode_handler::NaiveDatanodeHandler;
2302    use common_meta::test_util::{MockDatanodeManager, MockFlownodeHandler, MockFlownodeManager};
2303    use common_query::native_histogram::NATIVE_HISTOGRAM_FIELD;
2304    use common_query::prelude::{greptime_native_histogram, set_default_prefix};
2305    use datatypes::data_type::ConcreteDataType;
2306    use datatypes::schema::ColumnSchema;
2307    use moka::future::Cache;
2308    use session::context::QueryContext;
2309    use table::TableRef;
2310    use table::dist_table::DummyDataSource;
2311    use table::metadata::{TableInfoBuilder, TableMetaBuilder, TableType};
2312
2313    use crate::insert::*;
2314    use crate::test_util::{
2315        create_partition_rule_manager, new_test_table_info, prepare_mocked_backend,
2316    };
2317
2318    fn make_table_ref_with_schema(
2319        ts_name: &str,
2320        field_name: &str,
2321        field_type: ConcreteDataType,
2322    ) -> TableRef {
2323        let schema = datatypes::schema::SchemaBuilder::try_from_columns(vec![
2324            ColumnSchema::new(
2325                ts_name,
2326                ConcreteDataType::timestamp_millisecond_datatype(),
2327                false,
2328            )
2329            .with_time_index(true),
2330            ColumnSchema::new(field_name, field_type, true),
2331        ])
2332        .unwrap()
2333        .build()
2334        .unwrap();
2335        let meta = TableMetaBuilder::empty()
2336            .schema(Arc::new(schema))
2337            .primary_key_indices(vec![])
2338            .value_indices(vec![1])
2339            .engine("mito")
2340            .next_column_id(0)
2341            .options(Default::default())
2342            .created_on(Default::default())
2343            .build()
2344            .unwrap();
2345        let info = Arc::new(
2346            TableInfoBuilder::default()
2347                .table_id(1)
2348                .table_version(0)
2349                .name("test_table")
2350                .schema_name(DEFAULT_SCHEMA_NAME)
2351                .catalog_name(DEFAULT_CATALOG_NAME)
2352                .desc(None)
2353                .table_type(TableType::Base)
2354                .meta(meta)
2355                .build()
2356                .unwrap(),
2357        );
2358        Arc::new(table::Table::new(
2359            info,
2360            table::metadata::FilterPushDownType::Unsupported,
2361            Arc::new(DummyDataSource),
2362        ))
2363    }
2364
2365    fn make_metric_physical_table_ref_with_time_unit(unit: TimeUnit) -> TableRef {
2366        let schema = datatypes::schema::SchemaBuilder::try_from_columns(vec![
2367            ColumnSchema::new(
2368                greptime_timestamp(),
2369                ConcreteDataType::timestamp_datatype(unit),
2370                false,
2371            )
2372            .with_time_index(true),
2373            ColumnSchema::new(greptime_value(), ConcreteDataType::float64_datatype(), true),
2374        ])
2375        .unwrap()
2376        .build()
2377        .unwrap();
2378        let meta = TableMetaBuilder::empty()
2379            .schema(Arc::new(schema))
2380            .primary_key_indices(vec![])
2381            .value_indices(vec![1])
2382            .engine("metric")
2383            .next_column_id(0)
2384            .options(Default::default())
2385            .created_on(Default::default())
2386            .build()
2387            .unwrap();
2388        let info = Arc::new(
2389            TableInfoBuilder::default()
2390                .table_id(1)
2391                .table_version(0)
2392                .name("greptime_physical_table")
2393                .schema_name(DEFAULT_SCHEMA_NAME)
2394                .catalog_name(DEFAULT_CATALOG_NAME)
2395                .desc(None)
2396                .table_type(TableType::Base)
2397                .meta(meta)
2398                .build()
2399                .unwrap(),
2400        );
2401        Arc::new(table::Table::new(
2402            info,
2403            table::metadata::FilterPushDownType::Unsupported,
2404            Arc::new(DummyDataSource),
2405        ))
2406    }
2407
2408    fn ms_row_insert_request(timestamp_ms: i64) -> RowInsertRequest {
2409        ms_row_insert_request_named("my_metric", timestamp_ms)
2410    }
2411
2412    fn ms_row_insert_request_named(table: &str, timestamp_ms: i64) -> RowInsertRequest {
2413        RowInsertRequest {
2414            table_name: table.to_string(),
2415            rows: Some(Rows {
2416                schema: vec![
2417                    time_index_column_schema(
2418                        greptime_timestamp(),
2419                        ColumnDataType::TimestampMillisecond,
2420                    ),
2421                    field_column_schema(greptime_value(), ColumnDataType::Float64),
2422                ],
2423                rows: vec![Row {
2424                    values: vec![
2425                        Value {
2426                            value_data: Some(ValueData::TimestampMillisecondValue(timestamp_ms)),
2427                        },
2428                        Value {
2429                            value_data: Some(ValueData::F64Value(1.0)),
2430                        },
2431                    ],
2432                }],
2433            }),
2434        }
2435    }
2436
2437    /// Converts the time index column of each request to `target_unit`.
2438    fn convert_row_insert_requests_time_unit(
2439        requests: &mut RowInsertRequests,
2440        target_unit: TimeUnit,
2441    ) -> Result<()> {
2442        for request in &mut requests.inserts {
2443            let Some(rows) = request.rows.as_mut() else {
2444                continue;
2445            };
2446            convert_rows_time_unit(rows, target_unit)?;
2447        }
2448        Ok(())
2449    }
2450
2451    #[test]
2452    fn test_convert_row_insert_requests_time_unit_noop_when_matching() {
2453        let mut requests = RowInsertRequests {
2454            inserts: vec![ms_row_insert_request(123)],
2455        };
2456        convert_row_insert_requests_time_unit(&mut requests, TimeUnit::Millisecond).unwrap();
2457        let rows = requests.inserts[0].rows.as_ref().unwrap();
2458        assert_eq!(
2459            rows.schema[0].datatype,
2460            ColumnDataType::TimestampMillisecond as i32
2461        );
2462        assert!(matches!(
2463            rows.rows[0].values[0].value_data,
2464            Some(ValueData::TimestampMillisecondValue(123))
2465        ));
2466    }
2467
2468    #[test]
2469    fn test_convert_row_insert_requests_time_unit_widens_losslessly() {
2470        let mut requests = RowInsertRequests {
2471            inserts: vec![ms_row_insert_request(123)],
2472        };
2473        convert_row_insert_requests_time_unit(&mut requests, TimeUnit::Microsecond).unwrap();
2474        let rows = requests.inserts[0].rows.as_ref().unwrap();
2475        assert_eq!(
2476            rows.schema[0].datatype,
2477            ColumnDataType::TimestampMicrosecond as i32
2478        );
2479        assert!(matches!(
2480            rows.rows[0].values[0].value_data,
2481            Some(ValueData::TimestampMicrosecondValue(123_000))
2482        ));
2483    }
2484
2485    #[test]
2486    fn test_convert_row_insert_requests_time_unit_truncates_on_narrowing() {
2487        // 123_456_789 ns floors to 123_456 us and 123 ms; negative values
2488        // floor towards negative infinity, matching `Timestamp::convert_to`.
2489        let requests = |value_ns: i64| RowInsertRequests {
2490            inserts: vec![RowInsertRequest {
2491                table_name: "my_metric".to_string(),
2492                rows: Some(Rows {
2493                    schema: vec![
2494                        time_index_column_schema(
2495                            greptime_timestamp(),
2496                            ColumnDataType::TimestampNanosecond,
2497                        ),
2498                        field_column_schema(greptime_value(), ColumnDataType::Float64),
2499                    ],
2500                    rows: vec![Row {
2501                        values: vec![
2502                            Value {
2503                                value_data: Some(ValueData::TimestampNanosecondValue(value_ns)),
2504                            },
2505                            Value {
2506                                value_data: Some(ValueData::F64Value(1.0)),
2507                            },
2508                        ],
2509                    }],
2510                }),
2511            }],
2512        };
2513
2514        let mut reqs = requests(123_456_789);
2515        convert_row_insert_requests_time_unit(&mut reqs, TimeUnit::Microsecond).unwrap();
2516        let rows = reqs.inserts[0].rows.as_ref().unwrap();
2517        assert!(matches!(
2518            rows.rows[0].values[0].value_data,
2519            Some(ValueData::TimestampMicrosecondValue(123_456))
2520        ));
2521
2522        let mut reqs = requests(123_456_789);
2523        convert_row_insert_requests_time_unit(&mut reqs, TimeUnit::Millisecond).unwrap();
2524        let rows = reqs.inserts[0].rows.as_ref().unwrap();
2525        assert!(matches!(
2526            rows.rows[0].values[0].value_data,
2527            Some(ValueData::TimestampMillisecondValue(123))
2528        ));
2529
2530        let mut reqs = requests(-123_456_789);
2531        convert_row_insert_requests_time_unit(&mut reqs, TimeUnit::Microsecond).unwrap();
2532        let rows = reqs.inserts[0].rows.as_ref().unwrap();
2533        assert!(matches!(
2534            rows.rows[0].values[0].value_data,
2535            Some(ValueData::TimestampMicrosecondValue(-123_457))
2536        ));
2537    }
2538
2539    #[test]
2540    fn test_convert_row_insert_requests_time_unit_overflow() {
2541        let mut requests = RowInsertRequests {
2542            inserts: vec![ms_row_insert_request(i64::MAX)],
2543        };
2544        let err =
2545            convert_row_insert_requests_time_unit(&mut requests, TimeUnit::Nanosecond).unwrap_err();
2546        assert!(err.to_string().contains("overflows"), "{err}");
2547    }
2548
2549    #[test]
2550    fn test_table_time_index_unit() {
2551        assert_eq!(
2552            table_time_index_unit(&make_metric_physical_table_ref_with_time_unit(
2553                TimeUnit::Microsecond
2554            )),
2555            Some(TimeUnit::Microsecond)
2556        );
2557        assert_eq!(
2558            table_time_index_unit(&make_metric_physical_table_ref_with_time_unit(
2559                TimeUnit::Millisecond
2560            )),
2561            Some(TimeUnit::Millisecond)
2562        );
2563    }
2564
2565    #[tokio::test]
2566    async fn test_accommodate_existing_schema_and_reject_kind_changes() {
2567        let ts_name = "my_ts";
2568        let field_name = "my_field";
2569        let table =
2570            make_table_ref_with_schema(ts_name, field_name, ConcreteDataType::float64_datatype());
2571
2572        // The request uses different names for timestamp and field columns
2573        let mut req = RowInsertRequest {
2574            table_name: "test_table".to_string(),
2575            rows: Some(Rows {
2576                schema: vec![
2577                    time_index_column_schema("ts_wrong", ColumnDataType::TimestampMillisecond),
2578                    field_column_schema("field_wrong", ColumnDataType::Float64),
2579                ],
2580                rows: vec![api::v1::Row {
2581                    values: vec![Value::default(), Value::default()],
2582                }],
2583            }),
2584        };
2585        let ctx = Arc::new(QueryContext::with(
2586            DEFAULT_CATALOG_NAME,
2587            DEFAULT_SCHEMA_NAME,
2588        ));
2589
2590        let kv_backend = prepare_mocked_backend().await;
2591        let inserter = Inserter::new(
2592            catalog::memory::MemoryCatalogManager::new(),
2593            create_partition_rule_manager(kv_backend.clone()).await,
2594            Arc::new(MockDatanodeManager::new(NaiveDatanodeHandler)),
2595            Arc::new(new_table_flownode_set_cache(
2596                String::new(),
2597                Cache::new(100),
2598                kv_backend.clone(),
2599            )),
2600            true,
2601        );
2602        // Do not apply an absent-table plan to a table that appeared concurrently.
2603        assert!(
2604            inserter
2605                .get_alter_table_expr_on_demand(&mut req, &table, &ctx, false, false, false)
2606                .unwrap()
2607                .is_none()
2608        );
2609        let alter_expr = inserter
2610            .get_alter_table_expr_on_demand(&mut req, &table, &ctx, true, true, true)
2611            .unwrap();
2612        assert!(alter_expr.is_none());
2613
2614        // The request's schema should have updated names for timestamp and field columns
2615        let req_schema = req.rows.as_ref().unwrap().schema.clone();
2616        assert_eq!(req_schema[0].column_name, ts_name);
2617        assert_eq!(req_schema[1].column_name, field_name);
2618
2619        let (datatype, datatype_extension) =
2620            ColumnDataTypeWrapper::try_from(native_histogram_value_type().clone())
2621                .unwrap()
2622                .into_parts();
2623        let mut histogram_req = RowInsertRequest {
2624            table_name: "test_table".to_string(),
2625            rows: Some(Rows {
2626                schema: vec![
2627                    time_index_column_schema("ts", ColumnDataType::TimestampMillisecond),
2628                    api::v1::ColumnSchema {
2629                        column_name: greptime_native_histogram().to_string(),
2630                        datatype: datatype as i32,
2631                        semantic_type: SemanticType::Field as i32,
2632                        datatype_extension,
2633                        options: None,
2634                    },
2635                ],
2636                rows: vec![],
2637            }),
2638        };
2639        let error = inserter
2640            .get_alter_table_expr_on_demand(&mut histogram_req, &table, &ctx, false, true, true)
2641            .unwrap_err();
2642        assert!(
2643            error
2644                .to_string()
2645                .contains("cannot mix native histogram and float sample fields")
2646        );
2647
2648        let histogram_table = make_table_ref_with_schema(
2649            "ts",
2650            greptime_native_histogram(),
2651            native_histogram_value_type().clone(),
2652        );
2653        let mut sample_req = RowInsertRequest {
2654            table_name: "test_table".to_string(),
2655            rows: Some(Rows {
2656                schema: vec![
2657                    time_index_column_schema("ts", ColumnDataType::TimestampMillisecond),
2658                    field_column_schema(greptime_value(), ColumnDataType::Float64),
2659                ],
2660                rows: vec![],
2661            }),
2662        };
2663        let error = inserter
2664            .get_alter_table_expr_on_demand(
2665                &mut sample_req,
2666                &histogram_table,
2667                &ctx,
2668                false,
2669                true,
2670                true,
2671            )
2672            .unwrap_err();
2673        assert!(
2674            error
2675                .to_string()
2676                .contains("cannot mix native histogram and float sample fields")
2677        );
2678    }
2679
2680    #[test]
2681    fn test_native_histogram_detection_survives_prefix_change() {
2682        set_default_prefix(Some("custom")).unwrap();
2683        let table = make_table_ref_with_schema(
2684            "custom_timestamp",
2685            NATIVE_HISTOGRAM_FIELD,
2686            native_histogram_value_type().clone(),
2687        );
2688        let (datatype, datatype_extension) =
2689            ColumnDataTypeWrapper::try_from(native_histogram_value_type().clone())
2690                .unwrap()
2691                .into_parts();
2692        let request_schema = [api::v1::ColumnSchema {
2693            column_name: greptime_native_histogram().to_string(),
2694            datatype: datatype as i32,
2695            semantic_type: SemanticType::Field as i32,
2696            datatype_extension,
2697            options: None,
2698        }];
2699
2700        assert!(request_is_native_histogram(&request_schema));
2701        assert!(table_is_native_histogram(&table));
2702    }
2703
2704    // Keep global meter registration in one test, isolated by nextest's per-test process.
2705    #[tokio::test]
2706    async fn test_write_meter_admission() {
2707        use std::cell::Cell;
2708        use std::sync::Mutex;
2709        use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
2710
2711        use api::region::RegionResponse;
2712        use api::v1::region::region_request::Body;
2713        use arrow::array::{Int32Array, TimestampMillisecondArray};
2714        use arrow::record_batch::RecordBatch;
2715        use bytes::Bytes;
2716        use common_error::ext::{ErrorExt, RetryHint};
2717        use common_error::status_code::StatusCode;
2718        use common_grpc::flight::{FlightEncoder, FlightMessage};
2719        use common_meta::ddl::test_util::datanode_handler::DatanodeWatcher;
2720        use futures::future::BoxFuture;
2721        use meter_core::ItemCalculator;
2722        use meter_core::collect::{Collect, WriteRejected};
2723        use meter_core::data::MeterRecord;
2724        use meter_core::global::global_registry;
2725        use session::context::Channel;
2726
2727        const CATALOG: &str = "write_meter_test";
2728
2729        #[derive(Default)]
2730        struct Meter {
2731            reject: AtomicBool,
2732            attempts: Mutex<Vec<MeterRecord>>,
2733            accepted_value: AtomicU64,
2734        }
2735
2736        impl Collect for Meter {
2737            fn on_write(
2738                &self,
2739                record: MeterRecord,
2740            ) -> BoxFuture<'_, std::result::Result<(), WriteRejected>> {
2741                Box::pin(async move {
2742                    if record.catalog != CATALOG {
2743                        return Ok(());
2744                    }
2745                    let value = record.value;
2746                    self.attempts.lock().unwrap().push(record);
2747                    if self.reject.load(Ordering::Relaxed) {
2748                        return Err(WriteRejected::new("database row quota exhausted"));
2749                    }
2750                    self.accepted_value.fetch_add(value, Ordering::Relaxed);
2751                    Ok(())
2752                })
2753            }
2754
2755            fn on_read(&self, _: MeterRecord) {}
2756        }
2757
2758        impl ItemCalculator<InstantAndNormalInsertRequests> for Meter {
2759            fn calc(&self, _: &InstantAndNormalInsertRequests) -> u64 {
2760                17
2761            }
2762        }
2763
2764        let kv_backend = prepare_mocked_backend().await;
2765        let partition_manager = create_partition_rule_manager(kv_backend.clone()).await;
2766        let (sender, mut dispatched) = tokio::sync::mpsc::channel(16);
2767        let watcher = DatanodeWatcher::new(sender).with_handler(|_, request| {
2768            let rows = match request.body.unwrap() {
2769                Body::Inserts(requests) => requests
2770                    .requests
2771                    .iter()
2772                    .filter_map(|request| request.rows.as_ref())
2773                    .map(|rows| rows.rows.len())
2774                    .sum(),
2775                // The bulk batches below each contain two rows.
2776                Body::BulkInsert(_) => 2,
2777                body => panic!("unexpected request: {body:?}"),
2778            };
2779            Ok(RegionResponse::new(rows))
2780        });
2781        let flow_cache = Cache::new(100);
2782        let inserter = Inserter::new(
2783            catalog::memory::MemoryCatalogManager::new(),
2784            partition_manager,
2785            Arc::new(MockDatanodeManager::new(watcher)),
2786            Arc::new(new_table_flownode_set_cache(
2787                String::new(),
2788                flow_cache.clone(),
2789                kv_backend,
2790            )),
2791            true,
2792        );
2793        let mut table_info = new_test_table_info(1, "table_1", [1].into_iter());
2794        table_info.catalog_name = CATALOG.to_string();
2795        table_info.schema_name = "target_db".to_string();
2796        let table_info = Arc::new(table_info);
2797        let table_infos = HashMap::from_iter([(1, table_info.clone())]);
2798        let ctx = Arc::new(QueryContext::with_channel(
2799            DEFAULT_CATALOG_NAME,
2800            DEFAULT_SCHEMA_NAME,
2801            Channel::Postgres,
2802        ));
2803        let rows_request = || {
2804            let build = |num_rows| RegionInsertRequests {
2805                requests: vec![RegionInsertRequest {
2806                    region_id: RegionId::new(1, 1).as_u64(),
2807                    rows: Some(Rows {
2808                        schema: vec![],
2809                        rows: vec![api::v1::Row { values: vec![] }; num_rows],
2810                    }),
2811                    ..Default::default()
2812                }],
2813            };
2814            InstantAndNormalInsertRequests {
2815                normal_requests: build(3),
2816                instant_requests: build(2),
2817            }
2818        };
2819
2820        // No collector or calculator: ordinary OSS insertion still succeeds.
2821        let output = inserter
2822            .do_request(rows_request(), &table_infos, &ctx)
2823            .await
2824            .unwrap();
2825        assert_eq!(output.meta.cost, 0);
2826        assert!(matches!(output.data, OutputData::AffectedRows(3)));
2827        dispatched.try_recv().unwrap();
2828
2829        let meter = Arc::new(Meter::default());
2830        global_registry().set_collector(meter.clone());
2831        global_registry().register_calculator(meter.clone());
2832        // The dependency controls noop mode; exercise both builds with this test.
2833        let enabled = Cell::new(false);
2834        write_meter!({
2835            enabled.set(true);
2836            MeterRecord::new("probe".into(), "probe".into(), 0, 0, 0)
2837        })
2838        .await
2839        .unwrap();
2840        let enabled = enabled.get();
2841
2842        let output = inserter
2843            .do_request(rows_request(), &table_infos, &ctx)
2844            .await
2845            .unwrap();
2846        assert_eq!(output.meta.cost, if enabled { 17 } else { 0 });
2847        assert!(matches!(output.data, OutputData::AffectedRows(3)));
2848        dispatched.try_recv().unwrap();
2849        assert!(flow_cache.contains_key(&1));
2850        flow_cache.invalidate_all();
2851        meter.reject.store(true, Ordering::Relaxed);
2852        let result = inserter
2853            .do_request(rows_request(), &table_infos, &ctx)
2854            .await;
2855        if enabled {
2856            let error = result.unwrap_err();
2857            assert_eq!(error.status_code(), StatusCode::RateLimited);
2858            assert_eq!(error.retry_hint(), RetryHint::Retryable);
2859            assert!(error.to_string().contains("database row quota exhausted"));
2860            assert!(dispatched.try_recv().is_err());
2861            assert!(
2862                !flow_cache.contains_key(&1),
2863                "rejected write reached flow mirroring"
2864            );
2865            assert_eq!(meter.accepted_value.load(Ordering::Relaxed), 17);
2866            let attempts = meter.attempts.lock().unwrap();
2867            assert_eq!(attempts.len(), 2);
2868            for record in attempts.iter() {
2869                assert_eq!(record.catalog, CATALOG);
2870                assert_eq!(record.schema, "target_db");
2871                assert_eq!(
2872                    (record.rows, record.value, record.source),
2873                    (5, 17, Channel::Postgres as u8)
2874                );
2875            }
2876        } else {
2877            assert_eq!(result.unwrap().meta.cost, 0);
2878            dispatched.try_recv().unwrap();
2879            assert!(meter.attempts.lock().unwrap().is_empty());
2880        }
2881        meter.attempts.lock().unwrap().clear();
2882
2883        // Rejection must also precede admission to the table batcher's queue.
2884        if enabled {
2885            let batcher: Arc<dyn PendingRowsBatcher> = Arc::new(UnexpectedBatcher);
2886            let rows = Rows {
2887                schema: vec![
2888                    api::v1::helper::tag_column_schema("a", ColumnDataType::Int32),
2889                    time_index_column_schema("ts", ColumnDataType::TimestampMillisecond),
2890                    field_column_schema("b", ColumnDataType::Int32),
2891                ],
2892                rows: vec![api::v1::Row {
2893                    values: vec![
2894                        api::v1::value::ValueData::I32Value(60).into(),
2895                        Value {
2896                            value_data: Some(api::v1::value::ValueData::TimestampMillisecondValue(
2897                                0,
2898                            )),
2899                        },
2900                        api::v1::value::ValueData::I32Value(0).into(),
2901                    ],
2902                }],
2903            };
2904            let error = inserter
2905                .submit_table_rows(rows, table_info.clone(), ctx.clone(), &batcher)
2906                .await
2907                .unwrap_err();
2908            assert_eq!(error.status_code(), StatusCode::RateLimited);
2909            let mut attempts = meter.attempts.lock().unwrap();
2910            assert_eq!(attempts.len(), 1);
2911            assert_eq!(attempts[0].catalog, CATALOG);
2912            assert_eq!(attempts[0].schema, "target_db");
2913            assert_eq!(attempts[0].rows, 1);
2914            attempts.clear();
2915        }
2916
2917        let table = Arc::new(table::Table::new(
2918            table_info.clone(),
2919            table::metadata::FilterPushDownType::Unsupported,
2920            Arc::new(DummyDataSource),
2921        ));
2922        let batch = RecordBatch::try_new(
2923            table_info.meta.schema.arrow_schema().clone(),
2924            vec![
2925                Arc::new(Int32Array::from(vec![60, 70])),
2926                Arc::new(TimestampMillisecondArray::from(vec![0, 1])),
2927                Arc::new(Int32Array::from(vec![0, 0])),
2928            ],
2929        )
2930        .unwrap();
2931        let bulk_insert = |batch: RecordBatch| {
2932            let flight_data = FlightEncoder::default()
2933                .encode(FlightMessage::RecordBatch(batch.clone()))
2934                .into_iter()
2935                .next()
2936                .unwrap();
2937            inserter.handle_bulk_insert(
2938                table.clone(),
2939                flight_data,
2940                batch,
2941                Bytes::new(),
2942                false,
2943                Channel::Grpc,
2944            )
2945        };
2946        // Empty batches bypass admission even while the collector rejects.
2947        assert_eq!(bulk_insert(batch.slice(0, 0)).await.unwrap(), 0);
2948        assert!(meter.attempts.lock().unwrap().is_empty());
2949        assert!(dispatched.try_recv().is_err());
2950
2951        // A collector change between batches takes effect on the very next batch.
2952        for reject in [false, true, false] {
2953            meter.reject.store(reject, Ordering::Relaxed);
2954            let result = bulk_insert(batch.clone()).await;
2955            if enabled && reject {
2956                let error = result.unwrap_err();
2957                assert_eq!(error.status_code(), StatusCode::RateLimited);
2958                assert_eq!(error.retry_hint(), RetryHint::Retryable);
2959                assert!(dispatched.try_recv().is_err());
2960            } else {
2961                assert_eq!(result.unwrap(), 2);
2962                dispatched.try_recv().unwrap();
2963            }
2964        }
2965        {
2966            let attempts = meter.attempts.lock().unwrap();
2967            assert_eq!(attempts.len(), if enabled { 3 } else { 0 });
2968            for record in attempts.iter() {
2969                assert_eq!(record.catalog, CATALOG);
2970                assert_eq!(record.schema, "target_db");
2971                assert_eq!(
2972                    (record.rows, record.value, record.source),
2973                    (2, 0, Channel::Grpc as u8)
2974                );
2975            }
2976        }
2977        assert_eq!(
2978            meter.accepted_value.load(Ordering::Relaxed),
2979            if enabled { 17 } else { 0 }
2980        );
2981
2982        // Aggregate repeated database targets and retain admission across nested
2983        // batching without changing the caller's reusable context.
2984        meter.attempts.lock().unwrap().clear();
2985        let original = Arc::new(QueryContext::with_channel(
2986            CATALOG,
2987            "a",
2988            Channel::Prometheus,
2989        ));
2990        let mut batches = ["a", "b", "a"].map(|schema| {
2991            let ctx = if schema == "a" {
2992                original.clone()
2993            } else {
2994                Arc::new(QueryContext::with_channel(
2995                    CATALOG,
2996                    schema,
2997                    Channel::Prometheus,
2998                ))
2999            };
3000            (
3001                ctx,
3002                RowInsertRequests {
3003                    inserts: vec![RowInsertRequest {
3004                        table_name: "data".into(),
3005                        rows: Some(Rows {
3006                            schema: vec![],
3007                            rows: vec![api::v1::Row::default(); 2],
3008                        }),
3009                    }],
3010                },
3011            )
3012        });
3013        admit_row_insert_batches(&mut batches).await.unwrap();
3014        admit_row_insert_batches(&mut batches).await.unwrap();
3015        assert_eq!(original.write_rows_to_admit(CATALOG, "a", 4), 4);
3016        for (ctx, _) in &batches {
3017            assert_eq!(
3018                ctx.write_rows_to_admit(CATALOG, &ctx.current_schema(), 2),
3019                0
3020            );
3021            assert_eq!(ctx.channel(), Channel::Prometheus);
3022        }
3023        {
3024            let attempts = meter.attempts.lock().unwrap();
3025            let totals = attempts
3026                .iter()
3027                .map(|r| (r.schema.as_str(), r.rows, r.value))
3028                .collect::<Vec<_>>();
3029            assert_eq!(
3030                totals,
3031                if enabled {
3032                    vec![("a", 4, 0), ("b", 2, 0)]
3033                } else {
3034                    vec![]
3035                }
3036            );
3037        }
3038        meter.attempts.lock().unwrap().clear();
3039        meter.reject.store(false, Ordering::Relaxed);
3040        let ctx = Arc::new(QueryContext::with_channel(
3041            CATALOG,
3042            "logical",
3043            Channel::Otlp,
3044        ));
3045        let admitted = admit_write(2, &ctx).await.unwrap();
3046        let mut requests = RowInsertRequests {
3047            inserts: vec![RowInsertRequest {
3048                table_name: "metric".to_string(),
3049                rows: Some(Rows {
3050                    schema: vec![],
3051                    rows: vec![api::v1::Row::default(); 2],
3052                }),
3053            }],
3054        };
3055        let original = requests.clone();
3056        let cost = Inserter::meter_row_inserts(&mut requests, &admitted)
3057            .await
3058            .unwrap();
3059        assert_eq!(cost, if enabled { 17 } else { 0 });
3060        assert_eq!(requests, original);
3061        let records = meter
3062            .attempts
3063            .lock()
3064            .unwrap()
3065            .iter()
3066            .map(|record| (record.rows, record.value))
3067            .collect::<Vec<_>>();
3068        assert_eq!(
3069            records,
3070            if enabled {
3071                vec![(2, 0), (0, 17)]
3072            } else {
3073                vec![]
3074            }
3075        );
3076        meter.reject.store(true, Ordering::Relaxed);
3077        let result = Inserter::meter_row_inserts(&mut requests, &admitted).await;
3078        if enabled {
3079            assert_eq!(result.unwrap_err().status_code(), StatusCode::RateLimited);
3080        } else {
3081            assert_eq!(result.unwrap(), 0);
3082        }
3083        assert_eq!(requests, original);
3084    }
3085
3086    #[test]
3087    fn test_skip_wal_does_not_change_table_options() {
3088        check_skip_wal_does_not_change_table_options(false);
3089        check_skip_wal_does_not_change_table_options(true);
3090    }
3091
3092    fn check_skip_wal_does_not_change_table_options(skip_wal: bool) {
3093        let ctx = Arc::new(QueryContext::with(
3094            DEFAULT_CATALOG_NAME,
3095            DEFAULT_SCHEMA_NAME,
3096        ));
3097        ctx.set_skip_wal(skip_wal);
3098        let mut options = Default::default();
3099        fill_table_options_for_create(&mut options, &AutoCreateTableType::Physical, &ctx);
3100        assert!(!options.contains_key(session::hints::INSERT_SKIP_WAL_HINT));
3101        assert!(!options.contains_key("skip_wal"));
3102    }
3103
3104    #[test]
3105    fn test_last_non_null_create_options_preserve_default_without_append_mode() {
3106        let ctx = Arc::new(QueryContext::with(
3107            DEFAULT_CATALOG_NAME,
3108            DEFAULT_SCHEMA_NAME,
3109        ));
3110        let mut table_options = Default::default();
3111
3112        fill_table_options_for_create(&mut table_options, &AutoCreateTableType::LastNonNull, &ctx);
3113
3114        assert_eq!(
3115            Some("last_non_null"),
3116            table_options.get(MERGE_MODE_KEY).map(String::as_str)
3117        );
3118        assert!(!table_options.contains_key(APPEND_MODE_KEY));
3119    }
3120
3121    #[test]
3122    fn test_fill_table_options_copies_semantic_extensions() {
3123        use table::requests::{
3124            SEMANTIC_METRIC_TYPE, SEMANTIC_PER_TABLE_INDEX_KEY, SEMANTIC_SIGNAL_TYPE,
3125            SEMANTIC_SOURCE, SEMANTIC_SOURCE_VERSION, SIGNAL_TYPE_METRIC, SOURCE_OPENTELEMETRY,
3126        };
3127
3128        let mut ctx = QueryContext::with(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME);
3129        ctx.set_extension(SEMANTIC_SIGNAL_TYPE, SIGNAL_TYPE_METRIC);
3130        ctx.set_extension(SEMANTIC_SOURCE, SOURCE_OPENTELEMETRY);
3131        ctx.set_extension(SEMANTIC_SOURCE_VERSION, "2.0");
3132        ctx.set_extension(SEMANTIC_METRIC_TYPE, "bogus");
3133        // The internal transport key must NOT be copied into table options.
3134        ctx.set_extension(SEMANTIC_PER_TABLE_INDEX_KEY, "{}");
3135        let ctx = Arc::new(ctx);
3136        let mut table_options = Default::default();
3137
3138        fill_table_options_for_create(&mut table_options, &AutoCreateTableType::Physical, &ctx);
3139
3140        assert_eq!(
3141            Some(SIGNAL_TYPE_METRIC),
3142            table_options.get(SEMANTIC_SIGNAL_TYPE).map(String::as_str)
3143        );
3144        assert_eq!(
3145            Some(SOURCE_OPENTELEMETRY),
3146            table_options.get(SEMANTIC_SOURCE).map(String::as_str)
3147        );
3148        assert_eq!(
3149            Some("2.0"),
3150            table_options
3151                .get(SEMANTIC_SOURCE_VERSION)
3152                .map(String::as_str)
3153        );
3154        assert!(!table_options.contains_key(SEMANTIC_METRIC_TYPE));
3155        assert!(!table_options.contains_key(SEMANTIC_PER_TABLE_INDEX_KEY));
3156    }
3157
3158    #[test]
3159    fn test_apply_per_table_semantic_options() {
3160        use table::requests::{
3161            SEMANTIC_METRIC_TYPE, SEMANTIC_METRIC_UNIT, SEMANTIC_PER_TABLE_INDEX_KEY,
3162        };
3163
3164        let index = format!(
3165            r#"{{
3166            "{DEFAULT_SCHEMA_NAME}": {{
3167                "http_requests_total": {{
3168                    "greptime.semantic.metric.type": "counter",
3169                    "greptime.semantic.metric.unit": "By",
3170                    "greptime.semantic.metric.type_BOGUS": "x"
3171                }},
3172                "other_table": {{
3173                    "greptime.semantic.metric.type": "gauge"
3174                }}
3175            }},
3176            "other_schema": {{
3177                "http_requests_total": {{
3178                    "greptime.semantic.metric.type": "gauge"
3179                }}
3180            }}
3181        }}"#
3182        );
3183        let mut ctx = QueryContext::with(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME);
3184        ctx.set_extension(SEMANTIC_PER_TABLE_INDEX_KEY, index);
3185        let ctx = Arc::new(ctx);
3186
3187        let index = parse_per_table_semantic_index(&ctx);
3188        assert!(index.is_some());
3189        let index = index.as_ref();
3190
3191        let mut table_options = std::collections::HashMap::new();
3192        apply_per_table_semantic_options(
3193            &mut table_options,
3194            index,
3195            DEFAULT_SCHEMA_NAME,
3196            "http_requests_total",
3197        );
3198        // The write schema's entry applies — not other_schema's `gauge`.
3199        assert_eq!(
3200            table_options.get(SEMANTIC_METRIC_TYPE).map(String::as_str),
3201            Some("counter")
3202        );
3203        assert_eq!(
3204            table_options.get(SEMANTIC_METRIC_UNIT).map(String::as_str),
3205            Some("By")
3206        );
3207        // The unknown key is rejected by the vocabulary check; other tables' keys
3208        // never appear.
3209        assert!(!table_options.contains_key("greptime.semantic.metric.type_BOGUS"));
3210        assert_eq!(table_options.len(), 2);
3211
3212        let mut empty = std::collections::HashMap::new();
3213        apply_per_table_semantic_options(&mut empty, index, DEFAULT_SCHEMA_NAME, "not_in_index");
3214        assert!(empty.is_empty());
3215
3216        // A schema with no entry is a no-op even when the table name matches
3217        // elsewhere.
3218        let mut opts = std::collections::HashMap::new();
3219        apply_per_table_semantic_options(
3220            &mut opts,
3221            index,
3222            "schema_without_entry",
3223            "http_requests_total",
3224        );
3225        assert!(opts.is_empty());
3226
3227        // No extension at all parses to no index (e.g. logs / Prom RW v1).
3228        let bare = Arc::new(QueryContext::with(
3229            DEFAULT_CATALOG_NAME,
3230            DEFAULT_SCHEMA_NAME,
3231        ));
3232        assert!(parse_per_table_semantic_index(&bare).is_none());
3233        let mut opts = std::collections::HashMap::new();
3234        apply_per_table_semantic_options(
3235            &mut opts,
3236            None,
3237            DEFAULT_SCHEMA_NAME,
3238            "http_requests_total",
3239        );
3240        assert!(opts.is_empty());
3241    }
3242
3243    #[test]
3244    fn test_last_non_null_create_options_preserve_default_with_append_mode_false() {
3245        let mut ctx = QueryContext::with(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME);
3246        ctx.set_extension(APPEND_MODE_KEY, "false");
3247        let ctx = Arc::new(ctx);
3248        let mut table_options = Default::default();
3249
3250        fill_table_options_for_create(&mut table_options, &AutoCreateTableType::LastNonNull, &ctx);
3251
3252        assert!(!table_options.contains_key(APPEND_MODE_KEY));
3253        assert_eq!(
3254            Some("last_non_null"),
3255            table_options.get(MERGE_MODE_KEY).map(String::as_str)
3256        );
3257    }
3258
3259    #[test]
3260    fn test_last_non_null_create_options_use_configured_merge_mode() {
3261        let mut ctx = QueryContext::with(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME);
3262        ctx.set_extension(MERGE_MODE_KEY, "last_row");
3263        let ctx = Arc::new(ctx);
3264        let mut table_options = Default::default();
3265
3266        fill_table_options_for_create(&mut table_options, &AutoCreateTableType::LastNonNull, &ctx);
3267
3268        assert_eq!(
3269            Some("last_row"),
3270            table_options.get(MERGE_MODE_KEY).map(String::as_str)
3271        );
3272        assert!(!table_options.contains_key(APPEND_MODE_KEY));
3273    }
3274
3275    #[test]
3276    fn test_last_non_null_create_options_use_last_row_with_append_mode_true() {
3277        let mut ctx = QueryContext::with(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME);
3278        ctx.set_extension(APPEND_MODE_KEY, "true");
3279        let ctx = Arc::new(ctx);
3280        let mut table_options = Default::default();
3281
3282        fill_table_options_for_create(&mut table_options, &AutoCreateTableType::LastNonNull, &ctx);
3283
3284        assert_eq!(
3285            Some("true"),
3286            table_options.get(APPEND_MODE_KEY).map(String::as_str)
3287        );
3288        assert_eq!(
3289            Some("last_row"),
3290            table_options.get(MERGE_MODE_KEY).map(String::as_str)
3291        );
3292    }
3293
3294    struct UnexpectedBatcher;
3295
3296    #[async_trait::async_trait]
3297    impl PendingRowsBatcher for UnexpectedBatcher {
3298        async fn acquire(&self) -> Result<Arc<tokio::sync::OwnedSemaphorePermit>> {
3299            panic!("empty writes must not acquire batch admission")
3300        }
3301
3302        async fn submit(
3303            &self,
3304            _table_info: TableInfoRef,
3305            _batch: arrow::record_batch::RecordBatch,
3306            _ctx: QueryContextRef,
3307            _permit: Arc<tokio::sync::OwnedSemaphorePermit>,
3308        ) -> Result<Output> {
3309            panic!("empty writes must not submit a batch")
3310        }
3311    }
3312
3313    async fn batcher_test_inserter() -> Inserter {
3314        let kv_backend = prepare_mocked_backend().await;
3315        Inserter::new(
3316            catalog::memory::MemoryCatalogManager::new(),
3317            create_partition_rule_manager(kv_backend.clone()).await,
3318            Arc::new(MockDatanodeManager::new(NaiveDatanodeHandler)),
3319            Arc::new(new_table_flownode_set_cache(
3320                String::new(),
3321                Cache::new(100),
3322                kv_backend,
3323            )),
3324            true,
3325        )
3326    }
3327
3328    #[tokio::test]
3329    async fn test_logical_batcher_eligibility() {
3330        use catalog::RegisterTableRequest;
3331        use catalog::memory::MemoryCatalogManager;
3332        use common_meta::instruction::{CacheIdent, CreateFlow};
3333        use common_meta::kv_backend::KvBackendRef;
3334        use common_meta::kv_backend::memory::MemoryKvBackend;
3335        use datatypes::schema::{ColumnDefaultConstraint, SchemaBuilder};
3336        let requests = RowInsertRequests {
3337            inserts: vec![RowInsertRequest {
3338                table_name: "test_table".to_string(),
3339                rows: None,
3340            }],
3341        };
3342        let original =
3343            make_table_ref_with_schema("ts", "value", ConcreteDataType::float64_datatype())
3344                .table_info();
3345        for case in [
3346            "eligible",
3347            "physical",
3348            "ordinary",
3349            "instant",
3350            "disabled",
3351            "hint",
3352            "flow",
3353            "required_tag",
3354            "default",
3355        ] {
3356            let mut info = (*original).clone();
3357            info.meta.engine = METRIC_ENGINE_NAME.to_string();
3358            info.meta.options.extra_options.insert(
3359                LOGICAL_TABLE_METADATA_KEY.to_string(),
3360                "physical".to_string(),
3361            );
3362            let mut ctx = QueryContext::arc().fork();
3363            match case {
3364                "physical" => {
3365                    info.meta
3366                        .options
3367                        .extra_options
3368                        .insert(LOGICAL_TABLE_METADATA_KEY.to_string(), "other".to_string());
3369                }
3370                "ordinary" => info.meta.engine = "mito".to_string(),
3371                "instant" => info.meta.options.ttl = Some(common_time::ttl::TimeToLive::Instant),
3372                "hint" => ctx.set_extension(AUTO_CREATE_TABLE_KEY, "false"),
3373                "required_tag" | "default" => {
3374                    let mut columns = info.meta.schema.column_schemas().to_vec();
3375                    let mut tag = ColumnSchema::new(
3376                        "tag",
3377                        ConcreteDataType::string_datatype(),
3378                        case != "required_tag",
3379                    );
3380                    if case == "default" {
3381                        tag = tag
3382                            .with_default_constraint(Some(ColumnDefaultConstraint::null_value()))
3383                            .unwrap();
3384                    }
3385                    columns.push(tag);
3386                    info.meta.schema = Arc::new(
3387                        SchemaBuilder::try_from_columns(columns)
3388                            .unwrap()
3389                            .build()
3390                            .unwrap(),
3391                    );
3392                    info.meta.primary_key_indices = vec![2];
3393                }
3394                _ => {}
3395            }
3396            let catalog = MemoryCatalogManager::with_default_setup();
3397            let table = Arc::new(table::Table::new(
3398                Arc::new(info),
3399                table::metadata::FilterPushDownType::Unsupported,
3400                Arc::new(DummyDataSource),
3401            ));
3402            catalog
3403                .register_table_sync(RegisterTableRequest {
3404                    catalog: DEFAULT_CATALOG_NAME.to_string(),
3405                    schema: DEFAULT_SCHEMA_NAME.to_string(),
3406                    table_name: "test_table".to_string(),
3407                    table_id: 1,
3408                    table,
3409                })
3410                .unwrap();
3411            let mut inserter = batcher_test_inserter().await;
3412            inserter.catalog_manager = catalog;
3413            inserter.auto_create_table = case != "disabled";
3414            let kv_backend: KvBackendRef = Arc::new(MemoryKvBackend::default());
3415            inserter.table_flownode_set_cache = Arc::new(new_table_flownode_set_cache(
3416                String::new(),
3417                Cache::new(10),
3418                kv_backend,
3419            ));
3420            if case == "flow" {
3421                inserter
3422                    .table_flownode_set_cache
3423                    .invalidate(&[CacheIdent::CreateFlow(CreateFlow {
3424                        flow_id: 1,
3425                        source_table_ids: vec![1],
3426                        partition_to_peer_mapping: vec![(0, Peer::empty(1))],
3427                    })])
3428                    .await
3429                    .unwrap();
3430            }
3431            assert_eq!(
3432                inserter
3433                    .can_batch_metric_rows(&requests, &Arc::new(ctx), "physical")
3434                    .await
3435                    .unwrap(),
3436                case == "eligible",
3437                "{case}"
3438            );
3439        }
3440    }
3441
3442    #[tokio::test]
3443    async fn test_batcher_meter_preserves_request() {
3444        let mut requests = RowInsertRequests {
3445            inserts: vec![RowInsertRequest {
3446                table_name: "sample".to_string(),
3447                rows: Some(Rows {
3448                    schema: vec![],
3449                    rows: vec![api::v1::Row {
3450                        values: vec![Value {
3451                            value_data: Some(api::v1::value::ValueData::F64Value(1.5)),
3452                        }],
3453                    }],
3454                }),
3455            }],
3456        };
3457        let expected = requests.clone();
3458        Inserter::meter_row_inserts(&mut requests, &QueryContext::arc())
3459            .await
3460            .unwrap();
3461        assert_eq!(requests, expected);
3462    }
3463
3464    #[tokio::test]
3465    async fn test_instant_table_bypasses_batcher() {
3466        let batcher: Arc<dyn PendingRowsBatcher> = Arc::new(UnexpectedBatcher);
3467        let inserter = batcher_test_inserter()
3468            .await
3469            .with_pending_rows_batcher(Some(batcher));
3470        let mut ctx = session::context::QueryContextBuilder::default().build();
3471        ctx.set_batching_enabled(true);
3472        let ctx = Arc::new(ctx);
3473        let table = make_table_ref_with_schema("ts", "value", ConcreteDataType::float64_datatype())
3474            .table_info();
3475        assert!(inserter.table_batcher(&table, &ctx).is_some());
3476        let mut instant = (*table).clone();
3477        instant.meta.options.ttl = Some(common_time::ttl::TimeToLive::Instant);
3478        assert!(inserter.table_batcher(&Arc::new(instant), &ctx).is_none());
3479    }
3480
3481    #[tokio::test]
3482    async fn test_empty_prepared_rows_skip_batcher() {
3483        let inserter = batcher_test_inserter().await;
3484        let table = make_table_ref_with_schema("ts", "value", ConcreteDataType::float64_datatype())
3485            .table_info();
3486        let batcher: Arc<dyn PendingRowsBatcher> = Arc::new(UnexpectedBatcher);
3487        let ctx = QueryContext::arc();
3488        let output = inserter
3489            .submit_table_rows(
3490                Rows {
3491                    schema: vec![],
3492                    rows: vec![],
3493                },
3494                table.clone(),
3495                ctx.clone(),
3496                &batcher,
3497            )
3498            .await
3499            .unwrap();
3500        assert!(matches!(output.data, OutputData::AffectedRows(0)));
3501        let output = inserter
3502            .submit_pending_rows(
3503                RowInsertRequests {
3504                    inserts: vec![RowInsertRequest {
3505                        table_name: table.name.clone(),
3506                        rows: None,
3507                    }],
3508                },
3509                HashMap::from_iter([(table.table_id(), table)]),
3510                ctx,
3511                &batcher,
3512            )
3513            .await
3514            .unwrap();
3515        assert!(matches!(output.data, OutputData::AffectedRows(0)));
3516    }
3517
3518    /// Flownode handler that keeps mirror inserts in flight until the test
3519    /// releases the gate, so the pending mirror budget can be observed
3520    /// deterministically.
3521    #[derive(Clone)]
3522    struct GatedFlownodeHandler {
3523        gate: Arc<tokio::sync::Semaphore>,
3524        affected_rows: u64,
3525    }
3526
3527    #[async_trait::async_trait]
3528    impl MockFlownodeHandler for GatedFlownodeHandler {
3529        async fn handle_inserts(
3530            &self,
3531            _peer: &Peer,
3532            _requests: api::v1::region::InsertRequests,
3533        ) -> common_meta::error::Result<FlowResponse> {
3534            let _permit = self.gate.acquire().await.unwrap();
3535            Ok(FlowResponse {
3536                affected_rows: self.affected_rows,
3537                ..Default::default()
3538            })
3539        }
3540    }
3541
3542    fn flownode_peer() -> Peer {
3543        Peer {
3544            id: 1,
3545            addr: "127.0.0.1:4001".to_string(),
3546        }
3547    }
3548
3549    #[derive(Clone)]
3550    struct RecordingNodeManager {
3551        datanode_dispatch: tokio::sync::mpsc::UnboundedSender<usize>,
3552        flownode_dispatch: tokio::sync::mpsc::UnboundedSender<Peer>,
3553        gates: Arc<std::sync::Mutex<HashMap<u64, Arc<tokio::sync::Semaphore>>>>,
3554        fail_flownode: bool,
3555    }
3556
3557    #[async_trait::async_trait]
3558    impl common_meta::node_manager::DatanodeManager for RecordingNodeManager {
3559        async fn datanode(&self, peer: &Peer) -> common_meta::node_manager::DatanodeRef {
3560            Arc::new(RecordingNode {
3561                peer: peer.clone(),
3562                manager: self.clone(),
3563            })
3564        }
3565    }
3566
3567    #[async_trait::async_trait]
3568    impl common_meta::node_manager::FlownodeManager for RecordingNodeManager {
3569        async fn flownode(&self, peer: &Peer) -> common_meta::node_manager::FlownodeRef {
3570            Arc::new(RecordingNode {
3571                peer: peer.clone(),
3572                manager: self.clone(),
3573            })
3574        }
3575    }
3576
3577    struct RecordingNode {
3578        peer: Peer,
3579        manager: RecordingNodeManager,
3580    }
3581
3582    #[async_trait::async_trait]
3583    impl common_meta::node_manager::Datanode for RecordingNode {
3584        async fn handle(
3585            &self,
3586            _request: api::v1::region::RegionRequest,
3587        ) -> common_meta::error::Result<api::region::RegionResponse> {
3588            let rows = match _request.body.as_ref() {
3589                Some(api::v1::region::region_request::Body::Inserts(requests)) => requests
3590                    .requests
3591                    .iter()
3592                    .filter_map(|request| request.rows.as_ref())
3593                    .map(|rows| rows.rows.len())
3594                    .sum(),
3595                _ => 0,
3596            };
3597            let _ = self.manager.datanode_dispatch.send(rows);
3598            Ok(api::region::RegionResponse::new(rows))
3599        }
3600
3601        async fn handle_query(
3602            &self,
3603            _request: common_query::request::QueryRequest,
3604        ) -> common_meta::error::Result<common_recordbatch::SendableRecordBatchStream> {
3605            unreachable!()
3606        }
3607    }
3608
3609    #[async_trait::async_trait]
3610    impl common_meta::node_manager::Flownode for RecordingNode {
3611        async fn handle(
3612            &self,
3613            _request: api::v1::flow::FlowRequest,
3614        ) -> common_meta::error::Result<FlowResponse> {
3615            unreachable!()
3616        }
3617
3618        async fn handle_inserts(
3619            &self,
3620            _request: api::v1::region::InsertRequests,
3621        ) -> common_meta::error::Result<FlowResponse> {
3622            let _ = self.manager.flownode_dispatch.send(self.peer.clone());
3623            if self.manager.fail_flownode {
3624                return Err(common_meta::error::UnexpectedSnafu {
3625                    err_msg: "test flownode failure".to_string(),
3626                }
3627                .build());
3628            }
3629            let gate = self
3630                .manager
3631                .gates
3632                .lock()
3633                .unwrap()
3634                .get(&self.peer.id)
3635                .cloned();
3636            if let Some(gate) = gate {
3637                let _permit = gate.acquire().await.unwrap();
3638            }
3639            Ok(FlowResponse::default())
3640        }
3641
3642        async fn handle_mark_window_dirty(
3643            &self,
3644            _request: api::v1::flow::DirtyWindowRequests,
3645        ) -> common_meta::error::Result<FlowResponse> {
3646            unreachable!()
3647        }
3648    }
3649
3650    fn mirror_requests(peer: &Peer, num_rows: usize) -> HashMap<Peer, RegionInsertRequests> {
3651        HashMap::from_iter([(
3652            peer.clone(),
3653            RegionInsertRequests {
3654                requests: vec![RegionInsertRequest {
3655                    region_id: RegionId::new(1, 1).as_u64(),
3656                    rows: Some(Rows {
3657                        schema: vec![],
3658                        rows: vec![api::v1::Row { values: vec![] }; num_rows],
3659                    }),
3660                    ..Default::default()
3661                }],
3662            },
3663        )])
3664    }
3665
3666    /// Serializes the mirror metric tests below: the metrics are process-wide
3667    /// singletons, so these tests would otherwise observe each other's updates
3668    /// when the test binary runs tests in parallel (they are isolated per
3669    /// process only under nextest).
3670    static MIRROR_METRIC_TEST_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
3671
3672    async fn lock_mirror_metrics() -> tokio::sync::MutexGuard<'static, ()> {
3673        MIRROR_METRIC_TEST_LOCK.lock().await
3674    }
3675
3676    // The tests below share the process-wide mirror metrics, so they take
3677    // `MIRROR_METRIC_TEST_LOCK` and assert on the difference their own budget
3678    // publishes instead of on absolute values.
3679    #[tokio::test]
3680    async fn flow_mirror_saturation_preserves_instant_and_normal_contracts() {
3681        use common_error::ext::{ErrorExt, RetryHint};
3682        use common_error::status_code::StatusCode;
3683
3684        let _guard = lock_mirror_metrics().await;
3685        let kv_backend = prepare_mocked_backend().await;
3686        let partition_manager = create_partition_rule_manager(kv_backend.clone()).await;
3687        let (datanode_tx, mut datanode_rx) = tokio::sync::mpsc::unbounded_channel();
3688        let (flownode_tx, mut flownode_rx) = tokio::sync::mpsc::unbounded_channel();
3689        let node_manager = Arc::new(RecordingNodeManager {
3690            datanode_dispatch: datanode_tx,
3691            flownode_dispatch: flownode_tx,
3692            gates: Arc::new(std::sync::Mutex::new(HashMap::new())),
3693            fail_flownode: false,
3694        });
3695        let flow_cache = Cache::new(10);
3696        let flow_cache_backend = prepare_mocked_backend().await;
3697        let inserter = Inserter::new(
3698            catalog::memory::MemoryCatalogManager::new(),
3699            partition_manager,
3700            node_manager,
3701            Arc::new(new_table_flownode_set_cache(
3702                String::new(),
3703                flow_cache.clone(),
3704                flow_cache_backend,
3705            )),
3706            true,
3707        );
3708        let peer = flownode_peer();
3709        inserter
3710            .table_flownode_set_cache
3711            .invalidate(&[common_meta::instruction::CacheIdent::CreateFlow(
3712                common_meta::instruction::CreateFlow {
3713                    flow_id: 1,
3714                    source_table_ids: vec![1, 2],
3715                    partition_to_peer_mapping: vec![(0, peer)],
3716                },
3717            )])
3718            .await
3719            .unwrap();
3720        let mut normal_info = new_test_table_info(1, "normal_table", [1].into_iter());
3721        normal_info.catalog_name = DEFAULT_CATALOG_NAME.to_string();
3722        normal_info.schema_name = DEFAULT_SCHEMA_NAME.to_string();
3723        let mut instant_info = new_test_table_info(2, "instant_table", [1].into_iter());
3724        instant_info.meta.options.ttl = Some(common_time::ttl::TimeToLive::Instant);
3725        let table_infos =
3726            HashMap::from_iter([(1, Arc::new(normal_info)), (2, Arc::new(instant_info))]);
3727        let ctx = Arc::new(QueryContext::with(
3728            DEFAULT_CATALOG_NAME,
3729            DEFAULT_SCHEMA_NAME,
3730        ));
3731        let insert = |table_id, num_rows| RegionInsertRequests {
3732            requests: vec![RegionInsertRequest {
3733                region_id: RegionId::new(table_id, 1).as_u64(),
3734                rows: Some(Rows {
3735                    schema: vec![],
3736                    rows: vec![api::v1::Row::default(); num_rows],
3737                }),
3738                ..Default::default()
3739            }],
3740        };
3741        let make_request = |normal, instant| InstantAndNormalInsertRequests {
3742            normal_requests: if normal == 0 {
3743                RegionInsertRequests::default()
3744            } else {
3745                insert(1, normal)
3746            },
3747            instant_requests: if instant == 0 {
3748                RegionInsertRequests::default()
3749            } else {
3750                insert(2, instant)
3751            },
3752        };
3753
3754        let pending = inserter.mirror_pending_rows.clone();
3755        let gauge_before = crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.get();
3756        crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.add(MAX_MIRROR_PENDING_ROWS as i64);
3757        pending.store(MAX_MIRROR_PENDING_ROWS as u64, Ordering::Relaxed);
3758        let dropped_before = crate::metrics::DIST_MIRROR_DROPPED_ROW_COUNT.get();
3759        for request in [make_request(0, 1), make_request(1, 1)] {
3760            let error = inserter
3761                .do_request(request, &table_infos, &ctx)
3762                .await
3763                .unwrap_err();
3764            assert_eq!(error.status_code(), StatusCode::RateLimited);
3765            assert_eq!(error.retry_hint(), RetryHint::Retryable);
3766            assert!(datanode_rx.try_recv().is_err());
3767            assert!(flownode_rx.try_recv().is_err());
3768            assert_eq!(
3769                pending.load(Ordering::Relaxed),
3770                MAX_MIRROR_PENDING_ROWS as u64
3771            );
3772        }
3773        assert_eq!(
3774            crate::metrics::DIST_MIRROR_DROPPED_ROW_COUNT.get(),
3775            dropped_before
3776        );
3777
3778        for pending_before in [0, MAX_MIRROR_PENDING_ROWS as u64] {
3779            if pending_before == 0 {
3780                pending.store(0, Ordering::Relaxed);
3781                crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.sub(MAX_MIRROR_PENDING_ROWS as i64);
3782            }
3783            let gauge_before_rejection = crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.get();
3784            for (normal, instant) in [
3785                (0, MAX_MIRROR_PENDING_ROWS + 1),
3786                (MAX_MIRROR_PENDING_ROWS / 2, MAX_MIRROR_PENDING_ROWS / 2 + 1),
3787            ] {
3788                let error = inserter
3789                    .do_request(make_request(normal, instant), &table_infos, &ctx)
3790                    .await
3791                    .unwrap_err();
3792                assert_eq!(error.status_code(), StatusCode::InvalidArguments);
3793                assert_eq!(error.retry_hint(), RetryHint::NonRetryable);
3794                let message = error.to_string();
3795                assert!(message.contains("1000001"), "{message}");
3796                assert!(message.contains("1000000"), "{message}");
3797                assert!(message.contains("reduce the batch size"), "{message}");
3798                assert!(datanode_rx.try_recv().is_err());
3799                assert!(flownode_rx.try_recv().is_err());
3800                assert_eq!(pending.load(Ordering::Relaxed), pending_before);
3801                assert_eq!(
3802                    crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.get(),
3803                    gauge_before_rejection
3804                );
3805                assert_eq!(
3806                    crate::metrics::DIST_MIRROR_DROPPED_ROW_COUNT.get(),
3807                    dropped_before
3808                );
3809            }
3810            if pending_before == 0 {
3811                crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.add(MAX_MIRROR_PENDING_ROWS as i64);
3812                pending.store(MAX_MIRROR_PENDING_ROWS as u64, Ordering::Relaxed);
3813            }
3814        }
3815
3816        let result = inserter
3817            .do_request(make_request(1, 0), &table_infos, &ctx)
3818            .await
3819            .unwrap();
3820        assert!(matches!(result.data, OutputData::AffectedRows(1)));
3821        assert_eq!(1, datanode_rx.try_recv().unwrap());
3822        assert!(flownode_rx.try_recv().is_err());
3823        assert_eq!(
3824            crate::metrics::DIST_MIRROR_DROPPED_ROW_COUNT.get(),
3825            dropped_before + 1
3826        );
3827        pending.store(0, Ordering::Relaxed);
3828        crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.sub(MAX_MIRROR_PENDING_ROWS as i64);
3829        assert_eq!(
3830            crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.get(),
3831            gauge_before
3832        );
3833    }
3834
3835    #[tokio::test]
3836    async fn flow_mirror_new_counts_regions_and_peer_clones() {
3837        use common_meta::instruction::{CacheIdent, CreateFlow};
3838
3839        let peer1 = flownode_peer();
3840        let peer2 = Peer {
3841            id: 2,
3842            addr: "127.0.0.1:4002".to_string(),
3843        };
3844        let cache_backend = prepare_mocked_backend().await;
3845        let cache = Arc::new(new_table_flownode_set_cache(
3846            String::new(),
3847            Cache::new(10),
3848            cache_backend,
3849        ));
3850        cache
3851            .invalidate(&[CacheIdent::CreateFlow(CreateFlow {
3852                flow_id: 1,
3853                source_table_ids: vec![1],
3854                partition_to_peer_mapping: vec![(0, peer1), (1, peer2)],
3855            })])
3856            .await
3857            .unwrap();
3858        let request = |region_id| RegionInsertRequest {
3859            region_id,
3860            rows: Some(Rows {
3861                schema: vec![],
3862                rows: vec![api::v1::Row { values: vec![] }; 3],
3863            }),
3864            ..Default::default()
3865        };
3866        let requests = [
3867            request(RegionId::new(1, 0).as_u64()),
3868            request(RegionId::new(1, 1).as_u64()),
3869        ];
3870        let task = FlowMirrorTask::new(&cache, requests.iter()).await.unwrap();
3871        assert_eq!(12, task.pending_rows());
3872        assert_eq!(2, task.requests.len());
3873    }
3874
3875    #[tokio::test]
3876    async fn flow_mirror_dropped_when_pending_exceeds_limit() {
3877        let _guard = lock_mirror_metrics().await;
3878        let pending = Arc::new(AtomicU64::new(MAX_MIRROR_PENDING_ROWS as u64));
3879        let dropped_before = crate::metrics::DIST_MIRROR_DROPPED_ROW_COUNT.get();
3880        let gauge_before = crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.get();
3881        let num_rows = 100;
3882
3883        let task = FlowMirrorTask {
3884            requests: mirror_requests(&flownode_peer(), num_rows),
3885        };
3886        // The budget is already full, so nothing may be spawned for this batch.
3887        task.detach(
3888            Arc::new(MockDatanodeManager::new(NaiveDatanodeHandler)),
3889            pending.clone(),
3890            false,
3891        )
3892        .unwrap();
3893
3894        assert_eq!(
3895            num_rows as u64,
3896            crate::metrics::DIST_MIRROR_DROPPED_ROW_COUNT.get() - dropped_before
3897        );
3898        // The dropped batch must not leak its reservation.
3899        assert_eq!(
3900            MAX_MIRROR_PENDING_ROWS as u64,
3901            pending.load(Ordering::Relaxed)
3902        );
3903        // The reservation taken before the limit check is released again, so the
3904        // gauge nets out and keeps matching the pending budget.
3905        assert_eq!(
3906            gauge_before,
3907            crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.get(),
3908            "a dropped batch must not change the pending gauge"
3909        );
3910
3911        let task = FlowMirrorTask {
3912            requests: mirror_requests(&flownode_peer(), num_rows),
3913        };
3914        let error = task
3915            .detach(
3916                Arc::new(MockDatanodeManager::new(NaiveDatanodeHandler)),
3917                pending.clone(),
3918                true,
3919            )
3920            .err()
3921            .unwrap();
3922        use common_error::ext::{ErrorExt, RetryHint};
3923        use common_error::status_code::StatusCode;
3924        assert_eq!(error.status_code(), StatusCode::RateLimited);
3925        assert_eq!(error.retry_hint(), RetryHint::Retryable);
3926        assert_eq!(
3927            crate::metrics::DIST_MIRROR_DROPPED_ROW_COUNT.get(),
3928            dropped_before + num_rows as u64,
3929            "rejected instant rows are not successful best-effort drops"
3930        );
3931        assert_eq!(
3932            MAX_MIRROR_PENDING_ROWS as u64,
3933            pending.load(Ordering::Relaxed)
3934        );
3935        assert_eq!(
3936            gauge_before,
3937            crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.get()
3938        );
3939    }
3940
3941    /// A source table mapped to several flownodes clones its requests per peer,
3942    /// so the reservation must cover every clone and each spawned task must hand
3943    /// back exactly its own share; otherwise the shared budget drifts.
3944    #[tokio::test]
3945    async fn flow_mirror_reserves_and_releases_per_peer_shares() {
3946        let _guard = lock_mirror_metrics().await;
3947        let per_peer_rows = 10;
3948        let total_rows = (2 * per_peer_rows) as u64;
3949        // Exactly enough room for the whole batch, cloned rows included.
3950        let start = MAX_MIRROR_PENDING_ROWS as u64 - total_rows;
3951        let pending = Arc::new(AtomicU64::new(start));
3952        let gate = Arc::new(tokio::sync::Semaphore::new(0));
3953        let node_manager = Arc::new(MockFlownodeManager::new(GatedFlownodeHandler {
3954            gate: gate.clone(),
3955            affected_rows: per_peer_rows as u64,
3956        }));
3957
3958        let mut requests = mirror_requests(&flownode_peer(), per_peer_rows);
3959        requests.extend(mirror_requests(
3960            &Peer {
3961                id: 2,
3962                addr: "127.0.0.1:4002".to_string(),
3963            },
3964            per_peer_rows,
3965        ));
3966
3967        let gauge_before = crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.get();
3968        FlowMirrorTask { requests }
3969            .detach(node_manager, pending.clone(), false)
3970            .unwrap();
3971
3972        // Both per-peer payloads are counted, so the budget is exactly full
3973        // instead of over-reserved by the duplicated rows.
3974        assert_eq!(
3975            MAX_MIRROR_PENDING_ROWS as u64,
3976            pending.load(Ordering::Relaxed)
3977        );
3978
3979        // Let both tasks finish and check that the budget lands back on its
3980        // starting value: over-releasing a shared reservation would saturate it.
3981        // The gauge is checked too because `release_mirror_pending_rows` moves
3982        // the budget before the gauge; exiting on the budget alone could leak a
3983        // pending `gauge.sub` into the next test holding the metrics lock.
3984        gate.add_permits(2);
3985        let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10);
3986        while pending.load(Ordering::Relaxed) != start
3987            || crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.get() != gauge_before
3988        {
3989            assert!(
3990                std::time::Instant::now() < deadline,
3991                "the mirror tasks must release exactly their own shares"
3992            );
3993            std::thread::sleep(std::time::Duration::from_millis(1));
3994        }
3995    }
3996
3997    #[tokio::test]
3998    async fn flow_mirror_fast_peer_releases_only_its_share_before_drop() {
3999        let _guard = lock_mirror_metrics().await;
4000        let peer1 = flownode_peer();
4001        let peer2 = Peer {
4002            id: 2,
4003            addr: "127.0.0.1:4002".to_string(),
4004        };
4005        let gate1 = Arc::new(tokio::sync::Semaphore::new(0));
4006        let gate2 = Arc::new(tokio::sync::Semaphore::new(0));
4007        let gates = Arc::new(std::sync::Mutex::new(HashMap::from_iter([
4008            (peer1.id, gate1.clone()),
4009            (peer2.id, gate2.clone()),
4010        ])));
4011        let (datanode_tx, _datanode_rx) = tokio::sync::mpsc::unbounded_channel();
4012        let (flownode_tx, mut flownode_rx) = tokio::sync::mpsc::unbounded_channel();
4013        let node_manager = Arc::new(RecordingNodeManager {
4014            datanode_dispatch: datanode_tx,
4015            flownode_dispatch: flownode_tx,
4016            gates,
4017            fail_flownode: false,
4018        });
4019        let pending = Arc::new(AtomicU64::new(MAX_MIRROR_PENDING_ROWS as u64 - 20));
4020        let gauge_before = crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.get();
4021        let mut requests = mirror_requests(&peer1, 10);
4022        requests.extend(mirror_requests(&peer2, 10));
4023        let task = FlowMirrorTask { requests };
4024        task.detach(node_manager.clone(), pending.clone(), false)
4025            .unwrap();
4026        let first = tokio::time::timeout(Duration::from_secs(10), flownode_rx.recv())
4027            .await
4028            .unwrap()
4029            .unwrap();
4030        let second = tokio::time::timeout(Duration::from_secs(10), flownode_rx.recv())
4031            .await
4032            .unwrap()
4033            .unwrap();
4034        assert_ne!(first.id, second.id);
4035        let fast_gate = if first.id == peer1.id {
4036            gate1.clone()
4037        } else {
4038            gate2.clone()
4039        };
4040        let slow_gate = if first.id == peer1.id {
4041            gate2.clone()
4042        } else {
4043            gate1.clone()
4044        };
4045        fast_gate.add_permits(1);
4046        let deadline = std::time::Instant::now() + Duration::from_secs(10);
4047        while pending.load(Ordering::Relaxed) != MAX_MIRROR_PENDING_ROWS as u64 - 10 {
4048            assert!(
4049                std::time::Instant::now() < deadline,
4050                "fast peer did not release"
4051            );
4052            std::thread::sleep(Duration::from_millis(1));
4053        }
4054        let dropped_before = crate::metrics::DIST_MIRROR_DROPPED_ROW_COUNT.get();
4055        FlowMirrorTask {
4056            requests: mirror_requests(&peer1, 11),
4057        }
4058        .detach(node_manager, pending.clone(), false)
4059        .unwrap();
4060        assert_eq!(
4061            crate::metrics::DIST_MIRROR_DROPPED_ROW_COUNT.get(),
4062            dropped_before + 11
4063        );
4064        assert!(flownode_rx.try_recv().is_err());
4065        slow_gate.add_permits(1);
4066        while pending.load(Ordering::Relaxed) != MAX_MIRROR_PENDING_ROWS as u64 - 20
4067            || crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.get() != gauge_before
4068        {
4069            assert!(
4070                std::time::Instant::now() < deadline,
4071                "slow peer did not release"
4072            );
4073            std::thread::sleep(Duration::from_millis(1));
4074        }
4075    }
4076
4077    #[tokio::test]
4078    async fn flow_mirror_failure_releases_budget_and_gauge() {
4079        let _guard = lock_mirror_metrics().await;
4080        let (datanode_tx, _datanode_rx) = tokio::sync::mpsc::unbounded_channel();
4081        let (flownode_tx, mut flownode_rx) = tokio::sync::mpsc::unbounded_channel();
4082        let node_manager = Arc::new(RecordingNodeManager {
4083            datanode_dispatch: datanode_tx,
4084            flownode_dispatch: flownode_tx,
4085            gates: Arc::new(std::sync::Mutex::new(HashMap::new())),
4086            fail_flownode: true,
4087        });
4088        let pending = Arc::new(AtomicU64::new(0));
4089        let gauge_before = crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.get();
4090        FlowMirrorTask {
4091            requests: mirror_requests(&flownode_peer(), 5),
4092        }
4093        .detach(node_manager, pending.clone(), false)
4094        .unwrap();
4095        assert_eq!(
4096            tokio::time::timeout(Duration::from_secs(10), flownode_rx.recv())
4097                .await
4098                .unwrap()
4099                .unwrap()
4100                .id,
4101            flownode_peer().id
4102        );
4103        let deadline = std::time::Instant::now() + Duration::from_secs(10);
4104        while pending.load(Ordering::Relaxed) != 0
4105            || crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.get() != gauge_before
4106        {
4107            assert!(
4108                std::time::Instant::now() < deadline,
4109                "failed task leaked reservation"
4110            );
4111            std::thread::sleep(Duration::from_millis(1));
4112        }
4113    }
4114
4115    #[tokio::test]
4116    async fn flow_mirror_pending_counter_consistent() {
4117        let _guard = lock_mirror_metrics().await;
4118        let pending = Arc::new(AtomicU64::new(0));
4119        let gate = Arc::new(tokio::sync::Semaphore::new(0));
4120        let node_manager = Arc::new(MockFlownodeManager::new(GatedFlownodeHandler {
4121            gate: gate.clone(),
4122            affected_rows: 7,
4123        }));
4124
4125        // A batch without rows short-circuits and leaves the budget untouched.
4126        let gauge_before = crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.get();
4127        FlowMirrorTask {
4128            requests: HashMap::new(),
4129        }
4130        .detach(node_manager.clone(), pending.clone(), false)
4131        .unwrap();
4132        assert_eq!(0, pending.load(Ordering::Relaxed));
4133        assert_eq!(
4134            gauge_before,
4135            crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.get(),
4136            "an empty batch must not touch the pending gauge"
4137        );
4138
4139        let dropped_before = crate::metrics::DIST_MIRROR_DROPPED_ROW_COUNT.get();
4140        let num_rows = 7;
4141
4142        FlowMirrorTask {
4143            requests: mirror_requests(&flownode_peer(), num_rows),
4144        }
4145        .detach(node_manager, pending.clone(), false)
4146        .unwrap();
4147
4148        // The spawned task is still waiting on the gated flownode, so the
4149        // reservation made by `detach` is the value observable here.
4150        assert_eq!(num_rows as u64, pending.load(Ordering::Relaxed));
4151        assert_eq!(
4152            num_rows as i64,
4153            crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.get() - gauge_before,
4154            "the pending gauge must track the pending rows"
4155        );
4156        assert_eq!(
4157            dropped_before,
4158            crate::metrics::DIST_MIRROR_DROPPED_ROW_COUNT.get(),
4159            "a batch within the limit must not be counted as dropped"
4160        );
4161
4162        // Release the in-flight task and wait for it to hand its reservation
4163        // back, so no task is left holding the shared metrics after this test.
4164        gate.add_permits(1);
4165        let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10);
4166        while pending.load(Ordering::Relaxed) != 0
4167            || crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.get() != gauge_before
4168        {
4169            assert!(
4170                std::time::Instant::now() < deadline,
4171                "the mirror task did not release its pending reservation"
4172            );
4173            std::thread::sleep(std::time::Duration::from_millis(1));
4174        }
4175    }
4176
4177    #[test]
4178    fn mirror_logs_are_rate_limited() {
4179        let log = Arc::new(MirrorLog::new());
4180        assert_eq!(Some(1), log.claim_report(1_000));
4181        assert_eq!(None, log.claim_report(1_500));
4182        assert_eq!(None, log.claim_report(10_999));
4183        assert_eq!(Some(3), log.claim_report(11_000));
4184        assert_eq!(None, log.claim_report(11_001));
4185        assert_eq!(Some(2), log.claim_report(21_000));
4186
4187        let concurrent = Arc::new(MirrorLog::new());
4188        assert_eq!(Some(1), concurrent.claim_report(50_000));
4189        let barrier = Arc::new(std::sync::Barrier::new(8));
4190        let threads = (0..8)
4191            .map(|_| {
4192                let log = concurrent.clone();
4193                let barrier = barrier.clone();
4194                std::thread::spawn(move || {
4195                    barrier.wait();
4196                    log.claim_report(60_000)
4197                })
4198            })
4199            .collect::<Vec<_>>();
4200        let winners = threads
4201            .into_iter()
4202            .filter_map(|thread| thread.join().unwrap())
4203            .collect::<Vec<_>>();
4204        assert_eq!(winners.len(), 1);
4205        let residual = concurrent.events.swap(0, Ordering::Relaxed);
4206        assert_eq!(winners[0] + residual, 8);
4207    }
4208}