1use 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 auto_create_table: bool,
108 pending_rows_batcher: Option<Arc<dyn PendingRowsBatcher>>,
109 mirror_pending_rows: Arc<AtomicU64>,
113}
114
115pub type InserterRef = Arc<Inserter>;
116
117#[derive(Clone)]
119pub enum AutoCreateTableType {
120 Logical(String),
122 Physical,
124 Log,
126 LastNonNull,
128 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#[derive(Clone)]
160pub struct InstantAndNormalInsertRequests {
161 pub normal_requests: RegionInsertRequests,
163 pub instant_requests: RegionInsertRequests,
166}
167
168impl Inserter {
169 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 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 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 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 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 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 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 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 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 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 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 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 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 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 let permit = batcher.acquire().await?;
548 let submissions = prepared.into_iter().map(|(info, batch)| {
549 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 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 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 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 let physical_table_ref = self
589 .create_physical_table_on_demand(&ctx, physical_table.clone(), statement_executor)
590 .await?;
591
592 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
706pub async fn admit_write(rows: u64, ctx: &QueryContextRef) -> Result<QueryContextRef> {
709 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
723pub 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 let requests = fill_reqs_with_impure_default(table_infos, requests)?;
781
782 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 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 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 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 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 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 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 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 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 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 #[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 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 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 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 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 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 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 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 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 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 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 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 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 create_table
1305 .table_options
1306 .insert(APPEND_MODE_KEY.to_string(), "false".to_string());
1307 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 let partitions = if matches!(trace_table_partitions, Some(0) | Some(1)) {
1322 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 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 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 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 if let Some(table) = self
1420 .get_table(catalog_name, &schema_name, &physical_table)
1421 .await?
1422 {
1423 return Ok(table);
1424 }
1425
1426 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 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 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 METRIC_ENGINE_NAME
1515 } else {
1516 default_engine()
1517 };
1518
1519 let table_ref = TableReference::full(ctx.current_catalog(), &schema, &req.table_name);
1520 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 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 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 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 {
1591 let table_schema = table.schema();
1592 let ts_col_name = table_schema.timestamp_column().map(|c| c.name.clone());
1594 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 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 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 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
1749fn 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 rows.schema[ts_index].datatype_extension = None;
1782
1783 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 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
1857pub 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
1879pub 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 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 table_options.insert(
1916 COMPACTION_TYPE.to_string(),
1917 COMPACTION_TYPE_TWCS.to_string(),
1918 );
1919 }
1920 }
1921 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
1945pub type PerTableSemanticIndex = BTreeMap<String, BTreeMap<String, BTreeMap<String, String>>>;
1951
1952pub 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
1967pub 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
2001pub const SPLUNK_PK_METADATA_ORDER_KEY: &str = "splunk_pk_metadata_order";
2005
2006fn reorder_splunk_primary_keys(primary_keys: &mut [String]) {
2009 const LEAD: [&str; 3] = ["host", "source", "sourcetype"];
2010 primary_keys.sort_by_key(|name| {
2013 LEAD.iter()
2014 .position(|&lead| lead == name.as_str())
2015 .unwrap_or(LEAD.len())
2016 });
2017}
2018
2019struct CreateAlterTableResult {
2021 instant_table_ids: HashSet<TableId>,
2023 table_infos: HashMap<TableId, Arc<TableInfo>>,
2025}
2026
2027const MAX_MIRROR_PENDING_ROWS: usize = 1_000_000;
2035
2036const 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 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
2077fn 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 Some(None) => continue,
2102 _ => {
2103 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 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 inserts
2136 .entry(peers[0].clone())
2137 .or_default()
2138 .requests
2139 .extend(reqs.requests);
2140 continue;
2141 } else {
2142 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 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 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 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_mirror_pending_rows(&mirror_pending_rows, peer_rows);
2240 });
2241 }
2242
2243 Ok(())
2244 }
2245}
2246
2247fn 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
2258fn reserve_mirror_pending_rows(mirror_pending_rows: &AtomicU64, rows: u64) -> u64 {
2263 let pending = mirror_pending_rows.fetch_add(rows, Ordering::Relaxed) + rows;
2264 crate::metrics::DIST_MIRROR_PENDING_ROW_COUNT.add(rows as i64);
2267 pending
2268}
2269
2270fn 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 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 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 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 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 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 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 #[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 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 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 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 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 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 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 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 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 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 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 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 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 #[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 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 #[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 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 assert_eq!(
3900 MAX_MIRROR_PENDING_ROWS as u64,
3901 pending.load(Ordering::Relaxed)
3902 );
3903 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 #[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 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 assert_eq!(
3975 MAX_MIRROR_PENDING_ROWS as u64,
3976 pending.load(Ordering::Relaxed)
3977 );
3978
3979 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 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 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 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}