1use std::collections::{HashMap, HashSet};
18use std::fmt;
19use std::str::FromStr;
20
21use common_base::readable_size::ReadableSize;
22use common_datasource::object_store::oss::is_supported_in_oss;
23use common_datasource::object_store::s3::is_supported_in_s3;
24use common_query::AddColumnLocation;
25use common_time::TimeToLive;
26use common_time::range::TimestampRange;
27use datatypes::data_type::ConcreteDataType;
28use datatypes::json::JsonSettings;
29use datatypes::prelude::VectorRef;
30use datatypes::schema::{
31 ColumnDefaultConstraint, ColumnSchema, FulltextOptions, Schema, SkippingIndexOptions,
32};
33use greptime_proto::v1::region::{build_index_request, compact_request};
34use once_cell::sync::Lazy;
35use serde::{Deserialize, Serialize};
36use store_api::metric_engine_consts::{
37 LOGICAL_TABLE_METADATA_KEY, PHYSICAL_TABLE_METADATA_KEY, is_metric_engine_option_key,
38};
39use store_api::mito_engine_options::{
40 APPEND_MODE_KEY, COMPACTION_TYPE, EXPERIMENTAL_SST_FLOAT_FIELD_ENCODING, FloatFieldEncoding,
41 MEMTABLE_BULK_ENCODE_BYTES_THRESHOLD, MEMTABLE_BULK_ENCODE_ROW_THRESHOLD,
42 MEMTABLE_BULK_MAX_MERGE_GROUPS, MEMTABLE_BULK_MERGE_THRESHOLD, MEMTABLE_TYPE, MERGE_MODE_KEY,
43 SST_FORMAT_KEY, TWCS_ACTIVE_WINDOW_L1_MERGE_TRIGGER, TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM,
44 TWCS_FALLBACK_TO_LOCAL, TWCS_INACTIVE_WINDOW_L1_MERGE_TRIGGER,
45 TWCS_INACTIVE_WINDOW_TRIGGER_FILE_NUM, TWCS_MAX_OUTPUT_FILE_SIZE, TWCS_TIME_WINDOW,
46 TWCS_TRIGGER_FILE_NUM, is_mito_engine_option_key, normalize_twcs_trigger_options,
47};
48use store_api::region_request::{SetRegionOption, UnsetRegionOption};
49
50use crate::error::{ConflictingTableOptionsSnafu, ParseTableOptionSnafu, Result};
51use crate::metadata::{TableId, TableVersion};
52
53mod semantic;
54pub use semantic::*;
55
56pub const FILE_TABLE_META_KEY: &str = "__private.file_table_meta";
57pub const FILE_TABLE_LOCATION_KEY: &str = "location";
58pub const FILE_TABLE_PATTERN_KEY: &str = "pattern";
59pub const FILE_TABLE_FORMAT_KEY: &str = "format";
60
61pub const TABLE_DATA_MODEL: &str = "table_data_model";
62pub const TABLE_DATA_MODEL_TRACE_V1: &str = "greptime_trace_v1";
63pub const TABLE_DATA_MODEL_TRACE_V2: &str = "greptime_trace_v2";
65
66pub fn is_trace_table(table_info: &crate::metadata::TableInfo) -> bool {
69 let table_data_model = table_info.meta.options.data_model();
70 matches!(
71 table_data_model,
72 Some(TABLE_DATA_MODEL_TRACE_V1 | TABLE_DATA_MODEL_TRACE_V2)
73 )
74}
75
76pub const OTLP_METRIC_COMPAT_KEY: &str = "otlp_metric_compat";
77pub const OTLP_METRIC_COMPAT_PROM: &str = "prom";
78
79pub const VALID_TABLE_OPTION_KEYS: [&str; 15] = [
80 WRITE_BUFFER_SIZE_KEY,
82 TTL_KEY,
83 STORAGE_KEY,
84 COMMENT_KEY,
85 SKIP_WAL_KEY,
86 SST_FORMAT_KEY,
87 FILE_TABLE_LOCATION_KEY,
89 FILE_TABLE_FORMAT_KEY,
90 FILE_TABLE_PATTERN_KEY,
91 PHYSICAL_TABLE_METADATA_KEY,
93 LOGICAL_TABLE_METADATA_KEY,
94 TABLE_DATA_MODEL,
96 OTLP_METRIC_COMPAT_KEY,
97 REPARTITION_COLUMN_HINT_KEY,
98 REPARTITION_PARTITION_NUM_HINT_KEY,
99];
100
101pub const DDL_TIMEOUT: &str = "timeout";
102pub const DDL_WAIT: &str = "wait";
103
104pub const VALID_DDL_OPTION_KEYS: [&str; 2] = [DDL_TIMEOUT, DDL_WAIT];
105
106pub const INGEST_ROWS_RATE_LIMIT_KEY: &str = "ingest_rows_rate_limit";
108
109static VALID_DB_OPT_KEYS: Lazy<HashSet<&str>> = Lazy::new(|| {
111 let mut set = HashSet::new();
112 set.insert(TTL_KEY);
113 set.insert(STORAGE_KEY);
114 set.insert(MEMTABLE_TYPE);
115 set.insert(MEMTABLE_BULK_MERGE_THRESHOLD);
116 set.insert(MEMTABLE_BULK_ENCODE_ROW_THRESHOLD);
117 set.insert(MEMTABLE_BULK_ENCODE_BYTES_THRESHOLD);
118 set.insert(MEMTABLE_BULK_MAX_MERGE_GROUPS);
119 set.insert(APPEND_MODE_KEY);
120 set.insert(MERGE_MODE_KEY);
121 set.insert(SKIP_WAL_KEY);
122 set.insert(COMPACTION_TYPE);
123 set.insert(TWCS_FALLBACK_TO_LOCAL);
124 set.insert(TWCS_TIME_WINDOW);
125 set.insert(TWCS_TRIGGER_FILE_NUM);
126 set.insert(TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM);
127 set.insert(TWCS_ACTIVE_WINDOW_L1_MERGE_TRIGGER);
128 set.insert(TWCS_INACTIVE_WINDOW_TRIGGER_FILE_NUM);
129 set.insert(TWCS_INACTIVE_WINDOW_L1_MERGE_TRIGGER);
130 set.insert(TWCS_MAX_OUTPUT_FILE_SIZE);
131 set.insert(SST_FORMAT_KEY);
132 set.insert(INGEST_ROWS_RATE_LIMIT_KEY);
133 set
134});
135
136pub fn validate_database_option(key: &str) -> bool {
138 VALID_DB_OPT_KEYS.contains(&key)
139}
140
141pub fn validate_database_option_value(
143 key: &str,
144 value: Option<&str>,
145) -> std::result::Result<(), &'static str> {
146 if key == INGEST_ROWS_RATE_LIMIT_KEY {
147 return value
148 .and_then(|value| value.parse::<u64>().ok())
149 .map(|_| ())
150 .ok_or("expected a non-negative integer fitting in u64");
151 }
152 let (minimum, constraint) = match key {
153 TWCS_TRIGGER_FILE_NUM
154 | TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM
155 | TWCS_INACTIVE_WINDOW_TRIGGER_FILE_NUM => {
156 (0, "expected a non-negative integer fitting in usize")
157 }
158 TWCS_ACTIVE_WINDOW_L1_MERGE_TRIGGER | TWCS_INACTIVE_WINDOW_L1_MERGE_TRIGGER => {
159 (2, "expected an integer greater than or equal to 2")
160 }
161 _ => return Ok(()),
162 };
163 if value
164 .and_then(|value| value.parse::<usize>().ok())
165 .is_some_and(|files| files >= minimum)
166 {
167 Ok(())
168 } else {
169 Err(constraint)
170 }
171}
172
173pub fn validate_table_option(key: &str) -> bool {
175 if is_supported_in_s3(key) {
176 return true;
177 }
178
179 if is_supported_in_oss(key) {
180 return true;
181 }
182
183 if is_mito_engine_option_key(key) {
184 return true;
185 }
186
187 if is_metric_engine_option_key(key) {
188 return true;
189 }
190
191 if is_semantic_option_key(key) {
194 return true;
195 }
196
197 VALID_TABLE_OPTION_KEYS.contains(&key) || VALID_DDL_OPTION_KEYS.contains(&key)
198}
199
200#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
201#[serde(default)]
202pub struct TableOptions {
203 pub write_buffer_size: Option<ReadableSize>,
205 pub ttl: Option<TimeToLive>,
207 pub skip_wal: bool,
209 pub extra_options: HashMap<String, String>,
211}
212
213pub const WRITE_BUFFER_SIZE_KEY: &str = store_api::mito_engine_options::WRITE_BUFFER_SIZE_KEY;
214pub const TTL_KEY: &str = store_api::mito_engine_options::TTL_KEY;
215pub const STORAGE_KEY: &str = "storage";
216pub const COMMENT_KEY: &str = "comment";
217pub const AUTO_CREATE_TABLE_KEY: &str = "auto_create_table";
218pub const SKIP_WAL_KEY: &str = store_api::mito_engine_options::SKIP_WAL_KEY;
219pub const TRACE_TABLE_PARTITIONS_HINT_KEY: &str = "trace_table_partitions";
220pub const REPARTITION_COLUMN_HINT_KEY: &str = "repartition.column.hint";
221
222pub const REPARTITION_PARTITION_NUM_HINT_KEY: &str = "repartition.partition.num.hint";
224
225impl TableOptions {
226 pub fn data_model(&self) -> Option<&str> {
228 self.extra_options.get(TABLE_DATA_MODEL).map(String::as_str)
229 }
230
231 pub fn try_from_iter<T: ToString, U: IntoIterator<Item = (T, T)>>(
232 iter: U,
233 ) -> Result<TableOptions> {
234 let mut options = TableOptions::default();
235
236 let mut kvs: HashMap<String, String> = iter
237 .into_iter()
238 .map(|(k, v)| (k.to_string(), v.to_string()))
239 .collect();
240
241 normalize_twcs_trigger_options(&mut kvs).map_err(|conflict| {
242 ConflictingTableOptionsSnafu {
243 first_key: TWCS_TRIGGER_FILE_NUM,
244 first_value: conflict.legacy_value,
245 second_key: TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM,
246 second_value: conflict.canonical_value,
247 }
248 .build()
249 })?;
250
251 if let Some(write_buffer_size) = kvs.get(WRITE_BUFFER_SIZE_KEY) {
252 let size = ReadableSize::from_str(write_buffer_size).map_err(|_| {
253 ParseTableOptionSnafu {
254 key: WRITE_BUFFER_SIZE_KEY,
255 value: write_buffer_size,
256 }
257 .build()
258 })?;
259 options.write_buffer_size = Some(size)
260 }
261
262 if let Some(ttl) = kvs.get(TTL_KEY) {
263 let ttl_value = TimeToLive::from_humantime_or_str(ttl).map_err(|_| {
264 ParseTableOptionSnafu {
265 key: TTL_KEY,
266 value: ttl,
267 }
268 .build()
269 })?;
270 options.ttl = Some(ttl_value);
271 }
272
273 if let Some(encoding) = kvs.get(EXPERIMENTAL_SST_FLOAT_FIELD_ENCODING) {
274 encoding.parse::<FloatFieldEncoding>().map_err(|_| {
275 ParseTableOptionSnafu {
276 key: EXPERIMENTAL_SST_FLOAT_FIELD_ENCODING,
277 value: encoding,
278 }
279 .build()
280 })?;
281 }
282
283 if let Some(skip_wal) = kvs.get(SKIP_WAL_KEY) {
284 options.skip_wal = skip_wal.parse().map_err(|_| {
285 ParseTableOptionSnafu {
286 key: SKIP_WAL_KEY,
287 value: skip_wal,
288 }
289 .build()
290 })?;
291 }
292
293 options.extra_options = HashMap::from_iter(
294 kvs.into_iter()
295 .filter(|(k, _)| k != WRITE_BUFFER_SIZE_KEY && k != TTL_KEY),
296 );
297
298 Ok(options)
299 }
300}
301
302impl fmt::Display for TableOptions {
303 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
304 let mut key_vals = vec![];
305 if let Some(size) = self.write_buffer_size {
306 key_vals.push(format!("{}={}", WRITE_BUFFER_SIZE_KEY, size));
307 }
308
309 if let Some(ttl) = self.ttl.map(|ttl| ttl.to_string()) {
310 key_vals.push(format!("{}={}", TTL_KEY, ttl));
311 }
312
313 if self.skip_wal && !self.extra_options.contains_key(SKIP_WAL_KEY) {
314 key_vals.push(format!("{}={}", SKIP_WAL_KEY, self.skip_wal));
315 }
316
317 for (k, v) in &self.extra_options {
318 key_vals.push(format!("{}={}", k, v));
319 }
320
321 write!(f, "{}", key_vals.join(" "))
322 }
323}
324
325impl From<&TableOptions> for HashMap<String, String> {
326 fn from(opts: &TableOptions) -> Self {
327 let mut res = HashMap::with_capacity(3 + opts.extra_options.len());
328 if let Some(write_buffer_size) = opts.write_buffer_size {
329 let _ = res.insert(
330 WRITE_BUFFER_SIZE_KEY.to_string(),
331 write_buffer_size.to_string(),
332 );
333 }
334 if let Some(ttl_str) = opts.ttl.map(|ttl| ttl.to_string()) {
335 let _ = res.insert(TTL_KEY.to_string(), ttl_str);
336 }
337 if opts.skip_wal {
338 let _ = res.insert(SKIP_WAL_KEY.to_string(), true.to_string());
339 }
340 res.extend(opts.extra_options.clone());
341 res
342 }
343}
344
345#[derive(Debug, Clone, Serialize, Deserialize)]
347pub struct AlterTableRequest {
348 pub catalog_name: String,
349 pub schema_name: String,
350 pub table_name: String,
351 pub table_id: TableId,
352 pub alter_kind: AlterKind,
353 pub table_version: Option<TableVersion>,
355}
356
357#[derive(Debug, Clone, Serialize, Deserialize)]
359pub struct AddColumnRequest {
360 pub column_schema: ColumnSchema,
361 pub is_key: bool,
362 pub location: Option<AddColumnLocation>,
363 pub add_if_not_exists: bool,
365}
366
367#[derive(Debug, Clone, Serialize, Deserialize)]
369pub struct ModifyColumnTypeRequest {
370 pub column_name: String,
371 pub target_type: ConcreteDataType,
372}
373
374#[derive(Debug, Clone, Serialize, Deserialize)]
376pub struct SetJsonSettingsRequest {
377 pub column_name: String,
378 pub settings: JsonSettings,
379}
380
381#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
385pub enum AnnotationFamily {
386 Semantic,
388 RepartitionHint,
390}
391
392impl AnnotationFamily {
393 pub fn namespace(self) -> &'static str {
395 match self {
396 Self::Semantic => SEMANTIC_PREFIX,
397 Self::RepartitionHint => "repartition.",
398 }
399 }
400
401 pub fn of_key(key: &str) -> Option<Self> {
402 if key.starts_with(SEMANTIC_PREFIX) {
403 Some(Self::Semantic)
404 } else if matches!(
405 key,
406 REPARTITION_COLUMN_HINT_KEY | REPARTITION_PARTITION_NUM_HINT_KEY
407 ) {
408 Some(Self::RepartitionHint)
409 } else {
410 None
411 }
412 }
413
414 pub fn allows_logical_tables(self) -> bool {
418 match self {
419 Self::Semantic => true,
420 Self::RepartitionHint => false,
421 }
422 }
423
424 pub fn requires_unique_keys(self) -> bool {
426 matches!(self, Self::RepartitionHint)
427 }
428
429 pub fn mixed_batch_error(self) -> String {
431 match self {
432 Self::Semantic => format!(
433 "`{SEMANTIC_PREFIX}*` options must be altered separately from other table options"
434 ),
435 Self::RepartitionHint => {
436 "repartition hints must be altered separately from other table options".to_string()
437 }
438 }
439 }
440}
441
442#[derive(Debug, Clone, PartialEq, Eq)]
444pub enum AnnotationKeyError {
445 MixedFamilies { family: AnnotationFamily },
446 DuplicateKey,
447}
448
449impl fmt::Display for AnnotationKeyError {
450 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
451 match self {
452 Self::MixedFamilies { family } => f.write_str(&family.mixed_batch_error()),
453 Self::DuplicateKey => f.write_str("duplicate repartition hint keys"),
454 }
455 }
456}
457
458pub fn validate_annotation_keys<'a>(
464 keys: impl IntoIterator<Item = &'a str>,
465) -> std::result::Result<Option<AnnotationFamily>, AnnotationKeyError> {
466 let mut keys = keys.into_iter();
467 let Some(first) = keys.next() else {
468 return Ok(None);
469 };
470 let family = AnnotationFamily::of_key(first);
471 let reject_duplicates = family.is_some_and(|family| family.requires_unique_keys());
472 let mut seen = HashSet::new();
473 if reject_duplicates {
474 seen.insert(first);
475 }
476 for key in keys {
477 let this = AnnotationFamily::of_key(key);
478 if this != family
479 && let Some(family) = family.or(this)
480 {
481 return Err(AnnotationKeyError::MixedFamilies { family });
482 }
483 if reject_duplicates && !seen.insert(key) {
484 return Err(AnnotationKeyError::DuplicateKey);
485 }
486 }
487 Ok(family)
488}
489
490pub struct AnnotationContext<'a> {
492 pub data_model: Option<&'a str>,
493 pub schema: &'a Schema,
494 pub partition_key_indices: &'a [usize],
495}
496
497#[derive(Debug, Clone, PartialEq, Eq)]
502pub enum AnnotationValidationError {
503 UnknownKey {
504 key: String,
505 },
506 InvalidValue {
507 key: String,
508 value: String,
509 },
510 ColumnNotFound {
511 column: String,
512 },
513 ColumnNotStringForm {
514 key: String,
515 column: String,
516 ty: ConcreteDataType,
517 },
518 InvalidPartitionNumHint {
519 value: String,
520 },
521 NotSingleColumn,
522 PartitionMetadataConflict,
523 TimeIndexConflict,
524}
525
526impl fmt::Display for AnnotationValidationError {
527 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
528 match self {
529 Self::UnknownKey { key } => write!(f, "unknown semantic option `{key}`"),
530 Self::InvalidValue { key, value } => {
531 write!(f, "invalid value `{value}` for semantic option `{key}`")
532 }
533 Self::ColumnNotFound { column } => write!(f, "column `{column}` not found"),
534 Self::ColumnNotStringForm { key, column, ty } => write!(
535 f,
536 "entity column `{column}` (option `{key}`) has type `{ty}`, \
537 which cannot render as a string"
538 ),
539 Self::InvalidPartitionNumHint { value } => write!(
540 f,
541 "{REPARTITION_PARTITION_NUM_HINT_KEY} expects a positive integer within u32 range, got `{value}`"
542 ),
543 Self::NotSingleColumn => write!(
544 f,
545 "{REPARTITION_COLUMN_HINT_KEY} expects exactly one column name"
546 ),
547 Self::PartitionMetadataConflict => write!(
548 f,
549 "cannot set {REPARTITION_COLUMN_HINT_KEY} on a table with partition metadata"
550 ),
551 Self::TimeIndexConflict => write!(
552 f,
553 "cannot set {REPARTITION_COLUMN_HINT_KEY} to the time index column"
554 ),
555 }
556 }
557}
558
559pub(crate) fn validate_and_normalize_annotation(
563 family: AnnotationFamily,
564 cx: &AnnotationContext<'_>,
565 key: &str,
566 value: &str,
567) -> std::result::Result<String, AnnotationValidationError> {
568 match family {
569 AnnotationFamily::Semantic => {
570 if !is_semantic_option_key(key) {
571 return Err(AnnotationValidationError::UnknownKey {
572 key: key.to_string(),
573 });
574 }
575 if !validate_semantic_option(key, value) {
576 return Err(AnnotationValidationError::InvalidValue {
577 key: key.to_string(),
578 value: value.to_string(),
579 });
580 }
581 if parse_entity_option_key(key).is_some() {
582 for column in parse_entity_columns(value) {
583 if trace_v2_attribute(cx.schema, cx.data_model, &column).is_some() {
584 continue;
585 }
586 let schema = cx.schema.column_schema_by_name(&column).ok_or_else(|| {
587 AnnotationValidationError::ColumnNotFound {
588 column: column.clone(),
589 }
590 })?;
591 if !has_stable_string_form(&schema.data_type) {
592 return Err(AnnotationValidationError::ColumnNotStringForm {
593 key: key.to_string(),
594 column,
595 ty: schema.data_type.clone(),
596 });
597 }
598 }
599 }
600 Ok(value.to_string())
601 }
602 AnnotationFamily::RepartitionHint if key == REPARTITION_PARTITION_NUM_HINT_KEY => {
603 let value = value.trim();
604 if !matches!(value.parse::<u32>(), Ok(1..)) {
605 return Err(AnnotationValidationError::InvalidPartitionNumHint {
606 value: value.to_string(),
607 });
608 }
609 Ok(value.to_string())
610 }
611 AnnotationFamily::RepartitionHint => {
612 let column_name = value.trim();
613 if column_name.is_empty() || column_name.contains(',') {
614 return Err(AnnotationValidationError::NotSingleColumn);
615 }
616 if !cx.partition_key_indices.is_empty() {
617 return Err(AnnotationValidationError::PartitionMetadataConflict);
618 }
619 let column_index = cx.schema.column_index_by_name(column_name).ok_or_else(|| {
620 AnnotationValidationError::ColumnNotFound {
621 column: column_name.to_string(),
622 }
623 })?;
624 if cx.schema.timestamp_index() == Some(column_index) {
625 return Err(AnnotationValidationError::TimeIndexConflict);
626 }
627 Ok(column_name.to_string())
628 }
629 }
630}
631
632pub fn validate_and_normalize_annotation_options(
635 options: &mut TableOptions,
636 schema: &Schema,
637 partition_key_indices: &[usize],
638) -> std::result::Result<(), AnnotationValidationError> {
639 let cx = AnnotationContext {
640 data_model: options.data_model(),
641 schema,
642 partition_key_indices,
643 };
644 let mut normalized = Vec::new();
645 for (key, value) in &options.extra_options {
646 let Some(family) = AnnotationFamily::of_key(key) else {
647 continue;
648 };
649 let checked = validate_and_normalize_annotation(family, &cx, key, value)?;
650 if checked != *value {
651 normalized.push((key.clone(), checked));
652 }
653 }
654 for (key, value) in normalized {
655 options.extra_options.insert(key, value);
656 }
657 Ok(())
658}
659
660#[derive(Debug, Clone, Serialize, Deserialize)]
661pub enum AlterKind {
662 AddColumns {
663 columns: Vec<AddColumnRequest>,
664 },
665 DropColumns {
666 names: Vec<String>,
667 },
668 ModifyColumnTypes {
669 columns: Vec<ModifyColumnTypeRequest>,
670 },
671 SetJsonSettings {
672 request: SetJsonSettingsRequest,
673 },
674 RenameTable {
675 new_table_name: String,
676 },
677 SetTableOptions {
678 options: Vec<SetRegionOption>,
679 },
680 UnsetTableOptions {
681 keys: Vec<UnsetRegionOption>,
682 },
683 SetAnnotations {
684 family: AnnotationFamily,
685 options: Vec<(String, String)>,
686 },
687 UnsetAnnotations {
688 family: AnnotationFamily,
689 keys: Vec<String>,
690 },
691 SetIndexes {
692 options: Vec<SetIndexOption>,
693 },
694 UnsetIndexes {
695 options: Vec<UnsetIndexOption>,
696 },
697 DropDefaults {
698 names: Vec<String>,
699 },
700 SetDefaults {
701 defaults: Vec<SetDefaultRequest>,
702 },
703}
704
705#[derive(Debug, Clone, Serialize, Deserialize)]
706pub struct SetDefaultRequest {
707 pub column_name: String,
708 pub default_constraint: Option<ColumnDefaultConstraint>,
709}
710
711#[derive(Debug, Clone, Serialize, Deserialize)]
712pub enum SetIndexOption {
713 Fulltext {
714 column_name: String,
715 options: FulltextOptions,
716 },
717 Inverted {
718 column_name: String,
719 },
720 Skipping {
721 column_name: String,
722 options: SkippingIndexOptions,
723 },
724}
725
726impl SetIndexOption {
727 pub fn column_name(&self) -> &str {
729 match self {
730 SetIndexOption::Fulltext { column_name, .. } => column_name,
731 SetIndexOption::Inverted { column_name, .. } => column_name,
732 SetIndexOption::Skipping { column_name, .. } => column_name,
733 }
734 }
735}
736
737#[derive(Debug, Clone, Serialize, Deserialize)]
738pub enum UnsetIndexOption {
739 Fulltext { column_name: String },
740 Inverted { column_name: String },
741 Skipping { column_name: String },
742}
743
744impl UnsetIndexOption {
745 pub fn column_name(&self) -> &str {
747 match self {
748 UnsetIndexOption::Fulltext { column_name, .. } => column_name,
749 UnsetIndexOption::Inverted { column_name, .. } => column_name,
750 UnsetIndexOption::Skipping { column_name, .. } => column_name,
751 }
752 }
753}
754
755#[derive(Debug)]
756pub struct InsertRequest {
757 pub catalog_name: String,
758 pub schema_name: String,
759 pub table_name: String,
760 pub columns_values: HashMap<String, VectorRef>,
761 pub skip_wal: bool,
763}
764
765#[derive(Debug)]
767pub struct DeleteRequest {
768 pub catalog_name: String,
769 pub schema_name: String,
770 pub table_name: String,
771 pub key_column_values: HashMap<String, VectorRef>,
775}
776
777#[derive(Debug)]
778pub enum CopyDirection {
779 Export,
780 Import,
781}
782
783#[derive(Debug)]
785pub struct CopyTableRequest {
786 pub catalog_name: String,
787 pub schema_name: String,
788 pub table_name: String,
789 pub location: String,
790 pub with: HashMap<String, String>,
791 pub connection: HashMap<String, String>,
792 pub pattern: Option<String>,
793 pub direction: CopyDirection,
794 pub timestamp_range: Option<TimestampRange>,
795 pub limit: Option<u64>,
796}
797
798#[derive(Debug, Clone, Default)]
799pub struct FlushTableRequest {
800 pub catalog_name: String,
801 pub schema_name: String,
802 pub table_name: String,
803}
804
805#[derive(Debug, Clone, Default)]
806pub struct BuildIndexTableRequest {
807 pub options: Option<build_index_request::Options>,
809 pub catalog_name: String,
810 pub schema_name: String,
811 pub table_name: String,
812}
813
814#[derive(Debug, Clone, PartialEq)]
815pub struct CompactTableRequest {
816 pub catalog_name: String,
817 pub schema_name: String,
818 pub table_name: String,
819 pub compact_options: compact_request::Options,
820 pub parallelism: u32,
821 pub time_range: Option<TimestampRange>,
822}
823
824impl Default for CompactTableRequest {
825 fn default() -> Self {
826 Self {
827 catalog_name: Default::default(),
828 schema_name: Default::default(),
829 table_name: Default::default(),
830 compact_options: compact_request::Options::Regular(Default::default()),
831 parallelism: 1,
832 time_range: None,
833 }
834 }
835}
836
837#[derive(Debug, Clone, Default, Deserialize, Serialize)]
838pub struct CopyDatabaseRequest {
839 pub catalog_name: String,
840 pub schema_name: String,
841 pub location: String,
842 pub with: HashMap<String, String>,
843 pub connection: HashMap<String, String>,
844 pub time_range: Option<TimestampRange>,
845}
846
847#[derive(Debug, Clone, Default, Deserialize, Serialize)]
848pub struct CopyQueryToRequest {
849 pub location: String,
850 pub with: HashMap<String, String>,
851 pub connection: HashMap<String, String>,
852}
853
854#[cfg(test)]
855mod tests {
856 use std::time::Duration;
857
858 use common_error::ext::ErrorExt;
859 use common_error::status_code::StatusCode;
860
861 use super::*;
862
863 #[test]
864 fn test_validate_annotation_keys() {
865 let column = REPARTITION_COLUMN_HINT_KEY;
866 let count = REPARTITION_PARTITION_NUM_HINT_KEY;
867 for (keys, expected) in [
868 (vec![], None),
869 (vec![TTL_KEY, TTL_KEY], None),
870 (vec!["repartition.unknown.hint"], None),
871 (vec![column], Some(AnnotationFamily::RepartitionHint)),
872 (vec![count], Some(AnnotationFamily::RepartitionHint)),
873 (vec![column, count], Some(AnnotationFamily::RepartitionHint)),
874 (vec![count, column], Some(AnnotationFamily::RepartitionHint)),
875 (
876 vec!["greptime.semantic.source", "greptime.semantic.source"],
877 Some(AnnotationFamily::Semantic),
878 ),
879 ] {
880 assert_eq!(validate_annotation_keys(keys), Ok(expected));
881 }
882 for keys in [
883 vec![column, column],
884 vec![count, count],
885 vec![column, count, column],
886 ] {
887 assert_eq!(
888 validate_annotation_keys(keys),
889 Err(AnnotationKeyError::DuplicateKey)
890 );
891 }
892 for keys in [
893 vec![column, TTL_KEY],
894 vec![TTL_KEY, count],
895 vec![column, "greptime.semantic.source"],
896 ] {
897 assert_eq!(
898 validate_annotation_keys(keys),
899 Err(AnnotationKeyError::MixedFamilies {
900 family: AnnotationFamily::RepartitionHint,
901 })
902 );
903 }
904 }
905
906 #[test]
907 fn test_validate_table_option() {
908 assert!(validate_table_option(FILE_TABLE_LOCATION_KEY));
909 assert!(validate_table_option(FILE_TABLE_FORMAT_KEY));
910 assert!(validate_table_option(FILE_TABLE_PATTERN_KEY));
911 assert!(validate_table_option(TTL_KEY));
912 assert!(validate_table_option(WRITE_BUFFER_SIZE_KEY));
913 assert!(validate_table_option(STORAGE_KEY));
914 assert!(validate_table_option(MEMTABLE_BULK_MERGE_THRESHOLD));
915 assert!(validate_table_option(EXPERIMENTAL_SST_FLOAT_FIELD_ENCODING));
916 assert!(validate_table_option(REPARTITION_COLUMN_HINT_KEY));
917 assert!(validate_table_option(REPARTITION_PARTITION_NUM_HINT_KEY));
918 assert_eq!(AnnotationFamily::of_key("repartition.unknown.hint"), None);
919 assert!(!validate_table_option("foo"));
920
921 assert!(validate_table_option(SEMANTIC_SIGNAL_TYPE));
923 assert!(validate_table_option(SEMANTIC_METRIC_TYPE));
924 assert!(!validate_table_option("greptime.semantic.future.key"));
926 assert!(!validate_table_option("greptime.semanticx"));
927 assert!(!validate_table_option(SEMANTIC_PER_TABLE_INDEX_KEY));
928 }
929
930 #[test]
931 fn test_validate_database_option() {
932 assert!(validate_database_option(MEMTABLE_TYPE));
933 assert!(validate_database_option(MEMTABLE_BULK_MERGE_THRESHOLD));
934 assert!(validate_database_option(MEMTABLE_BULK_ENCODE_ROW_THRESHOLD));
935 assert!(validate_database_option(
936 MEMTABLE_BULK_ENCODE_BYTES_THRESHOLD
937 ));
938 assert!(validate_database_option(MEMTABLE_BULK_MAX_MERGE_GROUPS));
939 assert!(validate_database_option(
940 "compaction.twcs.active_window.trigger_file_num"
941 ));
942 assert!(validate_database_option(
943 "compaction.twcs.active_window.l1_merge_trigger"
944 ));
945 assert!(validate_database_option(
946 "compaction.twcs.inactive_window.trigger_file_num"
947 ));
948 assert!(validate_database_option(
949 "compaction.twcs.inactive_window.l1_merge_trigger"
950 ));
951 assert!(validate_database_option(INGEST_ROWS_RATE_LIMIT_KEY));
952 assert!(validate_database_option("ingest_rows_rate_limit"));
953 assert!(!validate_database_option("foo"));
954 }
955
956 #[test]
957 fn test_parse_float_field_encoding_option() {
958 let options = TableOptions::try_from_iter([(
959 EXPERIMENTAL_SST_FLOAT_FIELD_ENCODING,
960 "byte_stream_split",
961 )])
962 .unwrap();
963 assert_eq!(
964 options
965 .extra_options
966 .get(EXPERIMENTAL_SST_FLOAT_FIELD_ENCODING),
967 Some(&"byte_stream_split".to_string())
968 );
969 assert!(
970 TableOptions::try_from_iter([(EXPERIMENTAL_SST_FLOAT_FIELD_ENCODING, "invalid",)])
971 .is_err()
972 );
973 }
974
975 #[test]
976 fn test_database_trigger_value_boundaries() {
977 let maximum = usize::MAX.to_string();
978 let overflow = format!("{maximum}0");
979 for (key, minimum) in [
980 (TWCS_TRIGGER_FILE_NUM, 0),
981 (TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM, 0),
982 (TWCS_INACTIVE_WINDOW_TRIGGER_FILE_NUM, 0),
983 (TWCS_ACTIVE_WINDOW_L1_MERGE_TRIGGER, 2),
984 (TWCS_INACTIVE_WINDOW_L1_MERGE_TRIGGER, 2),
985 ] {
986 for invalid in [
987 None,
988 Some(""),
989 Some("invalid"),
990 Some("-1"),
991 Some(overflow.as_str()),
992 ] {
993 assert!(
994 validate_database_option_value(key, invalid).is_err(),
995 "{key}: {invalid:?}"
996 );
997 }
998 for valid in ["2", maximum.as_str()] {
999 assert!(
1000 validate_database_option_value(key, Some(valid)).is_ok(),
1001 "{key}: {valid}"
1002 );
1003 }
1004 for boundary in ["0", "1"] {
1005 assert_eq!(
1006 validate_database_option_value(key, Some(boundary)).is_ok(),
1007 minimum == 0,
1008 "{key}: {boundary}"
1009 );
1010 }
1011 }
1012 }
1013
1014 #[test]
1015 fn test_database_ingest_rate_limit_value_boundaries() {
1016 for invalid in [
1017 None,
1018 Some(""),
1019 Some("abc"),
1020 Some("1000/s"),
1021 Some("-1"),
1022 Some("1.5"),
1023 Some("18446744073709551616"),
1024 ] {
1025 assert!(validate_database_option_value(INGEST_ROWS_RATE_LIMIT_KEY, invalid).is_err());
1026 }
1027 let maximum = u64::MAX.to_string();
1028 for valid in ["0", "1", maximum.as_str()] {
1029 assert!(
1030 validate_database_option_value(INGEST_ROWS_RATE_LIMIT_KEY, Some(valid)).is_ok()
1031 );
1032 }
1033 }
1034
1035 #[test]
1036 fn test_serialize_table_options() {
1037 let options = TableOptions {
1038 write_buffer_size: None,
1039 ttl: Some(Duration::from_secs(1000).into()),
1040 extra_options: HashMap::new(),
1041 skip_wal: false,
1042 };
1043 let serialized = serde_json::to_string(&options).unwrap();
1044 let deserialized: TableOptions = serde_json::from_str(&serialized).unwrap();
1045 assert_eq!(options, deserialized);
1046 }
1047
1048 #[test]
1049 fn test_convert_hashmap_between_table_options() {
1050 let options = TableOptions {
1051 write_buffer_size: Some(ReadableSize::mb(128)),
1052 ttl: Some(Duration::from_secs(1000).into()),
1053 extra_options: HashMap::new(),
1054 skip_wal: false,
1055 };
1056 let serialized_map = HashMap::from(&options);
1057 let serialized = TableOptions::try_from_iter(&serialized_map).unwrap();
1058 assert_eq!(options, serialized);
1059
1060 let options = TableOptions {
1061 write_buffer_size: None,
1062 ttl: None,
1063 extra_options: HashMap::from([(SKIP_WAL_KEY.to_string(), true.to_string())]),
1064 skip_wal: true,
1065 };
1066 let serialized_map = HashMap::from(&options);
1067 assert_eq!(
1068 Some("true"),
1069 serialized_map.get(SKIP_WAL_KEY).map(String::as_str)
1070 );
1071 let serialized = TableOptions::try_from_iter(&serialized_map).unwrap();
1072 assert_eq!(options, serialized);
1073
1074 let options = TableOptions {
1075 write_buffer_size: None,
1076 ttl: Default::default(),
1077 extra_options: HashMap::new(),
1078 skip_wal: false,
1079 };
1080 let serialized_map = HashMap::from(&options);
1081 let serialized = TableOptions::try_from_iter(&serialized_map).unwrap();
1082 assert_eq!(options, serialized);
1083
1084 let options = TableOptions {
1085 write_buffer_size: Some(ReadableSize::mb(128)),
1086 ttl: Some(Duration::from_secs(1000).into()),
1087 extra_options: HashMap::from([("a".to_string(), "A".to_string())]),
1088 skip_wal: false,
1089 };
1090 let serialized_map = HashMap::from(&options);
1091 let serialized = TableOptions::try_from_iter(&serialized_map).unwrap();
1092 assert_eq!(options, serialized);
1093
1094 let options = TableOptions {
1095 extra_options: HashMap::from([(SKIP_WAL_KEY.to_string(), false.to_string())]),
1096 skip_wal: false,
1097 ..Default::default()
1098 };
1099 let serialized_map = HashMap::from(&options);
1100 let serialized = TableOptions::try_from_iter(&serialized_map).unwrap();
1101 assert_eq!(options, serialized);
1102 }
1103
1104 #[test]
1105 fn test_table_options_normalizes_twcs_trigger_aliases() {
1106 for options in [
1107 vec![(TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM, "4")],
1108 vec![
1109 (TWCS_TRIGGER_FILE_NUM, "4"),
1110 (TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM, "4"),
1111 ],
1112 ] {
1113 let table_options = TableOptions::try_from_iter(options).unwrap();
1114 assert_eq!(
1115 HashMap::from([(TWCS_TRIGGER_FILE_NUM.to_string(), "4".to_string())]),
1116 table_options.extra_options
1117 );
1118 }
1119 }
1120
1121 #[test]
1122 fn test_table_options_rejects_conflicting_twcs_trigger_aliases() {
1123 let error = TableOptions::try_from_iter([
1124 (TWCS_TRIGGER_FILE_NUM, "4"),
1125 (TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM, "8"),
1126 ])
1127 .unwrap_err();
1128 assert_eq!(StatusCode::InvalidArguments, error.status_code());
1129 assert_eq!(
1130 "Conflicting table options: compaction.twcs.trigger_file_num=4 and compaction.twcs.active_window.trigger_file_num=8",
1131 error.to_string()
1132 );
1133 }
1134
1135 #[test]
1136 fn test_table_options_to_string() {
1137 let options = TableOptions {
1138 write_buffer_size: Some(ReadableSize::mb(128)),
1139 ttl: Some(Duration::from_secs(1000).into()),
1140 extra_options: HashMap::new(),
1141 skip_wal: false,
1142 };
1143
1144 assert_eq!(
1145 "write_buffer_size=128.0MiB ttl=16m 40s",
1146 options.to_string()
1147 );
1148
1149 let options = TableOptions {
1150 write_buffer_size: Some(ReadableSize::mb(128)),
1151 ttl: Some(Duration::from_secs(1000).into()),
1152 extra_options: HashMap::from([("a".to_string(), "A".to_string())]),
1153 skip_wal: false,
1154 };
1155
1156 assert_eq!(
1157 "write_buffer_size=128.0MiB ttl=16m 40s a=A",
1158 options.to_string()
1159 );
1160
1161 let options = TableOptions {
1162 write_buffer_size: Some(ReadableSize::mb(128)),
1163 ttl: Some(Duration::from_secs(1000).into()),
1164 extra_options: HashMap::new(),
1165 skip_wal: true,
1166 };
1167 assert_eq!(
1168 "write_buffer_size=128.0MiB ttl=16m 40s skip_wal=true",
1169 options.to_string()
1170 );
1171
1172 let options = TableOptions {
1173 write_buffer_size: None,
1174 ttl: None,
1175 extra_options: HashMap::from([(SKIP_WAL_KEY.to_string(), "false".to_string())]),
1176 skip_wal: false,
1177 };
1178 assert_eq!("skip_wal=false", options.to_string());
1179 }
1180}