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