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