Skip to main content

table/
requests.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Table and TableEngine requests
16
17use 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";
63/// Table data model used by the JSON2-based OTLP trace pipeline.
64pub const TABLE_DATA_MODEL_TRACE_V2: &str = "greptime_trace_v2";
65
66/// Returns true for the Trace V1 and V2 data models supported
67/// by semantic graph derivation.
68pub 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    // common keys:
81    WRITE_BUFFER_SIZE_KEY,
82    TTL_KEY,
83    STORAGE_KEY,
84    COMMENT_KEY,
85    SKIP_WAL_KEY,
86    SST_FORMAT_KEY,
87    // file engine keys:
88    FILE_TABLE_LOCATION_KEY,
89    FILE_TABLE_FORMAT_KEY,
90    FILE_TABLE_PATTERN_KEY,
91    // metric engine keys:
92    PHYSICAL_TABLE_METADATA_KEY,
93    LOGICAL_TABLE_METADATA_KEY,
94    // table model info
95    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
106/// The key of ingest rows rate limit option (rows per second, cluster-wide) in database options.
107pub const INGEST_ROWS_RATE_LIMIT_KEY: &str = "ingest_rows_rate_limit";
108
109// Valid option keys when creating a db.
110static 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
136/// Returns true if the `key` is a valid key for database.
137pub fn validate_database_option(key: &str) -> bool {
138    VALID_DB_OPT_KEYS.contains(&key)
139}
140
141/// Validates a database option value, returning the violated constraint on error.
142pub 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
173/// Returns true if the `key` is a valid key for any engine or storage.
174pub 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    // Semantic-layer keys share a reserved prefix instead of a fixed allowlist so
192    // the vocabulary can grow without touching this gate. See `semantic` module.
193    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    /// Per-region write buffer stall threshold. Writes are rejected at twice this size.
204    pub write_buffer_size: Option<ReadableSize>,
205    /// Time-to-live of table. Expired data will be automatically purged.
206    pub ttl: Option<TimeToLive>,
207    /// Skip wal write for this table.
208    pub skip_wal: bool,
209    /// Extra options that may not applicable to all table engines.
210    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
222/// Table-level partition count hint consumed by the auto-repartition planner.
223pub const REPARTITION_PARTITION_NUM_HINT_KEY: &str = "repartition.partition.num.hint";
224
225impl TableOptions {
226    /// Returns the table data model, if specified.
227    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/// Alter table request
346#[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    // None in standalone.
354    pub table_version: Option<TableVersion>,
355}
356
357/// Add column request
358#[derive(Debug, Clone, Serialize, Deserialize)]
359pub struct AddColumnRequest {
360    pub column_schema: ColumnSchema,
361    pub is_key: bool,
362    pub location: Option<AddColumnLocation>,
363    /// Add column if not exists.
364    pub add_if_not_exists: bool,
365}
366
367/// Change column datatype request
368#[derive(Debug, Clone, Serialize, Deserialize)]
369pub struct ModifyColumnTypeRequest {
370    pub column_name: String,
371    pub target_type: ConcreteDataType,
372}
373
374/// Set JSON2 settings request.
375#[derive(Debug, Clone, Serialize, Deserialize)]
376pub struct SetJsonSettingsRequest {
377    pub column_name: String,
378    pub settings: JsonSettings,
379}
380
381/// A family of annotation table options: pure metadata markers that no region
382/// consumes. Setting or unsetting them only rewrites the table's
383/// `extra_options`, so the alter skips region dispatch entirely.
384#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
385pub enum AnnotationFamily {
386    /// `greptime.semantic.*` options (see the [`semantic`] module).
387    Semantic,
388    /// Column and partition count hints consumed by the auto-repartition planner.
389    RepartitionHint,
390}
391
392impl AnnotationFamily {
393    /// The key namespace used in diagnostics; accepted keys are classified by [`Self::of_key`].
394    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    /// Whether this family may be altered on logical metric tables. Only
415    /// families whose values nothing on the physical side consumes qualify;
416    /// the repartition hint drives physical region repartitioning.
417    pub fn allows_logical_tables(self) -> bool {
418        match self {
419            Self::Semantic => true,
420            Self::RepartitionHint => false,
421        }
422    }
423
424    /// Whether duplicate keys are rejected in this family's SET/UNSET batch.
425    pub fn requires_unique_keys(self) -> bool {
426        matches!(self, Self::RepartitionHint)
427    }
428
429    /// The error for a SET/UNSET batch mixing this family with other options.
430    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/// Why an annotation SET/UNSET key batch was rejected.
443#[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
458/// Validates a SET/UNSET batch and returns its annotation family.
459///
460/// Empty batches and batches containing only non-annotation keys return `None`.
461/// Annotation keys cannot share a batch with another family or region options.
462/// Duplicate keys are rejected only for families that require unique keys.
463pub 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
490/// Table shape an annotation option is validated against.
491pub 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/// Why an annotation option was rejected. Typed so each DDL entry point maps
498/// rules onto its existing error variants and status codes: ALTER keeps
499/// missing columns as `TableColumnNotFound` (4002), CREATE keeps its
500/// `InvalidArguments` family — the rules converge, the contracts do not.
501#[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
559/// Validates one annotation option and returns the value to store — the
560/// repartition hints are trimmed, semantic values pass
561/// through unchanged.
562pub(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
632/// CREATE-side entry: validates every annotation option present in `options`
633/// and writes normalized values back in place.
634pub 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    /// Returns the column name of the index option.
728    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    /// Returns the column name of the index option.
746    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    /// Whether this insert should skip WAL.
762    pub skip_wal: bool,
763}
764
765/// Delete (by primary key) request
766#[derive(Debug)]
767pub struct DeleteRequest {
768    pub catalog_name: String,
769    pub schema_name: String,
770    pub table_name: String,
771    /// Values of each column in this table's primary key and time index.
772    ///
773    /// The key is the column name, and the value is the column value.
774    pub key_column_values: HashMap<String, VectorRef>,
775}
776
777#[derive(Debug)]
778pub enum CopyDirection {
779    Export,
780    Import,
781}
782
783/// Copy table request
784#[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    /// The index build mode. Absent options select SST indexes.
808    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        // Only whitelisted semantic keys are accepted.
922        assert!(validate_table_option(SEMANTIC_SIGNAL_TYPE));
923        assert!(validate_table_option(SEMANTIC_METRIC_TYPE));
924        // Unknown semantic key, near-miss, and the internal transport key are rejected.
925        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}