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