Skip to main content

store_api/
region_request.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::HashMap;
16use std::fmt::{self, Display};
17use std::time::Duration;
18
19use api::helper::{ColumnDataTypeWrapper, from_pb_time_ranges, from_pb_time_unit};
20use api::v1::add_column_location::LocationType;
21use api::v1::column_def::{
22    as_fulltext_option_analyzer, as_fulltext_option_backend, as_skipping_index_type,
23};
24use api::v1::region::bulk_insert_request::Body;
25use api::v1::region::{
26    AlterRequest, AlterRequests, BuildIndexRequest, BulkInsertRequest,
27    CleanUpRequest as PbCleanUpRequest, CloseRequest, CompactRequest, CreateRequest,
28    CreateRequests, DeleteRequests, DropRequest, DropRequests, FlushRequest, InsertRequests,
29    OpenRequest, TruncateRequest, alter_request, compact_request, region_request, truncate_request,
30};
31use api::v1::{
32    self, Analyzer, ArrowIpc, FulltextBackend as PbFulltextBackend, Option as PbOption, Rows,
33    SemanticType, SkippingIndexType as PbSkippingIndexType, WriteHint,
34};
35pub use common_base::AffectedRows;
36use common_base::readable_size::ReadableSize;
37use common_grpc::flight::FlightDecoder;
38use common_recordbatch::DfRecordBatch;
39use common_time::range::TimestampRange;
40use common_time::{TimeToLive, Timestamp};
41use datatypes::prelude::ConcreteDataType;
42use datatypes::schema::{FulltextOptions, SkippingIndexOptions};
43use num_enum::TryFromPrimitive;
44use serde::{Deserialize, Serialize};
45use snafu::{OptionExt, ResultExt, ensure};
46use strum::{AsRefStr, IntoStaticStr};
47
48use crate::logstore::entry;
49use crate::metadata::{
50    ColumnMetadata, ConvertTimeRangesSnafu, DecodeProtoSnafu, FlightCodecSnafu,
51    InvalidIndexOptionSnafu, InvalidRawRegionRequestSnafu, InvalidRegionRequestSnafu,
52    InvalidSetRegionOptionRequestSnafu, InvalidUnsetRegionOptionRequestSnafu, MetadataError,
53    RegionMetadata, Result, UnexpectedSnafu,
54};
55use crate::metric_engine_consts::PHYSICAL_TABLE_METADATA_KEY;
56use crate::metrics;
57use crate::mito_engine_options::{
58    APPEND_MODE_KEY, AUTO_FLUSH_INTERVAL_KEY, MAX_ROW_GROUP_ROW_COUNT,
59    MAX_ROW_GROUP_ROW_COUNT_LIMIT, SKIP_WAL_KEY, SST_FORMAT_KEY, TTL_KEY,
60    TWCS_MAX_OUTPUT_FILE_SIZE, TWCS_TIME_WINDOW, TWCS_TRIGGER_FILE_NUM, WRITE_BUFFER_SIZE_KEY,
61};
62use crate::path_utils::table_dir;
63use crate::storage::{ColumnId, RegionId, ScanRequest};
64
65/// The type of path to generate.
66#[derive(Debug, Clone, Copy, PartialEq, TryFromPrimitive)]
67#[repr(u8)]
68pub enum PathType {
69    /// A bare path - the original path of an engine.
70    ///
71    /// The path prefix is `{table_dir}/{table_id}_{region_sequence}/`.
72    Bare,
73    /// A path for the data region of a metric engine table.
74    ///
75    /// The path prefix is `{table_dir}/{table_id}_{region_sequence}/data/`.
76    Data,
77    /// A path for the metadata region of a metric engine table.
78    ///
79    /// The path prefix is `{table_dir}/{table_id}_{region_sequence}/metadata/`.
80    Metadata,
81}
82
83#[derive(Debug, IntoStaticStr)]
84pub enum BatchRegionDdlRequest {
85    Create(Vec<(RegionId, RegionCreateRequest)>),
86    Drop(Vec<(RegionId, RegionDropRequest)>),
87    Alter(Vec<(RegionId, RegionAlterRequest)>),
88}
89
90impl BatchRegionDdlRequest {
91    /// Converts [Body](region_request::Body) to [`BatchRegionDdlRequest`].
92    pub fn try_from_request_body(body: region_request::Body) -> Result<Option<Self>> {
93        match body {
94            region_request::Body::Creates(creates) => {
95                let requests = creates
96                    .requests
97                    .into_iter()
98                    .map(parse_region_create)
99                    .collect::<Result<Vec<_>>>()?;
100                Ok(Some(Self::Create(requests)))
101            }
102            region_request::Body::Drops(drops) => {
103                let requests = drops
104                    .requests
105                    .into_iter()
106                    .map(parse_region_drop)
107                    .collect::<Result<Vec<_>>>()?;
108                Ok(Some(Self::Drop(requests)))
109            }
110            region_request::Body::Alters(alters) => {
111                let requests = alters
112                    .requests
113                    .into_iter()
114                    .map(parse_region_alter)
115                    .collect::<Result<Vec<_>>>()?;
116                Ok(Some(Self::Alter(requests)))
117            }
118            _ => Ok(None),
119        }
120    }
121
122    pub fn request_type(&self) -> &'static str {
123        self.into()
124    }
125
126    pub fn into_region_requests(self) -> Vec<(RegionId, RegionRequest)> {
127        match self {
128            Self::Create(requests) => requests
129                .into_iter()
130                .map(|(region_id, request)| (region_id, RegionRequest::Create(request)))
131                .collect(),
132            Self::Drop(requests) => requests
133                .into_iter()
134                .map(|(region_id, request)| (region_id, RegionRequest::Drop(request)))
135                .collect(),
136            Self::Alter(requests) => requests
137                .into_iter()
138                .map(|(region_id, request)| (region_id, RegionRequest::Alter(request)))
139                .collect(),
140        }
141    }
142}
143
144#[derive(Debug, IntoStaticStr)]
145pub enum RegionRequest {
146    Put(RegionPutRequest),
147    Delete(RegionDeleteRequest),
148    Create(RegionCreateRequest),
149    Drop(RegionDropRequest),
150    Open(RegionOpenRequest),
151    CleanUp(RegionCleanUpRequest),
152    Close(RegionCloseRequest),
153    Alter(RegionAlterRequest),
154    Flush(RegionFlushRequest),
155    Compact(RegionCompactRequest),
156    BuildIndex(RegionBuildIndexRequest),
157    Truncate(RegionTruncateRequest),
158    Catchup(RegionCatchupRequest),
159    BulkInserts(RegionBulkInsertsRequest),
160    EnterStaging(EnterStagingRequest),
161    ApplyStagingManifest(ApplyStagingManifestRequest),
162}
163
164impl RegionRequest {
165    /// Convert [Body](region_request::Body) to a group of [RegionRequest] with region id.
166    /// Inserts/Deletes request might become multiple requests. Others are one-to-one.
167    pub fn try_from_request_body(body: region_request::Body) -> Result<Vec<(RegionId, Self)>> {
168        match body {
169            region_request::Body::Inserts(inserts) => make_region_puts(inserts),
170            region_request::Body::Deletes(deletes) => make_region_deletes(deletes),
171            region_request::Body::Create(create) => make_region_create(create),
172            region_request::Body::Drop(drop) => make_region_drop(drop),
173            region_request::Body::Open(open) => make_region_open(open),
174            region_request::Body::CleanUp(clean_up) => make_region_clean_up(clean_up),
175            region_request::Body::Close(close) => make_region_close(close),
176            region_request::Body::Alter(alter) => make_region_alter(alter),
177            region_request::Body::Flush(flush) => make_region_flush(flush),
178            region_request::Body::Compact(compact) => make_region_compact(compact),
179            region_request::Body::BuildIndex(index) => make_region_build_index(index),
180            region_request::Body::Truncate(truncate) => make_region_truncate(truncate),
181            region_request::Body::Creates(creates) => make_region_creates(creates),
182            region_request::Body::Drops(drops) => make_region_drops(drops),
183            region_request::Body::Alters(alters) => make_region_alters(alters),
184            region_request::Body::BulkInsert(bulk) => make_region_bulk_inserts(bulk),
185            region_request::Body::Sync(_) => UnexpectedSnafu {
186                reason: "Sync request should be handled separately by RegionServer",
187            }
188            .fail(),
189            region_request::Body::ListMetadata(_) => UnexpectedSnafu {
190                reason: "ListMetadata request should be handled separately by RegionServer",
191            }
192            .fail(),
193            region_request::Body::RemoteDynFilter(_) => UnexpectedSnafu {
194                reason: "RemoteDynFilter request should be handled separately by RegionServer",
195            }
196            .fail(),
197            region_request::Body::ApplyStagingManifest(apply) => {
198                make_region_apply_staging_manifest(apply)
199            }
200        }
201    }
202
203    /// Returns the type name of the request.
204    pub fn request_type(&self) -> &'static str {
205        self.into()
206    }
207}
208
209fn make_region_puts(inserts: InsertRequests) -> Result<Vec<(RegionId, RegionRequest)>> {
210    let requests = inserts
211        .requests
212        .into_iter()
213        .filter_map(|r| {
214            let region_id = r.region_id.into();
215            r.rows.map(|rows| {
216                (
217                    region_id,
218                    RegionRequest::Put(RegionPutRequest {
219                        rows,
220                        hint: None,
221                        partition_expr_version: r.partition_expr_version.map(|v| v.value),
222                    }),
223                )
224            })
225        })
226        .collect();
227    Ok(requests)
228}
229
230fn make_region_deletes(deletes: DeleteRequests) -> Result<Vec<(RegionId, RegionRequest)>> {
231    let requests = deletes
232        .requests
233        .into_iter()
234        .filter_map(|r| {
235            let region_id = r.region_id.into();
236            r.rows.map(|rows| {
237                (
238                    region_id,
239                    RegionRequest::Delete(RegionDeleteRequest {
240                        rows,
241                        hint: None,
242                        partition_expr_version: r.partition_expr_version.map(|v| v.value),
243                    }),
244                )
245            })
246        })
247        .collect();
248    Ok(requests)
249}
250
251fn parse_region_create(create: CreateRequest) -> Result<(RegionId, RegionCreateRequest)> {
252    let column_metadatas = create
253        .column_defs
254        .into_iter()
255        .map(ColumnMetadata::try_from_column_def)
256        .collect::<Result<Vec<_>>>()?;
257    let region_id = RegionId::from(create.region_id);
258    let table_dir = table_dir(&create.path, region_id.table_id());
259    let partition_expr_json = create.partition.as_ref().map(|p| p.expression.clone());
260    Ok((
261        region_id,
262        RegionCreateRequest {
263            engine: create.engine,
264            column_metadatas,
265            primary_key: create.primary_key,
266            options: create.options,
267            table_dir,
268            path_type: PathType::Bare,
269            partition_expr_json,
270            requirements: create
271                .requirements
272                .map(RegionRequirements::from)
273                .unwrap_or_default(),
274        },
275    ))
276}
277
278fn make_region_create(create: CreateRequest) -> Result<Vec<(RegionId, RegionRequest)>> {
279    let (region_id, request) = parse_region_create(create)?;
280    Ok(vec![(region_id, RegionRequest::Create(request))])
281}
282
283fn make_region_creates(creates: CreateRequests) -> Result<Vec<(RegionId, RegionRequest)>> {
284    let mut requests = Vec::with_capacity(creates.requests.len());
285    for create in creates.requests {
286        requests.extend(make_region_create(create)?);
287    }
288    Ok(requests)
289}
290
291fn parse_region_drop(drop: DropRequest) -> Result<(RegionId, RegionDropRequest)> {
292    let region_id = drop.region_id.into();
293    Ok((
294        region_id,
295        RegionDropRequest {
296            fast_path: drop.fast_path,
297            force: drop.force,
298            partial_drop: drop.partial_drop,
299        },
300    ))
301}
302
303fn make_region_drop(drop: DropRequest) -> Result<Vec<(RegionId, RegionRequest)>> {
304    let (region_id, request) = parse_region_drop(drop)?;
305    Ok(vec![(region_id, RegionRequest::Drop(request))])
306}
307
308fn make_region_drops(drops: DropRequests) -> Result<Vec<(RegionId, RegionRequest)>> {
309    let mut requests = Vec::with_capacity(drops.requests.len());
310    for drop in drops.requests {
311        requests.extend(make_region_drop(drop)?);
312    }
313    Ok(requests)
314}
315
316fn make_region_open(open: OpenRequest) -> Result<Vec<(RegionId, RegionRequest)>> {
317    let region_id = RegionId::from(open.region_id);
318    let table_dir = table_dir(&open.path, region_id.table_id());
319    Ok(vec![(
320        region_id,
321        RegionRequest::Open(RegionOpenRequest {
322            engine: open.engine,
323            table_dir,
324            path_type: PathType::Bare,
325            options: open.options,
326            skip_wal_replay: false,
327            checkpoint: None,
328            requirements: Default::default(),
329        }),
330    )])
331}
332
333fn make_region_clean_up(clean_up: PbCleanUpRequest) -> Result<Vec<(RegionId, RegionRequest)>> {
334    let region_id = RegionId::from(clean_up.region_id);
335    let table_dir = table_dir(&clean_up.path, region_id.table_id());
336    Ok(vec![(
337        region_id,
338        RegionRequest::CleanUp(RegionCleanUpRequest {
339            engine: clean_up.engine,
340            table_dir,
341            path_type: PathType::Bare,
342            options: clean_up.options,
343        }),
344    )])
345}
346
347fn make_region_close(close: CloseRequest) -> Result<Vec<(RegionId, RegionRequest)>> {
348    let region_id = close.region_id.into();
349    Ok(vec![(
350        region_id,
351        RegionRequest::Close(RegionCloseRequest {
352            flush_on_close: close.flush_on_close,
353        }),
354    )])
355}
356
357fn parse_region_alter(alter: AlterRequest) -> Result<(RegionId, RegionAlterRequest)> {
358    let region_id = alter.region_id.into();
359    let request = RegionAlterRequest::try_from(alter)?;
360    Ok((region_id, request))
361}
362
363fn make_region_alter(alter: AlterRequest) -> Result<Vec<(RegionId, RegionRequest)>> {
364    let (region_id, request) = parse_region_alter(alter)?;
365    Ok(vec![(region_id, RegionRequest::Alter(request))])
366}
367
368fn make_region_alters(alters: AlterRequests) -> Result<Vec<(RegionId, RegionRequest)>> {
369    let mut requests = Vec::with_capacity(alters.requests.len());
370    for alter in alters.requests {
371        requests.extend(make_region_alter(alter)?);
372    }
373    Ok(requests)
374}
375
376fn make_region_flush(flush: FlushRequest) -> Result<Vec<(RegionId, RegionRequest)>> {
377    let region_id = flush.region_id.into();
378    Ok(vec![(
379        region_id,
380        RegionRequest::Flush(RegionFlushRequest::default()),
381    )])
382}
383
384fn make_region_compact(compact: CompactRequest) -> Result<Vec<(RegionId, RegionRequest)>> {
385    let region_id = compact.region_id.into();
386    let options = compact
387        .options
388        .unwrap_or(compact_request::Options::Regular(Default::default()));
389    // Convert parallelism: a value of 0 indicates no specific parallelism requested (None)
390    let parallelism = if compact.parallelism == 0 {
391        None
392    } else {
393        Some(compact.parallelism)
394    };
395    let time_range = compact
396        .time_range
397        .map(|range| {
398            let time_unit = v1::TimeUnit::try_from(range.time_unit).map_err(|_| {
399                InvalidRegionRequestSnafu {
400                    region_id,
401                    err: format!("invalid compaction time unit: {}", range.time_unit),
402                }
403                .build()
404            })?;
405            let time_unit = from_pb_time_unit(time_unit);
406            let start = Timestamp::new(range.start, time_unit);
407            let end = Timestamp::new(range.end, time_unit);
408            ensure!(
409                start < end,
410                InvalidRegionRequestSnafu {
411                    region_id,
412                    err: format!(
413                        "compaction start time must be earlier than end time: {} >= {}",
414                        range.start, range.end
415                    ),
416                }
417            );
418            TimestampRange::new(start, end).with_context(|| InvalidRegionRequestSnafu {
419                region_id,
420                err: "invalid compaction time range".to_string(),
421            })
422        })
423        .transpose()?;
424    Ok(vec![(
425        region_id,
426        RegionRequest::Compact(RegionCompactRequest {
427            options,
428            parallelism,
429            time_range,
430        }),
431    )])
432}
433
434fn make_region_build_index(index: BuildIndexRequest) -> Result<Vec<(RegionId, RegionRequest)>> {
435    let region_id = index.region_id.into();
436    Ok(vec![(
437        region_id,
438        RegionRequest::BuildIndex(RegionBuildIndexRequest {}),
439    )])
440}
441
442fn make_region_truncate(truncate: TruncateRequest) -> Result<Vec<(RegionId, RegionRequest)>> {
443    let region_id = truncate.region_id.into();
444    match truncate.kind {
445        None => InvalidRawRegionRequestSnafu {
446            err: "missing kind in TruncateRequest".to_string(),
447        }
448        .fail(),
449        Some(truncate_request::Kind::All(_)) => Ok(vec![(
450            region_id,
451            RegionRequest::Truncate(RegionTruncateRequest::All),
452        )]),
453        Some(truncate_request::Kind::TimeRanges(time_ranges)) => {
454            let time_ranges = from_pb_time_ranges(time_ranges).context(ConvertTimeRangesSnafu)?;
455
456            Ok(vec![(
457                region_id,
458                RegionRequest::Truncate(RegionTruncateRequest::ByTimeRanges { time_ranges }),
459            )])
460        }
461        Some(truncate_request::Kind::Unflushed(_)) => Ok(vec![(
462            region_id,
463            RegionRequest::Truncate(RegionTruncateRequest::Unflushed),
464        )]),
465    }
466}
467
468/// Convert [BulkInsertRequest] to [RegionRequest] and group by [RegionId].
469fn make_region_bulk_inserts(request: BulkInsertRequest) -> Result<Vec<(RegionId, RegionRequest)>> {
470    let region_id = request.region_id.into();
471    let partition_expr_version = request.partition_expr_version.map(|v| v.value);
472    let aligned_schema_version = request.aligned_schema_version.map(|v| v.schema_version);
473    let Some(Body::ArrowIpc(request)) = request.body else {
474        return Ok(vec![]);
475    };
476
477    let decoder_timer = metrics::CONVERT_REGION_BULK_REQUEST
478        .with_label_values(&["decode"])
479        .start_timer();
480    let mut decoder =
481        FlightDecoder::try_from_schema_bytes(&request.schema).context(FlightCodecSnafu)?;
482    let payload = decoder
483        .try_decode_record_batch(&request.data_header, &request.payload)
484        .context(FlightCodecSnafu)?;
485    decoder_timer.observe_duration();
486    Ok(vec![(
487        region_id,
488        RegionRequest::BulkInserts(RegionBulkInsertsRequest {
489            region_id,
490            payload,
491            raw_data: request,
492            partition_expr_version,
493            aligned_schema_version,
494        }),
495    )])
496}
497
498fn make_region_apply_staging_manifest(
499    api::v1::region::ApplyStagingManifestRequest {
500        region_id,
501        partition_expr,
502        central_region_id,
503        manifest_path,
504    }: api::v1::region::ApplyStagingManifestRequest,
505) -> Result<Vec<(RegionId, RegionRequest)>> {
506    let region_id = region_id.into();
507    Ok(vec![(
508        region_id,
509        RegionRequest::ApplyStagingManifest(ApplyStagingManifestRequest {
510            partition_expr,
511            central_region_id: central_region_id.into(),
512            manifest_path,
513        }),
514    )])
515}
516
517/// Request to put data into a region.
518#[derive(Debug)]
519pub struct RegionPutRequest {
520    /// Rows to put.
521    pub rows: Rows,
522    /// Write hint.
523    pub hint: Option<WriteHint>,
524    /// Partition expression version for the region.
525    pub partition_expr_version: Option<u64>,
526}
527
528#[derive(Debug)]
529pub struct RegionReadRequest {
530    pub request: ScanRequest,
531}
532
533/// Request to delete data from a region.
534#[derive(Debug)]
535pub struct RegionDeleteRequest {
536    /// Keys to rows to delete.
537    ///
538    /// Each row only contains primary key columns and a time index column.
539    pub rows: Rows,
540    /// Write hint.
541    pub hint: Option<WriteHint>,
542    /// Partition expression version for the region.
543    pub partition_expr_version: Option<u64>,
544}
545
546#[derive(Debug, Clone)]
547pub struct RegionCreateRequest {
548    /// Region engine name
549    pub engine: String,
550    /// Columns in this region.
551    pub column_metadatas: Vec<ColumnMetadata>,
552    /// Columns in the primary key.
553    pub primary_key: Vec<ColumnId>,
554    /// Options of the created region.
555    pub options: HashMap<String, String>,
556    /// Directory for table's data home. Usually is composed by catalog and table id
557    pub table_dir: String,
558    /// Path type for generating paths
559    pub path_type: PathType,
560    /// Partition expression JSON from table metadata. Set to empty string for a region without partition.
561    /// `Option` to keep compatibility with old clients.
562    pub partition_expr_json: Option<String>,
563    /// Requirements for creating the region.
564    pub requirements: RegionRequirements,
565}
566
567impl RegionCreateRequest {
568    /// Checks whether the request is valid, returns an error if it is invalid.
569    pub fn validate(&self) -> Result<()> {
570        // time index must exist
571        ensure!(
572            self.column_metadatas
573                .iter()
574                .any(|x| x.semantic_type == SemanticType::Timestamp),
575            InvalidRegionRequestSnafu {
576                region_id: RegionId::new(0, 0),
577                err: "missing timestamp column in create region request".to_string(),
578            }
579        );
580
581        // build column id to indices
582        let mut column_id_to_indices = HashMap::with_capacity(self.column_metadatas.len());
583        for (i, c) in self.column_metadatas.iter().enumerate() {
584            if let Some(previous) = column_id_to_indices.insert(c.column_id, i) {
585                return InvalidRegionRequestSnafu {
586                    region_id: RegionId::new(0, 0),
587                    err: format!(
588                        "duplicate column id {} (at position {} and {}) in create region request",
589                        c.column_id, previous, i
590                    ),
591                }
592                .fail();
593            }
594        }
595
596        // primary key must exist
597        for column_id in &self.primary_key {
598            ensure!(
599                column_id_to_indices.contains_key(column_id),
600                InvalidRegionRequestSnafu {
601                    region_id: RegionId::new(0, 0),
602                    err: format!(
603                        "missing primary key column {} in create region request",
604                        column_id
605                    ),
606                }
607            );
608        }
609
610        Ok(())
611    }
612
613    /// Returns true when the region belongs to the metric engine's physical table.
614    pub fn is_physical_table(&self) -> bool {
615        self.options.contains_key(PHYSICAL_TABLE_METADATA_KEY)
616    }
617}
618
619#[derive(Debug, Clone)]
620pub struct RegionDropRequest {
621    /// Enables fast-path drop optimizations for logical regions.
622    /// Only applicable to the Metric Engine; ignored by others.
623    pub fast_path: bool,
624
625    /// Forces the drop of a physical region and all its associated logical regions.
626    /// Only relevant for physical regions managed by the Metric Engine.
627    pub force: bool,
628
629    /// If true, indicates that only a portion of the region is being dropped, and files may still be referenced by other regions.
630    /// This is used to prevent deletion of files that are still in use by other regions.
631    pub partial_drop: bool,
632}
633
634/// Requirements for a region request.
635#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
636#[serde(default)]
637pub struct RegionRequirements {
638    /// Whether the region data must be backed by object storage.
639    pub object_storage: bool,
640}
641
642impl RegionRequirements {
643    /// Returns empty requirements.
644    pub fn empty() -> Self {
645        Self::default()
646    }
647
648    /// Returns requirements for object storage.
649    pub fn object_storage() -> Self {
650        Self {
651            object_storage: true,
652        }
653    }
654}
655
656impl From<api::v1::region::RegionRequirements> for RegionRequirements {
657    fn from(value: api::v1::region::RegionRequirements) -> Self {
658        Self {
659            object_storage: value.object_storage,
660        }
661    }
662}
663
664impl From<RegionRequirements> for api::v1::region::RegionRequirements {
665    fn from(value: RegionRequirements) -> Self {
666        Self {
667            object_storage: value.object_storage,
668        }
669    }
670}
671
672/// Open region request.
673#[derive(Debug, Clone)]
674pub struct RegionOpenRequest {
675    /// Region engine name
676    pub engine: String,
677    /// Directory for table's data home. Usually is composed by catalog and table id
678    pub table_dir: String,
679    /// Path type for generating paths
680    pub path_type: PathType,
681    /// Options of the opened region.
682    pub options: HashMap<String, String>,
683    /// To skip replaying the WAL.
684    pub skip_wal_replay: bool,
685    /// Replay checkpoint.
686    pub checkpoint: Option<ReplayCheckpoint>,
687    /// Requirements for opening the region.
688    pub requirements: RegionRequirements,
689}
690
691#[derive(Debug, Clone, Copy, PartialEq, Eq)]
692pub struct ReplayCheckpoint {
693    pub entry_id: u64,
694    pub metadata_entry_id: Option<u64>,
695}
696
697impl RegionOpenRequest {
698    /// Returns true when the region belongs to the metric engine's physical table.
699    pub fn is_physical_table(&self) -> bool {
700        self.options.contains_key(PHYSICAL_TABLE_METADATA_KEY)
701    }
702}
703
704/// Offline region cleanup request.
705#[derive(Debug, Clone)]
706pub struct RegionCleanUpRequest {
707    /// Region engine name
708    pub engine: String,
709    /// Directory for table's data home. Usually is composed by catalog and table id
710    pub table_dir: String,
711    /// Path type for generating paths
712    pub path_type: PathType,
713    /// Options of the cleaned region.
714    pub options: HashMap<String, String>,
715}
716
717impl RegionCleanUpRequest {
718    /// Returns true when the region belongs to the metric engine's physical table.
719    pub fn is_physical_table(&self) -> bool {
720        self.options.contains_key(PHYSICAL_TABLE_METADATA_KEY)
721    }
722}
723
724/// Close region request.
725#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
726pub struct RegionCloseRequest {
727    /// Whether to flush the region before closing it.
728    pub flush_on_close: bool,
729}
730
731/// Alter metadata of a region.
732#[derive(Debug, PartialEq, Eq, Clone)]
733pub struct RegionAlterRequest {
734    /// Kind of alteration to do.
735    pub kind: AlterKind,
736}
737
738impl RegionAlterRequest {
739    /// Checks whether the request is valid, returns an error if it is invalid.
740    pub fn validate(&self, metadata: &RegionMetadata) -> Result<()> {
741        self.kind.validate(metadata)?;
742
743        Ok(())
744    }
745
746    /// Returns true if we need to apply the request to the region.
747    ///
748    /// The `request` should be valid.
749    pub fn need_alter(&self, metadata: &RegionMetadata) -> bool {
750        debug_assert!(self.validate(metadata).is_ok());
751        self.kind.need_alter(metadata)
752    }
753}
754
755impl TryFrom<AlterRequest> for RegionAlterRequest {
756    type Error = MetadataError;
757
758    fn try_from(value: AlterRequest) -> Result<Self> {
759        let kind = value.kind.context(InvalidRawRegionRequestSnafu {
760            err: "missing kind in AlterRequest",
761        })?;
762
763        let kind = AlterKind::try_from(kind)?;
764        Ok(RegionAlterRequest { kind })
765    }
766}
767
768/// Kind of the alteration.
769#[derive(Debug, PartialEq, Eq, Clone, AsRefStr)]
770pub enum AlterKind {
771    /// Add columns to the region.
772    AddColumns {
773        /// Columns to add.
774        columns: Vec<AddColumn>,
775    },
776    /// Drop columns from the region, only fields are allowed to drop.
777    DropColumns {
778        /// Name of columns to drop.
779        names: Vec<String>,
780    },
781    /// Change columns datatype form the region, only fields are allowed to change.
782    ModifyColumnTypes {
783        /// Columns to change.
784        columns: Vec<ModifyColumnType>,
785    },
786    /// Set region options.
787    SetRegionOptions { options: Vec<SetRegionOption> },
788    /// Unset region options.
789    UnsetRegionOptions { keys: Vec<UnsetRegionOption> },
790    /// Set index options.
791    SetIndexes { options: Vec<SetIndexOption> },
792    /// Unset index options.
793    UnsetIndexes { options: Vec<UnsetIndexOption> },
794    /// Drop column default value.
795    DropDefaults {
796        /// Name of columns to drop.
797        names: Vec<String>,
798    },
799    /// Set column default value.
800    SetDefaults {
801        /// Columns to change.
802        columns: Vec<SetDefault>,
803    },
804    /// Sync column metadatas.
805    SyncColumns {
806        column_metadatas: Vec<ColumnMetadata>,
807    },
808}
809#[derive(Debug, PartialEq, Eq, Clone)]
810pub struct SetDefault {
811    pub name: String,
812    pub default_constraint: Vec<u8>,
813}
814
815#[derive(Debug, PartialEq, Eq, Clone)]
816pub enum SetIndexOption {
817    Fulltext {
818        column_name: String,
819        options: FulltextOptions,
820    },
821    Inverted {
822        column_name: String,
823    },
824    Skipping {
825        column_name: String,
826        options: SkippingIndexOptions,
827    },
828}
829
830impl SetIndexOption {
831    /// Returns the column name of the index option.
832    pub fn column_name(&self) -> &String {
833        match self {
834            SetIndexOption::Fulltext { column_name, .. } => column_name,
835            SetIndexOption::Inverted { column_name } => column_name,
836            SetIndexOption::Skipping { column_name, .. } => column_name,
837        }
838    }
839
840    /// Returns true if the index option is fulltext.
841    pub fn is_fulltext(&self) -> bool {
842        match self {
843            SetIndexOption::Fulltext { .. } => true,
844            SetIndexOption::Inverted { .. } => false,
845            SetIndexOption::Skipping { .. } => false,
846        }
847    }
848}
849
850impl TryFrom<v1::SetIndex> for SetIndexOption {
851    type Error = MetadataError;
852
853    fn try_from(value: v1::SetIndex) -> Result<Self> {
854        let option = value.options.context(InvalidRawRegionRequestSnafu {
855            err: "missing options in SetIndex",
856        })?;
857
858        let opt = match option {
859            v1::set_index::Options::Fulltext(x) => SetIndexOption::Fulltext {
860                column_name: x.column_name.clone(),
861                options: FulltextOptions::new(
862                    x.enable,
863                    as_fulltext_option_analyzer(
864                        Analyzer::try_from(x.analyzer).context(DecodeProtoSnafu)?,
865                    ),
866                    x.case_sensitive,
867                    as_fulltext_option_backend(
868                        PbFulltextBackend::try_from(x.backend).context(DecodeProtoSnafu)?,
869                    ),
870                    x.granularity as u32,
871                    x.false_positive_rate,
872                )
873                .context(InvalidIndexOptionSnafu)?,
874            },
875            v1::set_index::Options::Inverted(i) => SetIndexOption::Inverted {
876                column_name: i.column_name,
877            },
878            v1::set_index::Options::Skipping(s) => SetIndexOption::Skipping {
879                column_name: s.column_name,
880                options: SkippingIndexOptions::new(
881                    s.granularity as u32,
882                    s.false_positive_rate,
883                    as_skipping_index_type(
884                        PbSkippingIndexType::try_from(s.skipping_index_type)
885                            .context(DecodeProtoSnafu)?,
886                    ),
887                )
888                .context(InvalidIndexOptionSnafu)?,
889            },
890        };
891
892        Ok(opt)
893    }
894}
895
896#[derive(Debug, PartialEq, Eq, Clone)]
897pub enum UnsetIndexOption {
898    Fulltext { column_name: String },
899    Inverted { column_name: String },
900    Skipping { column_name: String },
901}
902
903impl UnsetIndexOption {
904    pub fn column_name(&self) -> &String {
905        match self {
906            UnsetIndexOption::Fulltext { column_name } => column_name,
907            UnsetIndexOption::Inverted { column_name } => column_name,
908            UnsetIndexOption::Skipping { column_name } => column_name,
909        }
910    }
911
912    pub fn is_fulltext(&self) -> bool {
913        match self {
914            UnsetIndexOption::Fulltext { .. } => true,
915            UnsetIndexOption::Inverted { .. } => false,
916            UnsetIndexOption::Skipping { .. } => false,
917        }
918    }
919}
920
921impl TryFrom<v1::UnsetIndex> for UnsetIndexOption {
922    type Error = MetadataError;
923
924    fn try_from(value: v1::UnsetIndex) -> Result<Self> {
925        let option = value.options.context(InvalidRawRegionRequestSnafu {
926            err: "missing options in UnsetIndex",
927        })?;
928
929        let opt = match option {
930            v1::unset_index::Options::Fulltext(f) => UnsetIndexOption::Fulltext {
931                column_name: f.column_name,
932            },
933            v1::unset_index::Options::Inverted(i) => UnsetIndexOption::Inverted {
934                column_name: i.column_name,
935            },
936            v1::unset_index::Options::Skipping(s) => UnsetIndexOption::Skipping {
937                column_name: s.column_name,
938            },
939        };
940
941        Ok(opt)
942    }
943}
944
945impl AlterKind {
946    /// Returns an error if the alter kind is invalid.
947    ///
948    /// It allows adding column if not exists and dropping column if exists.
949    pub fn validate(&self, metadata: &RegionMetadata) -> Result<()> {
950        match self {
951            AlterKind::AddColumns { columns } => {
952                for col_to_add in columns {
953                    col_to_add.validate(metadata)?;
954                }
955            }
956            AlterKind::DropColumns { names } => {
957                for name in names {
958                    Self::validate_column_to_drop(name, metadata)?;
959                }
960            }
961            AlterKind::ModifyColumnTypes { columns } => {
962                for col_to_change in columns {
963                    col_to_change.validate(metadata)?;
964                }
965            }
966            AlterKind::SetRegionOptions { .. } => {}
967            AlterKind::UnsetRegionOptions { .. } => {}
968            AlterKind::SetIndexes { options } => {
969                for option in options {
970                    Self::validate_column_alter_index_option(
971                        option.column_name(),
972                        metadata,
973                        option.is_fulltext(),
974                    )?;
975                }
976            }
977            AlterKind::UnsetIndexes { options } => {
978                for option in options {
979                    Self::validate_column_alter_index_option(
980                        option.column_name(),
981                        metadata,
982                        option.is_fulltext(),
983                    )?;
984                }
985            }
986            AlterKind::DropDefaults { names } => {
987                names
988                    .iter()
989                    .try_for_each(|name| Self::validate_column_existence(name, metadata))?;
990            }
991            AlterKind::SetDefaults { columns } => {
992                columns
993                    .iter()
994                    .try_for_each(|col| Self::validate_column_existence(&col.name, metadata))?;
995            }
996            AlterKind::SyncColumns { column_metadatas } => {
997                let new_primary_keys = column_metadatas
998                    .iter()
999                    .filter(|c| c.semantic_type == SemanticType::Tag)
1000                    .map(|c| (c.column_schema.name.as_str(), c.column_id))
1001                    .collect::<HashMap<_, _>>();
1002
1003                let old_primary_keys = metadata
1004                    .column_metadatas
1005                    .iter()
1006                    .filter(|c| c.semantic_type == SemanticType::Tag)
1007                    .map(|c| (c.column_schema.name.as_str(), c.column_id));
1008
1009                for (name, id) in old_primary_keys {
1010                    let primary_key =
1011                        new_primary_keys
1012                            .get(name)
1013                            .with_context(|| InvalidRegionRequestSnafu {
1014                                region_id: metadata.region_id,
1015                                err: format!("column {} is not a primary key", name),
1016                            })?;
1017
1018                    ensure!(
1019                        *primary_key == id,
1020                        InvalidRegionRequestSnafu {
1021                            region_id: metadata.region_id,
1022                            err: format!(
1023                                "column with same name {} has different id, existing: {}, got: {}",
1024                                name, id, primary_key
1025                            ),
1026                        }
1027                    );
1028                }
1029
1030                let new_ts_column = column_metadatas
1031                    .iter()
1032                    .find(|c| c.semantic_type == SemanticType::Timestamp)
1033                    .map(|c| (c.column_schema.name.as_str(), c.column_id))
1034                    .context(InvalidRegionRequestSnafu {
1035                        region_id: metadata.region_id,
1036                        err: "timestamp column not found",
1037                    })?;
1038
1039                // Safety: timestamp column must exist.
1040                let old_ts_column = metadata
1041                    .column_metadatas
1042                    .iter()
1043                    .find(|c| c.semantic_type == SemanticType::Timestamp)
1044                    .map(|c| (c.column_schema.name.as_str(), c.column_id))
1045                    .unwrap();
1046
1047                ensure!(
1048                    new_ts_column == old_ts_column,
1049                    InvalidRegionRequestSnafu {
1050                        region_id: metadata.region_id,
1051                        err: format!(
1052                            "timestamp column {} has different id, existing: {}, got: {}",
1053                            old_ts_column.0, old_ts_column.1, new_ts_column.1
1054                        ),
1055                    }
1056                );
1057            }
1058        }
1059        Ok(())
1060    }
1061
1062    /// Returns true if we need to apply the alteration to the region.
1063    pub fn need_alter(&self, metadata: &RegionMetadata) -> bool {
1064        debug_assert!(self.validate(metadata).is_ok());
1065        match self {
1066            AlterKind::AddColumns { columns } => columns
1067                .iter()
1068                .any(|col_to_add| col_to_add.need_alter(metadata)),
1069            AlterKind::DropColumns { names } => names
1070                .iter()
1071                .any(|name| metadata.column_by_name(name).is_some()),
1072            AlterKind::ModifyColumnTypes { columns } => columns
1073                .iter()
1074                .any(|col_to_change| col_to_change.need_alter(metadata)),
1075            AlterKind::SetRegionOptions { .. } => true,
1076            AlterKind::UnsetRegionOptions { .. } => true,
1077            AlterKind::SetIndexes { options, .. } => options
1078                .iter()
1079                .any(|option| metadata.column_by_name(option.column_name()).is_some()),
1080            AlterKind::UnsetIndexes { options } => options
1081                .iter()
1082                .any(|option| metadata.column_by_name(option.column_name()).is_some()),
1083            AlterKind::DropDefaults { names } => names
1084                .iter()
1085                .any(|name| metadata.column_by_name(name).is_some()),
1086
1087            AlterKind::SetDefaults { columns } => columns
1088                .iter()
1089                .any(|x| metadata.column_by_name(&x.name).is_some()),
1090            AlterKind::SyncColumns { column_metadatas } => {
1091                metadata.column_metadatas != *column_metadatas
1092            }
1093        }
1094    }
1095
1096    /// Returns an error if the column to drop is invalid.
1097    fn validate_column_to_drop(name: &str, metadata: &RegionMetadata) -> Result<()> {
1098        let Some(column) = metadata.column_by_name(name) else {
1099            return Ok(());
1100        };
1101        ensure!(
1102            column.semantic_type == SemanticType::Field,
1103            InvalidRegionRequestSnafu {
1104                region_id: metadata.region_id,
1105                err: format!("column {} is not a field and could not be dropped", name),
1106            }
1107        );
1108        Ok(())
1109    }
1110
1111    /// Returns an error if the column's alter index option is invalid.
1112    fn validate_column_alter_index_option(
1113        column_name: &String,
1114        metadata: &RegionMetadata,
1115        is_fulltext: bool,
1116    ) -> Result<()> {
1117        let column = metadata
1118            .column_by_name(column_name)
1119            .context(InvalidRegionRequestSnafu {
1120                region_id: metadata.region_id,
1121                err: format!("column {} not found", column_name),
1122            })?;
1123
1124        if is_fulltext {
1125            ensure!(
1126                column.column_schema.data_type.is_string(),
1127                InvalidRegionRequestSnafu {
1128                    region_id: metadata.region_id,
1129                    err: format!(
1130                        "cannot change alter index options for non-string column {}",
1131                        column_name
1132                    ),
1133                }
1134            );
1135        }
1136
1137        Ok(())
1138    }
1139
1140    /// Returns an error if the column isn't exist.
1141    fn validate_column_existence(column_name: &String, metadata: &RegionMetadata) -> Result<()> {
1142        metadata
1143            .column_by_name(column_name)
1144            .context(InvalidRegionRequestSnafu {
1145                region_id: metadata.region_id,
1146                err: format!("column {} not found", column_name),
1147            })?;
1148
1149        Ok(())
1150    }
1151}
1152
1153impl TryFrom<alter_request::Kind> for AlterKind {
1154    type Error = MetadataError;
1155
1156    fn try_from(kind: alter_request::Kind) -> Result<Self> {
1157        let alter_kind = match kind {
1158            alter_request::Kind::AddColumns(x) => {
1159                let columns = x
1160                    .add_columns
1161                    .into_iter()
1162                    .map(|x| x.try_into())
1163                    .collect::<Result<Vec<_>>>()?;
1164                AlterKind::AddColumns { columns }
1165            }
1166            alter_request::Kind::ModifyColumnTypes(x) => {
1167                let columns = x
1168                    .modify_column_types
1169                    .into_iter()
1170                    .map(|x| x.into())
1171                    .collect::<Vec<_>>();
1172                AlterKind::ModifyColumnTypes { columns }
1173            }
1174            alter_request::Kind::DropColumns(x) => {
1175                let names = x.drop_columns.into_iter().map(|x| x.name).collect();
1176                AlterKind::DropColumns { names }
1177            }
1178            alter_request::Kind::SetTableOptions(options) => AlterKind::SetRegionOptions {
1179                options: options
1180                    .table_options
1181                    .iter()
1182                    .map(TryFrom::try_from)
1183                    .collect::<Result<Vec<_>>>()?,
1184            },
1185            alter_request::Kind::UnsetTableOptions(options) => AlterKind::UnsetRegionOptions {
1186                keys: options
1187                    .keys
1188                    .iter()
1189                    .map(|key| UnsetRegionOption::try_from(key.as_str()))
1190                    .collect::<Result<Vec<_>>>()?,
1191            },
1192            alter_request::Kind::SetIndex(o) => AlterKind::SetIndexes {
1193                options: vec![SetIndexOption::try_from(o)?],
1194            },
1195            alter_request::Kind::UnsetIndex(o) => AlterKind::UnsetIndexes {
1196                options: vec![UnsetIndexOption::try_from(o)?],
1197            },
1198            alter_request::Kind::SetIndexes(o) => AlterKind::SetIndexes {
1199                options: o
1200                    .set_indexes
1201                    .into_iter()
1202                    .map(SetIndexOption::try_from)
1203                    .collect::<Result<Vec<_>>>()?,
1204            },
1205            alter_request::Kind::UnsetIndexes(o) => AlterKind::UnsetIndexes {
1206                options: o
1207                    .unset_indexes
1208                    .into_iter()
1209                    .map(UnsetIndexOption::try_from)
1210                    .collect::<Result<Vec<_>>>()?,
1211            },
1212            alter_request::Kind::DropDefaults(x) => AlterKind::DropDefaults {
1213                names: x.drop_defaults.into_iter().map(|x| x.column_name).collect(),
1214            },
1215            alter_request::Kind::SetDefaults(x) => AlterKind::SetDefaults {
1216                columns: x
1217                    .set_defaults
1218                    .into_iter()
1219                    .map(|x| {
1220                        Ok(SetDefault {
1221                            name: x.column_name,
1222                            default_constraint: x.default_constraint.clone(),
1223                        })
1224                    })
1225                    .collect::<Result<Vec<_>>>()?,
1226            },
1227            alter_request::Kind::SyncColumns(x) => AlterKind::SyncColumns {
1228                column_metadatas: x
1229                    .column_defs
1230                    .into_iter()
1231                    .map(ColumnMetadata::try_from_column_def)
1232                    .collect::<Result<Vec<_>>>()?,
1233            },
1234        };
1235
1236        Ok(alter_kind)
1237    }
1238}
1239
1240/// Adds a column.
1241#[derive(Debug, PartialEq, Eq, Clone)]
1242pub struct AddColumn {
1243    /// Metadata of the column to add.
1244    pub column_metadata: ColumnMetadata,
1245    /// Location to add the column. If location is None, the region adds
1246    /// the column to the last.
1247    pub location: Option<AddColumnLocation>,
1248}
1249
1250impl AddColumn {
1251    /// Returns an error if the column to add is invalid.
1252    ///
1253    /// It allows adding existing columns. However, the existing column must have the same metadata
1254    /// and the location must be None.
1255    pub fn validate(&self, metadata: &RegionMetadata) -> Result<()> {
1256        ensure!(
1257            self.column_metadata.column_schema.is_nullable()
1258                || self
1259                    .column_metadata
1260                    .column_schema
1261                    .default_constraint()
1262                    .is_some(),
1263            InvalidRegionRequestSnafu {
1264                region_id: metadata.region_id,
1265                err: format!(
1266                    "no default value for column {}",
1267                    self.column_metadata.column_schema.name
1268                ),
1269            }
1270        );
1271
1272        if let Some(existing_column) =
1273            metadata.column_by_name(&self.column_metadata.column_schema.name)
1274        {
1275            // If the column already exists.
1276            ensure!(
1277                *existing_column == self.column_metadata,
1278                InvalidRegionRequestSnafu {
1279                    region_id: metadata.region_id,
1280                    err: format!(
1281                        "column {} already exists with different metadata, existing: {:?}, got: {:?}",
1282                        self.column_metadata.column_schema.name,
1283                        existing_column,
1284                        self.column_metadata,
1285                    ),
1286                }
1287            );
1288            ensure!(
1289                self.location.is_none(),
1290                InvalidRegionRequestSnafu {
1291                    region_id: metadata.region_id,
1292                    err: format!(
1293                        "column {} already exists, but location is specified",
1294                        self.column_metadata.column_schema.name
1295                    ),
1296                }
1297            );
1298        }
1299
1300        if let Some(existing_column) = metadata.column_by_id(self.column_metadata.column_id) {
1301            // Ensures the existing column has the same name.
1302            ensure!(
1303                existing_column.column_schema.name == self.column_metadata.column_schema.name,
1304                InvalidRegionRequestSnafu {
1305                    region_id: metadata.region_id,
1306                    err: format!(
1307                        "column id {} already exists with different name {}",
1308                        self.column_metadata.column_id, existing_column.column_schema.name
1309                    ),
1310                }
1311            );
1312        }
1313
1314        Ok(())
1315    }
1316
1317    /// Returns true if no column to add to the region.
1318    pub fn need_alter(&self, metadata: &RegionMetadata) -> bool {
1319        debug_assert!(self.validate(metadata).is_ok());
1320        metadata
1321            .column_by_name(&self.column_metadata.column_schema.name)
1322            .is_none()
1323    }
1324}
1325
1326impl TryFrom<v1::region::AddColumn> for AddColumn {
1327    type Error = MetadataError;
1328
1329    fn try_from(add_column: v1::region::AddColumn) -> Result<Self> {
1330        let column_def = add_column
1331            .column_def
1332            .context(InvalidRawRegionRequestSnafu {
1333                err: "missing column_def in AddColumn",
1334            })?;
1335
1336        let column_metadata = ColumnMetadata::try_from_column_def(column_def)?;
1337        let location = add_column
1338            .location
1339            .map(AddColumnLocation::try_from)
1340            .transpose()?;
1341
1342        Ok(AddColumn {
1343            column_metadata,
1344            location,
1345        })
1346    }
1347}
1348
1349/// Location to add a column.
1350#[derive(Debug, PartialEq, Eq, Clone)]
1351pub enum AddColumnLocation {
1352    /// Add the column to the first position of columns.
1353    First,
1354    /// Add the column after specific column.
1355    After {
1356        /// Add the column after this column.
1357        column_name: String,
1358    },
1359}
1360
1361impl TryFrom<v1::AddColumnLocation> for AddColumnLocation {
1362    type Error = MetadataError;
1363
1364    fn try_from(location: v1::AddColumnLocation) -> Result<Self> {
1365        let location_type = LocationType::try_from(location.location_type)
1366            .map_err(|e| InvalidRawRegionRequestSnafu { err: e.to_string() }.build())?;
1367        let add_column_location = match location_type {
1368            LocationType::First => AddColumnLocation::First,
1369            LocationType::After => AddColumnLocation::After {
1370                column_name: location.after_column_name,
1371            },
1372        };
1373
1374        Ok(add_column_location)
1375    }
1376}
1377
1378/// Change a column's datatype.
1379#[derive(Debug, PartialEq, Eq, Clone)]
1380pub struct ModifyColumnType {
1381    /// Schema of the column to modify.
1382    pub column_name: String,
1383    /// Column will be changed to this type.
1384    pub target_type: ConcreteDataType,
1385}
1386
1387impl ModifyColumnType {
1388    /// Returns an error if the column's datatype to change is invalid.
1389    pub fn validate(&self, metadata: &RegionMetadata) -> Result<()> {
1390        let column_meta = metadata
1391            .column_by_name(&self.column_name)
1392            .with_context(|| InvalidRegionRequestSnafu {
1393                region_id: metadata.region_id,
1394                err: format!("column {} not found", self.column_name),
1395            })?;
1396
1397        ensure!(
1398            matches!(column_meta.semantic_type, SemanticType::Field),
1399            InvalidRegionRequestSnafu {
1400                region_id: metadata.region_id,
1401                err: "'timestamp' or 'tag' column cannot change type".to_string()
1402            }
1403        );
1404        ensure!(
1405            column_meta
1406                .column_schema
1407                .data_type
1408                .can_arrow_type_cast_to(&self.target_type),
1409            InvalidRegionRequestSnafu {
1410                region_id: metadata.region_id,
1411                err: format!(
1412                    "column '{}' cannot be cast automatically to type '{}'",
1413                    self.column_name, self.target_type
1414                ),
1415            }
1416        );
1417
1418        Ok(())
1419    }
1420
1421    /// Returns true if no column's datatype to change to the region.
1422    pub fn need_alter(&self, metadata: &RegionMetadata) -> bool {
1423        debug_assert!(self.validate(metadata).is_ok());
1424        metadata.column_by_name(&self.column_name).is_some()
1425    }
1426}
1427
1428impl From<v1::ModifyColumnType> for ModifyColumnType {
1429    fn from(modify_column_type: v1::ModifyColumnType) -> Self {
1430        let target_type = ColumnDataTypeWrapper::new(
1431            modify_column_type.target_type(),
1432            modify_column_type.target_type_extension,
1433        )
1434        .into();
1435
1436        ModifyColumnType {
1437            column_name: modify_column_type.column_name,
1438            target_type,
1439        }
1440    }
1441}
1442
1443/// Region option changes used by ALTER requests.
1444///
1445/// This type currently derives serde for request persistence. Keep future changes
1446/// backward compatible with previously serialized variants.
1447#[derive(Debug, Eq, PartialEq, Clone, Serialize, Deserialize)]
1448pub enum SetRegionOption {
1449    WriteBufferSize(Option<ReadableSize>),
1450    Ttl(Option<TimeToLive>),
1451    // Modifying TwscOptions with values as (option name, new value).
1452    Twsc(String, String),
1453    // Modifying the SST format.
1454    Format(String),
1455    // Modifying the append mode.
1456    AppendMode(bool),
1457    // Modifying the per-region auto flush interval override.
1458    AutoFlushInterval(Option<Duration>),
1459    // Modifying the max number of rows in a parquet row group.
1460    MaxRowGroupRowCount(Option<usize>),
1461    // Stops writing new WAL entries. This operation is irreversible.
1462    SkipWal,
1463}
1464
1465impl TryFrom<&PbOption> for SetRegionOption {
1466    type Error = MetadataError;
1467
1468    fn try_from(value: &PbOption) -> std::result::Result<Self, Self::Error> {
1469        let PbOption { key, value } = value;
1470        match key.as_str() {
1471            WRITE_BUFFER_SIZE_KEY => {
1472                let size = value
1473                    .parse::<ReadableSize>()
1474                    .map_err(|_| InvalidSetRegionOptionRequestSnafu { key, value }.build())?;
1475                Ok(Self::WriteBufferSize(Some(size)))
1476            }
1477            TTL_KEY => {
1478                let ttl = TimeToLive::from_humantime_or_str(value)
1479                    .map_err(|_| InvalidSetRegionOptionRequestSnafu { key, value }.build())?;
1480
1481                Ok(Self::Ttl(Some(ttl)))
1482            }
1483            TWCS_TRIGGER_FILE_NUM | TWCS_MAX_OUTPUT_FILE_SIZE | TWCS_TIME_WINDOW => {
1484                Ok(Self::Twsc(key.clone(), value.clone()))
1485            }
1486            SST_FORMAT_KEY => Ok(Self::Format(value.clone())),
1487            APPEND_MODE_KEY => {
1488                let append_mode = value
1489                    .parse::<bool>()
1490                    .map_err(|_| InvalidSetRegionOptionRequestSnafu { key, value }.build())?;
1491                Ok(Self::AppendMode(append_mode))
1492            }
1493            AUTO_FLUSH_INTERVAL_KEY => {
1494                if value.is_empty() {
1495                    // SET 'auto_flush_interval' = NULL comes through as an empty
1496                    // string; treat it as clearing the override (fall back to
1497                    // the global default), same as Ttl.
1498                    return Ok(Self::AutoFlushInterval(None));
1499                }
1500                let interval = humantime::parse_duration(value)
1501                    .map_err(|_| InvalidSetRegionOptionRequestSnafu { key, value }.build())?;
1502                if interval <= Duration::ZERO {
1503                    return InvalidSetRegionOptionRequestSnafu { key, value }.fail();
1504                }
1505                Ok(Self::AutoFlushInterval(Some(interval)))
1506            }
1507            MAX_ROW_GROUP_ROW_COUNT => {
1508                if value.is_empty() {
1509                    return Ok(Self::MaxRowGroupRowCount(None));
1510                }
1511                let row_count = value
1512                    .parse::<usize>()
1513                    .ok()
1514                    .filter(|row_count| {
1515                        *row_count > 0 && *row_count <= MAX_ROW_GROUP_ROW_COUNT_LIMIT
1516                    })
1517                    .ok_or_else(|| InvalidSetRegionOptionRequestSnafu { key, value }.build())?;
1518                Ok(Self::MaxRowGroupRowCount(Some(row_count)))
1519            }
1520            SKIP_WAL_KEY if value == "true" => Ok(Self::SkipWal),
1521            _ => InvalidSetRegionOptionRequestSnafu { key, value }.fail(),
1522        }
1523    }
1524}
1525
1526impl From<&UnsetRegionOption> for SetRegionOption {
1527    fn from(unset_option: &UnsetRegionOption) -> Self {
1528        match unset_option {
1529            UnsetRegionOption::TwcsTriggerFileNum => {
1530                SetRegionOption::Twsc(unset_option.to_string(), String::new())
1531            }
1532            UnsetRegionOption::TwcsMaxOutputFileSize => {
1533                SetRegionOption::Twsc(unset_option.to_string(), String::new())
1534            }
1535            UnsetRegionOption::TwcsTimeWindow => {
1536                SetRegionOption::Twsc(unset_option.to_string(), String::new())
1537            }
1538            UnsetRegionOption::Ttl => SetRegionOption::Ttl(Default::default()),
1539            UnsetRegionOption::MaxRowGroupRowCount => SetRegionOption::MaxRowGroupRowCount(None),
1540            UnsetRegionOption::WriteBufferSize => SetRegionOption::WriteBufferSize(None),
1541        }
1542    }
1543}
1544
1545impl TryFrom<&str> for UnsetRegionOption {
1546    type Error = MetadataError;
1547
1548    fn try_from(key: &str) -> Result<Self> {
1549        match key.to_ascii_lowercase().as_str() {
1550            TTL_KEY => Ok(Self::Ttl),
1551            WRITE_BUFFER_SIZE_KEY => Ok(Self::WriteBufferSize),
1552            TWCS_TRIGGER_FILE_NUM => Ok(Self::TwcsTriggerFileNum),
1553            TWCS_MAX_OUTPUT_FILE_SIZE => Ok(Self::TwcsMaxOutputFileSize),
1554            TWCS_TIME_WINDOW => Ok(Self::TwcsTimeWindow),
1555            MAX_ROW_GROUP_ROW_COUNT => Ok(Self::MaxRowGroupRowCount),
1556            _ => InvalidUnsetRegionOptionRequestSnafu { key }.fail(),
1557        }
1558    }
1559}
1560
1561#[derive(Debug, Eq, PartialEq, Clone, Serialize, Deserialize)]
1562pub enum UnsetRegionOption {
1563    TwcsTriggerFileNum,
1564    TwcsMaxOutputFileSize,
1565    TwcsTimeWindow,
1566    Ttl,
1567    MaxRowGroupRowCount,
1568    WriteBufferSize,
1569}
1570
1571impl UnsetRegionOption {
1572    pub fn as_str(&self) -> &str {
1573        match self {
1574            Self::Ttl => TTL_KEY,
1575            Self::WriteBufferSize => WRITE_BUFFER_SIZE_KEY,
1576            Self::TwcsTriggerFileNum => TWCS_TRIGGER_FILE_NUM,
1577            Self::TwcsMaxOutputFileSize => TWCS_MAX_OUTPUT_FILE_SIZE,
1578            Self::TwcsTimeWindow => TWCS_TIME_WINDOW,
1579            Self::MaxRowGroupRowCount => MAX_ROW_GROUP_ROW_COUNT,
1580        }
1581    }
1582}
1583
1584impl Display for UnsetRegionOption {
1585    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1586        write!(f, "{}", self.as_str())
1587    }
1588}
1589
1590#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1591pub enum RegionFlushReason {
1592    /// Flush triggered before region migration.
1593    RegionMigration,
1594    /// Flush triggered by repartition procedure.
1595    Repartition,
1596    /// Flush triggered by remote WAL pruning.
1597    RemoteWalPrune,
1598    /// Flush region before closing region.
1599    Closing,
1600    /// Flush region before downgrading region.
1601    Downgrading,
1602}
1603
1604#[derive(Debug, Clone, Default)]
1605pub struct RegionFlushRequest {
1606    pub row_group_size: Option<usize>,
1607    pub reason: Option<RegionFlushReason>,
1608}
1609
1610#[derive(Debug)]
1611pub struct RegionCompactRequest {
1612    pub options: compact_request::Options,
1613    pub parallelism: Option<u32>,
1614    pub time_range: Option<TimestampRange>,
1615}
1616
1617impl Default for RegionCompactRequest {
1618    fn default() -> Self {
1619        Self {
1620            // Default to regular compaction.
1621            options: compact_request::Options::Regular(Default::default()),
1622            parallelism: None,
1623            time_range: None,
1624        }
1625    }
1626}
1627
1628#[derive(Debug, Clone, Default)]
1629pub struct RegionBuildIndexRequest {}
1630
1631/// Truncate region request.
1632#[derive(Debug)]
1633pub enum RegionTruncateRequest {
1634    /// Truncate all data in the region.
1635    All,
1636    /// Discard all unflushed data while preserving persisted SST files.
1637    ///
1638    /// This destroys the region's in-memory data irreversibly. Persisted SST files and
1639    /// the writable state are preserved, and the WAL is obsoleted up to the discard point
1640    /// so a restart won't replay the discarded data.
1641    ///
1642    /// An error may be returned after the data has already been discarded, because the
1643    /// WAL is obsoleted last. Retrying is safe: a request against a region with nothing
1644    /// left to discard only re-attempts the WAL obsoletion.
1645    Unflushed,
1646    ByTimeRanges {
1647        /// Time ranges to truncate. Both bound are inclusive.
1648        /// only files that are fully contained in the time range will be truncated.
1649        /// so no guarantee that all data in the time range will be truncated.
1650        time_ranges: Vec<(Timestamp, Timestamp)>,
1651    },
1652}
1653
1654/// Catchup region request.
1655///
1656/// Makes a readonly region to catch up to leader region changes.
1657/// There is no effect if it operating on a leader region.
1658#[derive(Debug, Clone, Copy, Default)]
1659pub struct RegionCatchupRequest {
1660    /// Sets it to writable if it's available after it has caught up with all changes.
1661    pub set_writable: bool,
1662    /// The `entry_id` that was expected to reply to.
1663    /// `None` stands replaying to latest.
1664    pub entry_id: Option<entry::Id>,
1665    /// Used for metrics metadata region.
1666    /// The `entry_id` that was expected to reply to.
1667    /// `None` stands replaying to latest.
1668    pub metadata_entry_id: Option<entry::Id>,
1669    /// The hint for replaying memtable.
1670    pub location_id: Option<u64>,
1671    /// Replay checkpoint.
1672    pub checkpoint: Option<ReplayCheckpoint>,
1673}
1674
1675#[derive(Debug, Clone)]
1676pub struct RegionBulkInsertsRequest {
1677    pub region_id: RegionId,
1678    pub payload: DfRecordBatch,
1679    pub raw_data: ArrowIpc,
1680    pub partition_expr_version: Option<u64>,
1681    pub aligned_schema_version: Option<u64>,
1682}
1683
1684impl RegionBulkInsertsRequest {
1685    pub fn estimated_size(&self) -> usize {
1686        self.payload.get_array_memory_size()
1687    }
1688}
1689
1690/// Request to stage a region with a new partition directive.
1691///
1692/// This request transitions a region into the staging mode.
1693/// It first flushes the memtable for the old partition expression if it is not
1694/// empty, then enters the staging mode with the new directive.
1695#[derive(Debug, Clone, PartialEq, Eq)]
1696pub enum StagingPartitionDirective {
1697    UpdatePartitionExpr(String),
1698    RejectAllWrites,
1699}
1700
1701impl StagingPartitionDirective {
1702    /// Returns the partition expression carried by this directive, if any.
1703    pub fn partition_expr(&self) -> Option<&str> {
1704        match self {
1705            Self::UpdatePartitionExpr(expr) => Some(expr),
1706            Self::RejectAllWrites => None,
1707        }
1708    }
1709}
1710
1711#[derive(Debug, Clone)]
1712pub struct EnterStagingRequest {
1713    /// The staging partition directive of the region.
1714    pub partition_directive: StagingPartitionDirective,
1715}
1716
1717impl EnterStagingRequest {
1718    /// Builds an enter-staging request with a partition expression directive.
1719    pub fn with_partition_expr(partition_expr: String) -> Self {
1720        Self {
1721            partition_directive: StagingPartitionDirective::UpdatePartitionExpr(partition_expr),
1722        }
1723    }
1724}
1725
1726/// This request is used as part of the region repartition.
1727///
1728/// After a region has entered staging mode with a new partition expression
1729/// expression) and a separate process (for example, `remap_manifests`) has
1730/// generated the new file assignments for the staging region, this request
1731/// applies that generated manifest to the region.
1732///
1733/// In practice, this means:
1734/// - The `partition_expr` identifies the staging partition expression that the manifest
1735///   was generated for.
1736/// - `central_region_id` specifies which region holds the staging blob storage
1737///   where the manifest was written during the `remap_manifests` operation.
1738/// - `manifest_path` is the relative path within the central region's staging
1739///   blob storage to fetch the generated manifest.
1740///
1741/// It should typically be called **after** the staging region has been
1742/// initialized by [`EnterStagingRequest`] and the new file layout has been
1743/// computed, to finalize the repartition operation.
1744#[derive(Debug, Clone)]
1745pub struct ApplyStagingManifestRequest {
1746    /// The partition expression of the staging region.
1747    pub partition_expr: String,
1748    /// The region that stores the staging manifests in its staging blob storage.
1749    pub central_region_id: RegionId,
1750    /// The relative path to the staging manifest within the central region's
1751    /// staging blob storage.
1752    pub manifest_path: String,
1753}
1754
1755impl fmt::Display for RegionRequest {
1756    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1757        match self {
1758            RegionRequest::Put(_) => write!(f, "Put"),
1759            RegionRequest::Delete(_) => write!(f, "Delete"),
1760            RegionRequest::Create(_) => write!(f, "Create"),
1761            RegionRequest::Drop(_) => write!(f, "Drop"),
1762            RegionRequest::Open(_) => write!(f, "Open"),
1763            RegionRequest::CleanUp(_) => write!(f, "CleanUp"),
1764            RegionRequest::Close(_) => write!(f, "Close"),
1765            RegionRequest::Alter(_) => write!(f, "Alter"),
1766            RegionRequest::Flush(_) => write!(f, "Flush"),
1767            RegionRequest::Compact(_) => write!(f, "Compact"),
1768            RegionRequest::BuildIndex(_) => write!(f, "BuildIndex"),
1769            RegionRequest::Truncate(_) => write!(f, "Truncate"),
1770            RegionRequest::Catchup(_) => write!(f, "Catchup"),
1771            RegionRequest::BulkInserts(_) => write!(f, "BulkInserts"),
1772            RegionRequest::EnterStaging(_) => write!(f, "EnterStaging"),
1773            RegionRequest::ApplyStagingManifest(_) => write!(f, "ApplyStagingManifest"),
1774        }
1775    }
1776}
1777
1778#[cfg(test)]
1779mod tests {
1780
1781    use api::v1::region::RegionColumnDef;
1782    use api::v1::{ColumnDataType, ColumnDef};
1783    use common_time::range::TimestampRange;
1784    use datatypes::prelude::ConcreteDataType;
1785    use datatypes::schema::{ColumnSchema, FulltextAnalyzer, FulltextBackend};
1786
1787    use super::*;
1788    use crate::metadata::RegionMetadataBuilder;
1789
1790    #[test]
1791    fn test_make_region_compact_with_time_range() {
1792        let requests = make_region_compact(CompactRequest {
1793            region_id: 42,
1794            time_range: Some(api::v1::region::CompactionTimeRange {
1795                start: 1_000,
1796                end: 2_000,
1797                time_unit: api::v1::TimeUnit::Microsecond as i32,
1798            }),
1799            ..Default::default()
1800        })
1801        .unwrap();
1802
1803        let RegionRequest::Compact(request) = &requests[0].1 else {
1804            unreachable!();
1805        };
1806        assert_eq!(
1807            Some(
1808                TimestampRange::new(
1809                    Timestamp::new_microsecond(1_000),
1810                    Timestamp::new_microsecond(2_000),
1811                )
1812                .unwrap()
1813            ),
1814            request.time_range
1815        );
1816    }
1817
1818    #[test]
1819    fn test_make_region_truncate_unflushed() {
1820        let region_id = RegionId::new(42, 3);
1821        let requests =
1822            RegionRequest::try_from_request_body(region_request::Body::Truncate(TruncateRequest {
1823                region_id: region_id.as_u64(),
1824                kind: Some(truncate_request::Kind::Unflushed(
1825                    api::v1::region::Unflushed {},
1826                )),
1827            }))
1828            .unwrap();
1829
1830        assert_eq!(region_id, requests[0].0);
1831        assert!(matches!(
1832            requests[0].1,
1833            RegionRequest::Truncate(RegionTruncateRequest::Unflushed)
1834        ));
1835    }
1836
1837    #[test]
1838    fn test_make_region_truncate_requires_kind() {
1839        let error =
1840            RegionRequest::try_from_request_body(region_request::Body::Truncate(TruncateRequest {
1841                region_id: RegionId::new(42, 3).as_u64(),
1842                kind: None,
1843            }))
1844            .unwrap_err();
1845
1846        assert!(
1847            error
1848                .to_string()
1849                .contains("missing kind in TruncateRequest")
1850        );
1851    }
1852
1853    #[test]
1854    fn test_from_proto_location() {
1855        let proto_location = v1::AddColumnLocation {
1856            location_type: LocationType::First as i32,
1857            after_column_name: String::default(),
1858        };
1859        let location = AddColumnLocation::try_from(proto_location).unwrap();
1860        assert_eq!(location, AddColumnLocation::First);
1861
1862        let proto_location = v1::AddColumnLocation {
1863            location_type: 10,
1864            after_column_name: String::default(),
1865        };
1866        AddColumnLocation::try_from(proto_location).unwrap_err();
1867
1868        let proto_location = v1::AddColumnLocation {
1869            location_type: LocationType::After as i32,
1870            after_column_name: "a".to_string(),
1871        };
1872        let location = AddColumnLocation::try_from(proto_location).unwrap();
1873        assert_eq!(
1874            location,
1875            AddColumnLocation::After {
1876                column_name: "a".to_string()
1877            }
1878        );
1879    }
1880
1881    #[test]
1882    fn test_from_none_proto_add_column() {
1883        AddColumn::try_from(v1::region::AddColumn {
1884            column_def: None,
1885            location: None,
1886        })
1887        .unwrap_err();
1888    }
1889
1890    #[test]
1891    fn test_set_region_option_auto_flush_interval_try_from() {
1892        use std::time::Duration;
1893
1894        // Valid duration
1895        let pb = PbOption {
1896            key: "auto_flush_interval".to_string(),
1897            value: "5m".to_string(),
1898        };
1899        let opt = SetRegionOption::try_from(&pb).unwrap();
1900        assert_eq!(
1901            opt,
1902            SetRegionOption::AutoFlushInterval(Some(Duration::from_secs(300)))
1903        );
1904
1905        // Empty value clears the override (mirrors Ttl's behaviour for `SET key = NULL`).
1906        let pb = PbOption {
1907            key: "auto_flush_interval".to_string(),
1908            value: String::new(),
1909        };
1910        let opt = SetRegionOption::try_from(&pb).unwrap();
1911        assert_eq!(opt, SetRegionOption::AutoFlushInterval(None));
1912
1913        // Zero is rejected up front (engine invariant).
1914        let pb = PbOption {
1915            key: "auto_flush_interval".to_string(),
1916            value: "0s".to_string(),
1917        };
1918        assert!(SetRegionOption::try_from(&pb).is_err());
1919
1920        // Garbage value is rejected.
1921        let pb = PbOption {
1922            key: "auto_flush_interval".to_string(),
1923            value: "not_a_duration".to_string(),
1924        };
1925        assert!(SetRegionOption::try_from(&pb).is_err());
1926    }
1927
1928    #[test]
1929    fn test_set_region_option_skip_wal_try_from() {
1930        let pb = PbOption {
1931            key: SKIP_WAL_KEY.to_string(),
1932            value: "true".to_string(),
1933        };
1934        assert_eq!(
1935            SetRegionOption::SkipWal,
1936            SetRegionOption::try_from(&pb).unwrap()
1937        );
1938
1939        for value in ["false", "", "invalid"] {
1940            let pb = PbOption {
1941                key: SKIP_WAL_KEY.to_string(),
1942                value: value.to_string(),
1943            };
1944            assert!(SetRegionOption::try_from(&pb).is_err());
1945        }
1946
1947        assert!(UnsetRegionOption::try_from(SKIP_WAL_KEY).is_err());
1948    }
1949
1950    #[test]
1951    fn test_set_region_option_max_row_group_row_count_try_from() {
1952        let pb = PbOption {
1953            key: MAX_ROW_GROUP_ROW_COUNT.to_string(),
1954            value: "512".to_string(),
1955        };
1956        assert_eq!(
1957            SetRegionOption::MaxRowGroupRowCount(Some(512)),
1958            SetRegionOption::try_from(&pb).unwrap()
1959        );
1960
1961        let pb = PbOption {
1962            key: MAX_ROW_GROUP_ROW_COUNT.to_string(),
1963            value: String::new(),
1964        };
1965        assert_eq!(
1966            SetRegionOption::MaxRowGroupRowCount(None),
1967            SetRegionOption::try_from(&pb).unwrap()
1968        );
1969
1970        for value in [
1971            "0".to_string(),
1972            (MAX_ROW_GROUP_ROW_COUNT_LIMIT + 1).to_string(),
1973            "invalid".to_string(),
1974        ] {
1975            let pb = PbOption {
1976                key: MAX_ROW_GROUP_ROW_COUNT.to_string(),
1977                value,
1978            };
1979            assert!(SetRegionOption::try_from(&pb).is_err());
1980        }
1981
1982        assert_eq!(
1983            UnsetRegionOption::MaxRowGroupRowCount,
1984            UnsetRegionOption::try_from(MAX_ROW_GROUP_ROW_COUNT).unwrap()
1985        );
1986    }
1987
1988    #[test]
1989    fn test_from_proto_alter_request() {
1990        RegionAlterRequest::try_from(AlterRequest {
1991            region_id: 0,
1992            schema_version: 1,
1993            kind: None,
1994        })
1995        .unwrap_err();
1996
1997        let request = RegionAlterRequest::try_from(AlterRequest {
1998            region_id: 0,
1999            schema_version: 1,
2000            kind: Some(alter_request::Kind::AddColumns(v1::region::AddColumns {
2001                add_columns: vec![v1::region::AddColumn {
2002                    column_def: Some(RegionColumnDef {
2003                        column_def: Some(ColumnDef {
2004                            name: "a".to_string(),
2005                            data_type: ColumnDataType::String as i32,
2006                            is_nullable: true,
2007                            default_constraint: vec![],
2008                            semantic_type: SemanticType::Field as i32,
2009                            comment: String::new(),
2010                            ..Default::default()
2011                        }),
2012                        column_id: 1,
2013                    }),
2014                    location: Some(v1::AddColumnLocation {
2015                        location_type: LocationType::First as i32,
2016                        after_column_name: String::default(),
2017                    }),
2018                }],
2019            })),
2020        })
2021        .unwrap();
2022
2023        assert_eq!(
2024            request,
2025            RegionAlterRequest {
2026                kind: AlterKind::AddColumns {
2027                    columns: vec![AddColumn {
2028                        column_metadata: ColumnMetadata {
2029                            column_schema: ColumnSchema::new(
2030                                "a",
2031                                ConcreteDataType::string_datatype(),
2032                                true,
2033                            ),
2034                            semantic_type: SemanticType::Field,
2035                            column_id: 1,
2036                        },
2037                        location: Some(AddColumnLocation::First),
2038                    }]
2039                },
2040            }
2041        );
2042    }
2043
2044    #[test]
2045    fn test_write_buffer_size_region_options() {
2046        let option = PbOption {
2047            key: WRITE_BUFFER_SIZE_KEY.to_string(),
2048            value: "128MiB".to_string(),
2049        };
2050        assert_eq!(
2051            SetRegionOption::WriteBufferSize(Some(ReadableSize::mb(128))),
2052            SetRegionOption::try_from(&option).unwrap()
2053        );
2054
2055        let option = PbOption {
2056            key: WRITE_BUFFER_SIZE_KEY.to_string(),
2057            value: "invalid".to_string(),
2058        };
2059        SetRegionOption::try_from(&option).unwrap_err();
2060
2061        // Clearing the option uses UNSET. Empty SET values, including SQL NULL,
2062        // remain invalid because the SQL layer does not distinguish NULL from ''.
2063        let option = PbOption {
2064            key: WRITE_BUFFER_SIZE_KEY.to_string(),
2065            value: String::new(),
2066        };
2067        SetRegionOption::try_from(&option).unwrap_err();
2068
2069        assert_eq!(
2070            UnsetRegionOption::WriteBufferSize,
2071            UnsetRegionOption::try_from(WRITE_BUFFER_SIZE_KEY).unwrap()
2072        );
2073        assert_eq!(
2074            SetRegionOption::WriteBufferSize(None),
2075            SetRegionOption::from(&UnsetRegionOption::WriteBufferSize)
2076        );
2077    }
2078
2079    /// Returns a new region metadata for testing. Metadata:
2080    /// `[(ts, ms, 1), (tag_0, string, 2), (field_0, string, 3), (field_1, bool, 4)]`
2081    fn new_metadata() -> RegionMetadata {
2082        let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 1));
2083        builder
2084            .push_column_metadata(ColumnMetadata {
2085                column_schema: ColumnSchema::new(
2086                    "ts",
2087                    ConcreteDataType::timestamp_millisecond_datatype(),
2088                    false,
2089                ),
2090                semantic_type: SemanticType::Timestamp,
2091                column_id: 1,
2092            })
2093            .push_column_metadata(ColumnMetadata {
2094                column_schema: ColumnSchema::new(
2095                    "tag_0",
2096                    ConcreteDataType::string_datatype(),
2097                    true,
2098                ),
2099                semantic_type: SemanticType::Tag,
2100                column_id: 2,
2101            })
2102            .push_column_metadata(ColumnMetadata {
2103                column_schema: ColumnSchema::new(
2104                    "field_0",
2105                    ConcreteDataType::string_datatype(),
2106                    true,
2107                ),
2108                semantic_type: SemanticType::Field,
2109                column_id: 3,
2110            })
2111            .push_column_metadata(ColumnMetadata {
2112                column_schema: ColumnSchema::new(
2113                    "field_1",
2114                    ConcreteDataType::boolean_datatype(),
2115                    true,
2116                ),
2117                semantic_type: SemanticType::Field,
2118                column_id: 4,
2119            })
2120            .primary_key(vec![2]);
2121        builder.build().unwrap()
2122    }
2123
2124    #[test]
2125    fn test_add_column_validate() {
2126        let metadata = new_metadata();
2127        let add_column = AddColumn {
2128            column_metadata: ColumnMetadata {
2129                column_schema: ColumnSchema::new(
2130                    "tag_1",
2131                    ConcreteDataType::string_datatype(),
2132                    true,
2133                ),
2134                semantic_type: SemanticType::Tag,
2135                column_id: 5,
2136            },
2137            location: None,
2138        };
2139        add_column.validate(&metadata).unwrap();
2140        assert!(add_column.need_alter(&metadata));
2141
2142        // Add not null column.
2143        AddColumn {
2144            column_metadata: ColumnMetadata {
2145                column_schema: ColumnSchema::new(
2146                    "tag_1",
2147                    ConcreteDataType::string_datatype(),
2148                    false,
2149                ),
2150                semantic_type: SemanticType::Tag,
2151                column_id: 5,
2152            },
2153            location: None,
2154        }
2155        .validate(&metadata)
2156        .unwrap_err();
2157
2158        // Add existing column.
2159        let add_column = AddColumn {
2160            column_metadata: ColumnMetadata {
2161                column_schema: ColumnSchema::new(
2162                    "tag_0",
2163                    ConcreteDataType::string_datatype(),
2164                    true,
2165                ),
2166                semantic_type: SemanticType::Tag,
2167                column_id: 2,
2168            },
2169            location: None,
2170        };
2171        add_column.validate(&metadata).unwrap();
2172        assert!(!add_column.need_alter(&metadata));
2173    }
2174
2175    #[test]
2176    fn test_add_duplicate_columns() {
2177        let kind = AlterKind::AddColumns {
2178            columns: vec![
2179                AddColumn {
2180                    column_metadata: ColumnMetadata {
2181                        column_schema: ColumnSchema::new(
2182                            "tag_1",
2183                            ConcreteDataType::string_datatype(),
2184                            true,
2185                        ),
2186                        semantic_type: SemanticType::Tag,
2187                        column_id: 5,
2188                    },
2189                    location: None,
2190                },
2191                AddColumn {
2192                    column_metadata: ColumnMetadata {
2193                        column_schema: ColumnSchema::new(
2194                            "tag_1",
2195                            ConcreteDataType::string_datatype(),
2196                            true,
2197                        ),
2198                        semantic_type: SemanticType::Field,
2199                        column_id: 6,
2200                    },
2201                    location: None,
2202                },
2203            ],
2204        };
2205        let metadata = new_metadata();
2206        kind.validate(&metadata).unwrap();
2207        assert!(kind.need_alter(&metadata));
2208    }
2209
2210    #[test]
2211    fn test_add_existing_column_different_metadata() {
2212        let metadata = new_metadata();
2213
2214        // Add existing column with different id.
2215        let kind = AlterKind::AddColumns {
2216            columns: vec![AddColumn {
2217                column_metadata: ColumnMetadata {
2218                    column_schema: ColumnSchema::new(
2219                        "tag_0",
2220                        ConcreteDataType::string_datatype(),
2221                        true,
2222                    ),
2223                    semantic_type: SemanticType::Tag,
2224                    column_id: 4,
2225                },
2226                location: None,
2227            }],
2228        };
2229        kind.validate(&metadata).unwrap_err();
2230
2231        // Add existing column with different type.
2232        let kind = AlterKind::AddColumns {
2233            columns: vec![AddColumn {
2234                column_metadata: ColumnMetadata {
2235                    column_schema: ColumnSchema::new(
2236                        "tag_0",
2237                        ConcreteDataType::int64_datatype(),
2238                        true,
2239                    ),
2240                    semantic_type: SemanticType::Tag,
2241                    column_id: 2,
2242                },
2243                location: None,
2244            }],
2245        };
2246        kind.validate(&metadata).unwrap_err();
2247
2248        // Add existing column with different name.
2249        let kind = AlterKind::AddColumns {
2250            columns: vec![AddColumn {
2251                column_metadata: ColumnMetadata {
2252                    column_schema: ColumnSchema::new(
2253                        "tag_1",
2254                        ConcreteDataType::string_datatype(),
2255                        true,
2256                    ),
2257                    semantic_type: SemanticType::Tag,
2258                    column_id: 2,
2259                },
2260                location: None,
2261            }],
2262        };
2263        kind.validate(&metadata).unwrap_err();
2264    }
2265
2266    #[test]
2267    fn test_add_existing_column_with_location() {
2268        let metadata = new_metadata();
2269        let kind = AlterKind::AddColumns {
2270            columns: vec![AddColumn {
2271                column_metadata: ColumnMetadata {
2272                    column_schema: ColumnSchema::new(
2273                        "tag_0",
2274                        ConcreteDataType::string_datatype(),
2275                        true,
2276                    ),
2277                    semantic_type: SemanticType::Tag,
2278                    column_id: 2,
2279                },
2280                location: Some(AddColumnLocation::First),
2281            }],
2282        };
2283        kind.validate(&metadata).unwrap_err();
2284    }
2285
2286    #[test]
2287    fn test_validate_drop_column() {
2288        let metadata = new_metadata();
2289        let kind = AlterKind::DropColumns {
2290            names: vec!["xxxx".to_string()],
2291        };
2292        kind.validate(&metadata).unwrap();
2293        assert!(!kind.need_alter(&metadata));
2294
2295        AlterKind::DropColumns {
2296            names: vec!["tag_0".to_string()],
2297        }
2298        .validate(&metadata)
2299        .unwrap_err();
2300
2301        let kind = AlterKind::DropColumns {
2302            names: vec!["field_0".to_string()],
2303        };
2304        kind.validate(&metadata).unwrap();
2305        assert!(kind.need_alter(&metadata));
2306    }
2307
2308    #[test]
2309    fn test_validate_modify_column_type() {
2310        let metadata = new_metadata();
2311        AlterKind::ModifyColumnTypes {
2312            columns: vec![ModifyColumnType {
2313                column_name: "xxxx".to_string(),
2314                target_type: ConcreteDataType::string_datatype(),
2315            }],
2316        }
2317        .validate(&metadata)
2318        .unwrap_err();
2319
2320        AlterKind::ModifyColumnTypes {
2321            columns: vec![ModifyColumnType {
2322                column_name: "field_1".to_string(),
2323                target_type: ConcreteDataType::date_datatype(),
2324            }],
2325        }
2326        .validate(&metadata)
2327        .unwrap_err();
2328
2329        AlterKind::ModifyColumnTypes {
2330            columns: vec![ModifyColumnType {
2331                column_name: "ts".to_string(),
2332                target_type: ConcreteDataType::date_datatype(),
2333            }],
2334        }
2335        .validate(&metadata)
2336        .unwrap_err();
2337
2338        AlterKind::ModifyColumnTypes {
2339            columns: vec![ModifyColumnType {
2340                column_name: "tag_0".to_string(),
2341                target_type: ConcreteDataType::date_datatype(),
2342            }],
2343        }
2344        .validate(&metadata)
2345        .unwrap_err();
2346
2347        let kind = AlterKind::ModifyColumnTypes {
2348            columns: vec![ModifyColumnType {
2349                column_name: "field_0".to_string(),
2350                target_type: ConcreteDataType::int32_datatype(),
2351            }],
2352        };
2353        kind.validate(&metadata).unwrap();
2354        assert!(kind.need_alter(&metadata));
2355    }
2356
2357    #[test]
2358    fn test_validate_add_columns() {
2359        let kind = AlterKind::AddColumns {
2360            columns: vec![
2361                AddColumn {
2362                    column_metadata: ColumnMetadata {
2363                        column_schema: ColumnSchema::new(
2364                            "tag_1",
2365                            ConcreteDataType::string_datatype(),
2366                            true,
2367                        ),
2368                        semantic_type: SemanticType::Tag,
2369                        column_id: 5,
2370                    },
2371                    location: None,
2372                },
2373                AddColumn {
2374                    column_metadata: ColumnMetadata {
2375                        column_schema: ColumnSchema::new(
2376                            "field_2",
2377                            ConcreteDataType::string_datatype(),
2378                            true,
2379                        ),
2380                        semantic_type: SemanticType::Field,
2381                        column_id: 6,
2382                    },
2383                    location: None,
2384                },
2385            ],
2386        };
2387        let request = RegionAlterRequest { kind };
2388        let mut metadata = new_metadata();
2389        metadata.schema_version = 1;
2390        request.validate(&metadata).unwrap();
2391    }
2392
2393    #[test]
2394    fn test_validate_create_region() {
2395        let column_metadatas = vec![
2396            ColumnMetadata {
2397                column_schema: ColumnSchema::new(
2398                    "ts",
2399                    ConcreteDataType::timestamp_millisecond_datatype(),
2400                    false,
2401                ),
2402                semantic_type: SemanticType::Timestamp,
2403                column_id: 1,
2404            },
2405            ColumnMetadata {
2406                column_schema: ColumnSchema::new(
2407                    "tag_0",
2408                    ConcreteDataType::string_datatype(),
2409                    true,
2410                ),
2411                semantic_type: SemanticType::Tag,
2412                column_id: 2,
2413            },
2414            ColumnMetadata {
2415                column_schema: ColumnSchema::new(
2416                    "field_0",
2417                    ConcreteDataType::string_datatype(),
2418                    true,
2419                ),
2420                semantic_type: SemanticType::Field,
2421                column_id: 3,
2422            },
2423        ];
2424        let create = RegionCreateRequest {
2425            engine: "mito".to_string(),
2426            column_metadatas,
2427            primary_key: vec![3, 4],
2428            options: HashMap::new(),
2429            table_dir: "path".to_string(),
2430            path_type: PathType::Bare,
2431            partition_expr_json: Some("".to_string()),
2432            requirements: Default::default(),
2433        };
2434
2435        assert!(create.validate().is_err());
2436    }
2437
2438    #[test]
2439    fn test_parse_create_region_requirements_defaults_to_empty() {
2440        let create = CreateRequest {
2441            region_id: RegionId::new(42, 0).as_u64(),
2442            engine: "mito".to_string(),
2443            column_defs: vec![],
2444            primary_key: vec![],
2445            path: "test".to_string(),
2446            options: HashMap::new(),
2447            partition: None,
2448            requirements: None,
2449        };
2450
2451        let requests =
2452            RegionRequest::try_from_request_body(region_request::Body::Create(create)).unwrap();
2453        let RegionRequest::Create(request) = &requests[0].1 else {
2454            unreachable!()
2455        };
2456
2457        assert_eq!(request.requirements, RegionRequirements::empty());
2458    }
2459
2460    #[test]
2461    fn test_parse_create_region_requirements_from_proto() {
2462        let create = CreateRequest {
2463            region_id: RegionId::new(42, 0).as_u64(),
2464            engine: "mito".to_string(),
2465            column_defs: vec![],
2466            primary_key: vec![],
2467            path: "test".to_string(),
2468            options: HashMap::new(),
2469            partition: None,
2470            requirements: Some(api::v1::region::RegionRequirements {
2471                object_storage: true,
2472            }),
2473        };
2474
2475        let requests =
2476            RegionRequest::try_from_request_body(region_request::Body::Create(create)).unwrap();
2477        let RegionRequest::Create(request) = &requests[0].1 else {
2478            unreachable!()
2479        };
2480
2481        assert_eq!(request.requirements, RegionRequirements::object_storage());
2482    }
2483
2484    #[test]
2485    fn test_parse_region_cleanup_from_proto() {
2486        let clean_up = api::v1::region::CleanUpRequest {
2487            region_id: RegionId::new(42, 3).as_u64(),
2488            engine: "mito".to_string(),
2489            path: "test".to_string(),
2490            options: HashMap::from([("k".to_string(), "v".to_string())]),
2491        };
2492
2493        let requests =
2494            RegionRequest::try_from_request_body(region_request::Body::CleanUp(clean_up)).unwrap();
2495        let RegionRequest::CleanUp(request) = &requests[0].1 else {
2496            unreachable!()
2497        };
2498
2499        assert_eq!(requests[0].0, RegionId::new(42, 3));
2500        assert_eq!(request.engine, "mito");
2501        assert_eq!(request.table_dir, "data/test/42/");
2502        assert_eq!(request.path_type, PathType::Bare);
2503        assert_eq!(request.options.get("k"), Some(&"v".to_string()));
2504    }
2505
2506    #[test]
2507    fn test_validate_modify_column_fulltext_options() {
2508        let kind = AlterKind::SetIndexes {
2509            options: vec![SetIndexOption::Fulltext {
2510                column_name: "tag_0".to_string(),
2511                options: FulltextOptions::new_unchecked(
2512                    true,
2513                    FulltextAnalyzer::Chinese,
2514                    false,
2515                    FulltextBackend::Bloom,
2516                    1000,
2517                    0.01,
2518                ),
2519            }],
2520        };
2521        let request = RegionAlterRequest { kind };
2522        let mut metadata = new_metadata();
2523        metadata.schema_version = 1;
2524        request.validate(&metadata).unwrap();
2525
2526        let kind = AlterKind::UnsetIndexes {
2527            options: vec![UnsetIndexOption::Fulltext {
2528                column_name: "tag_0".to_string(),
2529            }],
2530        };
2531        let request = RegionAlterRequest { kind };
2532        let mut metadata = new_metadata();
2533        metadata.schema_version = 1;
2534        request.validate(&metadata).unwrap();
2535    }
2536
2537    #[test]
2538    fn test_validate_sync_columns() {
2539        let metadata = new_metadata();
2540        let kind = AlterKind::SyncColumns {
2541            column_metadatas: vec![
2542                ColumnMetadata {
2543                    column_schema: ColumnSchema::new(
2544                        "tag_1",
2545                        ConcreteDataType::string_datatype(),
2546                        true,
2547                    ),
2548                    semantic_type: SemanticType::Tag,
2549                    column_id: 5,
2550                },
2551                ColumnMetadata {
2552                    column_schema: ColumnSchema::new(
2553                        "field_2",
2554                        ConcreteDataType::string_datatype(),
2555                        true,
2556                    ),
2557                    semantic_type: SemanticType::Field,
2558                    column_id: 6,
2559                },
2560            ],
2561        };
2562        let err = kind.validate(&metadata).unwrap_err();
2563        assert!(err.to_string().contains("not a primary key"));
2564
2565        // Change the timestamp column name.
2566        let mut column_metadatas_with_different_ts_column = metadata.column_metadatas.clone();
2567        let ts_column = column_metadatas_with_different_ts_column
2568            .iter_mut()
2569            .find(|c| c.semantic_type == SemanticType::Timestamp)
2570            .unwrap();
2571        ts_column.column_schema.name = "ts1".to_string();
2572
2573        let kind = AlterKind::SyncColumns {
2574            column_metadatas: column_metadatas_with_different_ts_column,
2575        };
2576        let err = kind.validate(&metadata).unwrap_err();
2577        assert!(
2578            err.to_string()
2579                .contains("timestamp column ts has different id")
2580        );
2581
2582        // Change the primary key column name.
2583        let mut column_metadatas_with_different_pk_column = metadata.column_metadatas.clone();
2584        let pk_column = column_metadatas_with_different_pk_column
2585            .iter_mut()
2586            .find(|c| c.column_schema.name == "tag_0")
2587            .unwrap();
2588        pk_column.column_id = 100;
2589        let kind = AlterKind::SyncColumns {
2590            column_metadatas: column_metadatas_with_different_pk_column,
2591        };
2592        let err = kind.validate(&metadata).unwrap_err();
2593        assert!(
2594            err.to_string()
2595                .contains("column with same name tag_0 has different id")
2596        );
2597
2598        // Add a new field column.
2599        let mut column_metadatas_with_new_field_column = metadata.column_metadatas.clone();
2600        column_metadatas_with_new_field_column.push(ColumnMetadata {
2601            column_schema: ColumnSchema::new("field_2", ConcreteDataType::string_datatype(), true),
2602            semantic_type: SemanticType::Field,
2603            column_id: 4,
2604        });
2605        let kind = AlterKind::SyncColumns {
2606            column_metadatas: column_metadatas_with_new_field_column,
2607        };
2608        kind.validate(&metadata).unwrap();
2609    }
2610
2611    #[test]
2612    fn test_cast_path_type_to_primitive() {
2613        assert_eq!(PathType::Bare as u8, 0);
2614        assert_eq!(PathType::Data as u8, 1);
2615        assert_eq!(PathType::Metadata as u8, 2);
2616        assert_eq!(
2617            PathType::try_from(PathType::Bare as u8).unwrap(),
2618            PathType::Bare
2619        );
2620        assert_eq!(
2621            PathType::try_from(PathType::Data as u8).unwrap(),
2622            PathType::Data
2623        );
2624        assert_eq!(
2625            PathType::try_from(PathType::Metadata as u8).unwrap(),
2626            PathType::Metadata
2627        );
2628    }
2629}