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