1use std::collections::HashMap;
18use std::sync::Arc;
19use std::time::Instant;
20
21use api::helper::{
22 ColumnDataTypeWrapper, is_column_type_value_eq, is_semantic_type_eq, proto_value_type,
23};
24use api::v1::column_def::options_from_column_schema;
25use api::v1::{ColumnDataType, ColumnSchema, OpType, Rows, SemanticType, Value, WriteHint};
26use common_telemetry::info;
27use datatypes::prelude::DataType;
28use partition::expr::PartitionExpr;
29use prometheus::HistogramTimer;
30use prost::Message;
31use smallvec::SmallVec;
32use snafu::{OptionExt, ResultExt, ensure};
33use store_api::ManifestVersion;
34use store_api::codec::{PrimaryKeyEncoding, infer_primary_key_encoding_from_hint};
35use store_api::metadata::{ColumnMetadata, RegionMetadata, RegionMetadataRef};
36use store_api::region_engine::{
37 MitoCopyRegionFromResponse, SetRegionRoleStateResponse, SettableRegionRoleState,
38};
39use store_api::region_request::{
40 AffectedRows, ApplyStagingManifestRequest, EnterStagingRequest, RegionAlterRequest,
41 RegionBuildIndexRequest, RegionBulkInsertsRequest, RegionCatchupRequest, RegionCleanUpRequest,
42 RegionCloseRequest, RegionCompactRequest, RegionCreateRequest, RegionDropRequest,
43 RegionFlushRequest, RegionOpenRequest, RegionRequest, RegionTruncateRequest,
44 StagingPartitionDirective,
45};
46use store_api::storage::{FileId, RegionId, SequenceNumber};
47use tokio::sync::oneshot::{self, Receiver, Sender};
48
49use crate::compaction::{CompactionExecution, CompactionPickFinished};
50use crate::error::{
51 CompactRegionSnafu, CompactionCancelledSnafu, ConvertColumnDataTypeSnafu, CreateDefaultSnafu,
52 Error, FillDefaultSnafu, FlushRegionSnafu, InvalidPartitionExprSnafu, InvalidRequestSnafu,
53 MissingPartitionExprSnafu, Result, UnexpectedSnafu,
54};
55use crate::flush::FlushReason;
56use crate::manifest::action::{RegionEdit, TruncateKind};
57use crate::memtable::MemtableId;
58use crate::memtable::bulk::part::BulkPart;
59use crate::metrics::COMPACTION_ELAPSED_TOTAL;
60use crate::region::options::RegionOptions;
61use crate::sst::file::FileMeta;
62use crate::sst::index::IndexBuildType;
63use crate::wal::EntryId;
64use crate::wal::entry_distributor::WalEntryReceiver;
65
66#[derive(Debug)]
68pub struct WriteRequest {
69 pub region_id: RegionId,
71 pub op_type: OpType,
73 pub rows: Rows,
75 pub name_to_index: HashMap<String, usize>,
77 pub has_null: Vec<bool>,
79 pub hint: Option<WriteHint>,
81 pub(crate) region_metadata: Option<RegionMetadataRef>,
83 pub partition_expr_version: Option<u64>,
85}
86
87impl WriteRequest {
88 pub fn new(
92 region_id: RegionId,
93 op_type: OpType,
94 rows: Rows,
95 region_metadata: Option<RegionMetadataRef>,
96 ) -> Result<WriteRequest> {
97 let mut name_to_index = HashMap::with_capacity(rows.schema.len());
98 for (index, column) in rows.schema.iter().enumerate() {
99 ensure!(
100 name_to_index
101 .insert(column.column_name.clone(), index)
102 .is_none(),
103 InvalidRequestSnafu {
104 region_id,
105 reason: format!("duplicate column {}", column.column_name),
106 }
107 );
108 }
109
110 let mut has_null = vec![false; rows.schema.len()];
111 for row in &rows.rows {
112 ensure!(
113 row.values.len() == rows.schema.len(),
114 InvalidRequestSnafu {
115 region_id,
116 reason: format!(
117 "row has {} columns but schema has {}",
118 row.values.len(),
119 rows.schema.len()
120 ),
121 }
122 );
123
124 for (i, (value, column_schema)) in row.values.iter().zip(&rows.schema).enumerate() {
125 validate_proto_value(region_id, value, column_schema)?;
126
127 if value.value_data.is_none() {
128 has_null[i] = true;
129 }
130 }
131 }
132
133 Ok(WriteRequest {
134 region_id,
135 op_type,
136 rows,
137 name_to_index,
138 has_null,
139 hint: None,
140 region_metadata,
141 partition_expr_version: None,
142 })
143 }
144
145 pub fn with_hint(mut self, hint: Option<WriteHint>) -> Self {
147 self.hint = hint;
148 self
149 }
150
151 pub fn with_partition_expr_version(mut self, partition_expr_version: Option<u64>) -> Self {
152 self.partition_expr_version = partition_expr_version;
153 self
154 }
155
156 pub fn primary_key_encoding(&self) -> PrimaryKeyEncoding {
158 infer_primary_key_encoding_from_hint(self.hint.as_ref())
159 }
160
161 pub(crate) fn estimated_size(&self) -> usize {
163 let row_size = self
164 .rows
165 .rows
166 .first()
167 .map(|row| row.encoded_len())
168 .unwrap_or(0);
169 row_size * self.rows.rows.len()
170 }
171
172 pub fn column_index_by_name(&self, name: &str) -> Option<usize> {
174 self.name_to_index.get(name).copied()
175 }
176
177 pub(crate) fn check_schema(&self, metadata: &RegionMetadata) -> Result<()> {
182 debug_assert_eq!(self.region_id, metadata.region_id);
183
184 let region_id = self.region_id;
185 let mut rows_columns: HashMap<_, _> = self
187 .rows
188 .schema
189 .iter()
190 .map(|column| (&column.column_name, column))
191 .collect();
192
193 let mut need_fill_default = false;
194 for column in &metadata.column_metadatas {
196 if let Some(input_col) = rows_columns.remove(&column.column_schema.name) {
197 ensure!(
199 is_column_type_value_eq(
200 input_col.datatype,
201 input_col.datatype_extension.clone(),
202 &column.column_schema.data_type
203 ),
204 InvalidRequestSnafu {
205 region_id,
206 reason: format!(
207 "column {} expect type {:?}, given: {}({})",
208 column.column_schema.name,
209 column.column_schema.data_type,
210 ColumnDataType::try_from(input_col.datatype)
211 .map(|v| v.as_str_name())
212 .unwrap_or("Unknown"),
213 input_col.datatype,
214 )
215 }
216 );
217
218 ensure!(
220 is_semantic_type_eq(input_col.semantic_type, column.semantic_type),
221 InvalidRequestSnafu {
222 region_id,
223 reason: format!(
224 "column {} has semantic type {:?}, given: {}({})",
225 column.column_schema.name,
226 column.semantic_type,
227 api::v1::SemanticType::try_from(input_col.semantic_type)
228 .map(|v| v.as_str_name())
229 .unwrap_or("Unknown"),
230 input_col.semantic_type
231 ),
232 }
233 );
234
235 let has_null = self.has_null[self.name_to_index[&column.column_schema.name]];
238 ensure!(
239 !has_null || column.column_schema.is_nullable(),
240 InvalidRequestSnafu {
241 region_id,
242 reason: format!(
243 "column {} is not null but input has null",
244 column.column_schema.name
245 ),
246 }
247 );
248 } else {
249 self.check_missing_column(column)?;
251
252 need_fill_default = true;
253 }
254 }
255
256 if !rows_columns.is_empty() {
258 let names: Vec<_> = rows_columns.into_keys().collect();
259 return InvalidRequestSnafu {
260 region_id,
261 reason: format!("unknown columns: {:?}", names),
262 }
263 .fail();
264 }
265
266 ensure!(!need_fill_default, FillDefaultSnafu { region_id });
268
269 Ok(())
270 }
271
272 pub(crate) fn fill_missing_columns(&mut self, metadata: &RegionMetadata) -> Result<()> {
277 debug_assert_eq!(self.region_id, metadata.region_id);
278
279 let mut columns_to_fill = vec![];
280 for column in &metadata.column_metadatas {
281 if !self.name_to_index.contains_key(&column.column_schema.name) {
282 columns_to_fill.push(column);
283 }
284 }
285 self.fill_columns(columns_to_fill)?;
286
287 Ok(())
288 }
289
290 pub(crate) fn maybe_fill_missing_columns(&mut self, metadata: &RegionMetadata) -> Result<()> {
292 if let Err(e) = self.check_schema(metadata) {
293 if e.is_fill_default() {
294 self.fill_missing_columns(metadata)?;
298 } else {
299 return Err(e);
300 }
301 }
302
303 Ok(())
304 }
305
306 fn fill_columns(&mut self, columns: Vec<&ColumnMetadata>) -> Result<()> {
308 let mut default_values = Vec::with_capacity(columns.len());
309 let mut columns_to_fill = Vec::with_capacity(columns.len());
310 for column in columns {
311 let default_value = self.column_default_value(column)?;
312 if default_value.value_data.is_some() {
313 default_values.push(default_value);
314 columns_to_fill.push(column);
315 }
316 }
317
318 for row in &mut self.rows.rows {
319 row.values.extend(default_values.iter().cloned());
320 }
321
322 for column in columns_to_fill {
323 let (datatype, datatype_ext) =
324 ColumnDataTypeWrapper::try_from(column.column_schema.data_type.clone())
325 .with_context(|_| ConvertColumnDataTypeSnafu {
326 reason: format!(
327 "no protobuf type for column {} ({:?})",
328 column.column_schema.name, column.column_schema.data_type
329 ),
330 })?
331 .to_parts();
332 self.rows.schema.push(ColumnSchema {
333 column_name: column.column_schema.name.clone(),
334 datatype: datatype as i32,
335 semantic_type: column.semantic_type as i32,
336 datatype_extension: datatype_ext,
337 options: options_from_column_schema(&column.column_schema),
338 });
339 }
340
341 Ok(())
342 }
343
344 fn check_missing_column(&self, column: &ColumnMetadata) -> Result<()> {
346 if self.op_type == OpType::Delete {
347 if column.semantic_type == SemanticType::Field {
348 return Ok(());
351 } else {
352 return InvalidRequestSnafu {
353 region_id: self.region_id,
354 reason: format!("delete requests need column {}", column.column_schema.name),
355 }
356 .fail();
357 }
358 }
359
360 ensure!(
362 column.column_schema.is_nullable()
363 || column.column_schema.default_constraint().is_some(),
364 InvalidRequestSnafu {
365 region_id: self.region_id,
366 reason: format!("missing column {}", column.column_schema.name),
367 }
368 );
369
370 Ok(())
371 }
372
373 fn column_default_value(&self, column: &ColumnMetadata) -> Result<Value> {
375 let default_value = match self.op_type {
376 OpType::Delete => {
377 ensure!(
378 column.semantic_type == SemanticType::Field,
379 InvalidRequestSnafu {
380 region_id: self.region_id,
381 reason: format!(
382 "delete requests need column {}",
383 column.column_schema.name
384 ),
385 }
386 );
387
388 if column.column_schema.is_nullable() {
393 datatypes::value::Value::Null
394 } else {
395 column.column_schema.data_type.default_value()
396 }
397 }
398 OpType::Put => {
399 if column.column_schema.is_default_impure() {
401 UnexpectedSnafu {
402 reason: format!(
403 "unexpected impure default value with region_id: {}, column: {}, default_value: {:?}",
404 self.region_id,
405 column.column_schema.name,
406 column.column_schema.default_constraint(),
407 ),
408 }
409 .fail()?
410 }
411 column
412 .column_schema
413 .create_default()
414 .context(CreateDefaultSnafu {
415 region_id: self.region_id,
416 column: &column.column_schema.name,
417 })?
418 .with_context(|| InvalidRequestSnafu {
420 region_id: self.region_id,
421 reason: format!(
422 "column {} does not have default value",
423 column.column_schema.name
424 ),
425 })?
426 }
427 };
428
429 Ok(api::helper::to_grpc_value(default_value))
431 }
432}
433
434pub(crate) fn validate_proto_value(
436 region_id: RegionId,
437 value: &Value,
438 column_schema: &ColumnSchema,
439) -> Result<()> {
440 if let Some(value_type) = proto_value_type(value) {
441 let column_type = ColumnDataType::try_from(column_schema.datatype).map_err(|_| {
442 InvalidRequestSnafu {
443 region_id,
444 reason: format!(
445 "column {} has unknown type {}",
446 column_schema.column_name, column_schema.datatype
447 ),
448 }
449 .build()
450 })?;
451 ensure!(
452 proto_value_type_match(column_type, value_type),
453 InvalidRequestSnafu {
454 region_id,
455 reason: format!(
456 "value has type {:?}, but column {} has type {:?}({})",
457 value_type, column_schema.column_name, column_type, column_schema.datatype,
458 ),
459 }
460 );
461 }
462
463 Ok(())
464}
465
466fn proto_value_type_match(column_type: ColumnDataType, value_type: ColumnDataType) -> bool {
467 match (column_type, value_type) {
468 (ct, vt) if ct == vt => true,
469 (ColumnDataType::Vector, ColumnDataType::Binary) => true,
470 (ColumnDataType::Json, ColumnDataType::Binary) => true,
471 _ => false,
472 }
473}
474
475#[derive(Debug)]
477pub struct OutputTx(Sender<Result<AffectedRows>>);
478
479impl OutputTx {
480 pub(crate) fn new(sender: Sender<Result<AffectedRows>>) -> OutputTx {
482 OutputTx(sender)
483 }
484
485 pub(crate) fn send(self, result: Result<AffectedRows>) {
487 let _ = self.0.send(result);
489 }
490}
491
492#[derive(Debug)]
494pub(crate) struct OptionOutputTx(Option<OutputTx>);
495
496impl OptionOutputTx {
497 pub(crate) fn new(sender: Option<OutputTx>) -> OptionOutputTx {
499 OptionOutputTx(sender)
500 }
501
502 pub(crate) fn none() -> OptionOutputTx {
504 OptionOutputTx(None)
505 }
506
507 pub(crate) fn send_mut(&mut self, result: Result<AffectedRows>) {
509 if let Some(sender) = self.0.take() {
510 sender.send(result);
511 }
512 }
513
514 pub(crate) fn send(mut self, result: Result<AffectedRows>) {
516 if let Some(sender) = self.0.take() {
517 sender.send(result);
518 }
519 }
520
521 pub(crate) fn take_inner(&mut self) -> Option<OutputTx> {
523 self.0.take()
524 }
525}
526
527impl From<Sender<Result<AffectedRows>>> for OptionOutputTx {
528 fn from(sender: Sender<Result<AffectedRows>>) -> Self {
529 Self::new(Some(OutputTx::new(sender)))
530 }
531}
532
533impl OnFailure for OptionOutputTx {
534 fn on_failure(&mut self, err: Error) {
535 self.send_mut(Err(err));
536 }
537}
538
539pub(crate) trait OnFailure {
541 fn on_failure(&mut self, err: Error);
543}
544
545#[derive(Debug)]
547pub(crate) struct SenderWriteRequest {
548 pub(crate) sender: OptionOutputTx,
550 pub(crate) request: WriteRequest,
551}
552
553pub(crate) struct SenderBulkRequest {
554 pub(crate) sender: OptionOutputTx,
555 pub(crate) region_id: RegionId,
556 pub(crate) request: BulkPart,
557 pub(crate) region_metadata: Option<RegionMetadataRef>,
558 pub(crate) partition_expr_version: Option<u64>,
559}
560
561#[derive(Debug)]
562pub(crate) struct BulkInsertRequest {
563 pub(crate) metadata: Option<RegionMetadataRef>,
564 pub(crate) request: RegionBulkInsertsRequest,
565 pub(crate) sender: OptionOutputTx,
566}
567
568#[derive(Debug)]
570pub(crate) struct WorkerRequestWithTime {
571 pub(crate) request: WorkerRequest,
572 pub(crate) created_at: Instant,
573}
574
575impl WorkerRequestWithTime {
576 pub(crate) fn new(request: WorkerRequest) -> Self {
577 Self {
578 request,
579 created_at: Instant::now(),
580 }
581 }
582}
583
584#[derive(Debug)]
586pub(crate) enum WorkerRequest {
587 Write(SenderWriteRequest),
589
590 Ddl(SenderDdlRequest),
592
593 Background {
595 region_id: RegionId,
597 notify: BackgroundNotify,
599 },
600
601 SetRegionRoleStateGracefully {
603 region_id: RegionId,
605 region_role_state: SettableRegionRoleState,
607 sender: Sender<SetRegionRoleStateResponse>,
609 },
610
611 Stop,
613
614 EditRegion(RegionEditRequest),
616
617 SyncRegion(RegionSyncRequest),
619
620 BulkInserts(BulkInsertRequest),
622
623 RemapManifests(RemapManifestsRequest),
625
626 CopyRegionFrom(CopyRegionFromRequest),
628}
629
630impl WorkerRequest {
631 pub(crate) fn new_open_region_request(
633 region_id: RegionId,
634 request: RegionOpenRequest,
635 entry_receiver: Option<WalEntryReceiver>,
636 ) -> (WorkerRequest, Receiver<Result<AffectedRows>>) {
637 let (sender, receiver) = oneshot::channel();
638
639 let worker_request = WorkerRequest::Ddl(SenderDdlRequest {
640 region_id,
641 sender: sender.into(),
642 request: DdlRequest::Open((request, entry_receiver)),
643 });
644
645 (worker_request, receiver)
646 }
647
648 pub(crate) fn new_catchup_region_request(
650 region_id: RegionId,
651 request: RegionCatchupRequest,
652 entry_receiver: Option<WalEntryReceiver>,
653 ) -> (WorkerRequest, Receiver<Result<AffectedRows>>) {
654 let (sender, receiver) = oneshot::channel();
655 let worker_request = WorkerRequest::Ddl(SenderDdlRequest {
656 region_id,
657 sender: sender.into(),
658 request: DdlRequest::Catchup((request, entry_receiver)),
659 });
660 (worker_request, receiver)
661 }
662
663 pub(crate) fn try_from_region_request(
665 region_id: RegionId,
666 value: RegionRequest,
667 region_metadata: Option<RegionMetadataRef>,
668 ) -> Result<(WorkerRequest, Receiver<Result<AffectedRows>>)> {
669 let (sender, receiver) = oneshot::channel();
670 let worker_request = match value {
671 RegionRequest::Put(v) => {
672 let mut write_request =
673 WriteRequest::new(region_id, OpType::Put, v.rows, region_metadata.clone())?
674 .with_hint(v.hint)
675 .with_partition_expr_version(v.partition_expr_version);
676 if write_request.primary_key_encoding() == PrimaryKeyEncoding::Dense
677 && let Some(region_metadata) = ®ion_metadata
678 {
679 write_request.maybe_fill_missing_columns(region_metadata)?;
680 }
681 WorkerRequest::Write(SenderWriteRequest {
682 sender: sender.into(),
683 request: write_request,
684 })
685 }
686 RegionRequest::Delete(v) => {
687 let mut write_request =
688 WriteRequest::new(region_id, OpType::Delete, v.rows, region_metadata.clone())?
689 .with_hint(v.hint)
690 .with_partition_expr_version(v.partition_expr_version);
691 if write_request.primary_key_encoding() == PrimaryKeyEncoding::Dense
692 && let Some(region_metadata) = ®ion_metadata
693 {
694 write_request.maybe_fill_missing_columns(region_metadata)?;
695 }
696 WorkerRequest::Write(SenderWriteRequest {
697 sender: sender.into(),
698 request: write_request,
699 })
700 }
701 RegionRequest::Create(v) => WorkerRequest::Ddl(SenderDdlRequest {
702 region_id,
703 sender: sender.into(),
704 request: DdlRequest::Create(v),
705 }),
706 RegionRequest::Drop(v) => WorkerRequest::Ddl(SenderDdlRequest {
707 region_id,
708 sender: sender.into(),
709 request: DdlRequest::Drop(v),
710 }),
711 RegionRequest::Open(v) => WorkerRequest::Ddl(SenderDdlRequest {
712 region_id,
713 sender: sender.into(),
714 request: DdlRequest::Open((v, None)),
715 }),
716 RegionRequest::CleanUp(v) => WorkerRequest::Ddl(SenderDdlRequest {
717 region_id,
718 sender: sender.into(),
719 request: DdlRequest::OfflineCleanup(v),
720 }),
721 RegionRequest::Close(v) => WorkerRequest::Ddl(SenderDdlRequest {
722 region_id,
723 sender: sender.into(),
724 request: DdlRequest::Close(v),
725 }),
726 RegionRequest::Alter(v) => WorkerRequest::Ddl(SenderDdlRequest {
727 region_id,
728 sender: sender.into(),
729 request: DdlRequest::Alter(v),
730 }),
731 RegionRequest::Flush(v) => WorkerRequest::Ddl(SenderDdlRequest {
732 region_id,
733 sender: sender.into(),
734 request: DdlRequest::Flush(v),
735 }),
736 RegionRequest::Compact(v) => WorkerRequest::Ddl(SenderDdlRequest {
737 region_id,
738 sender: sender.into(),
739 request: DdlRequest::Compact(v),
740 }),
741 RegionRequest::BuildIndex(v) => WorkerRequest::Ddl(SenderDdlRequest {
742 region_id,
743 sender: sender.into(),
744 request: DdlRequest::BuildIndex(v),
745 }),
746 RegionRequest::Truncate(v) => WorkerRequest::Ddl(SenderDdlRequest {
747 region_id,
748 sender: sender.into(),
749 request: DdlRequest::Truncate(v),
750 }),
751 RegionRequest::Catchup(v) => WorkerRequest::Ddl(SenderDdlRequest {
752 region_id,
753 sender: sender.into(),
754 request: DdlRequest::Catchup((v, None)),
755 }),
756 RegionRequest::EnterStaging(v) => WorkerRequest::Ddl(SenderDdlRequest {
757 region_id,
758 sender: sender.into(),
759 request: DdlRequest::EnterStaging(v),
760 }),
761 RegionRequest::BulkInserts(region_bulk_inserts_request) => {
762 WorkerRequest::BulkInserts(BulkInsertRequest {
763 metadata: region_metadata,
764 sender: sender.into(),
765 request: region_bulk_inserts_request,
766 })
767 }
768 RegionRequest::ApplyStagingManifest(v) => WorkerRequest::Ddl(SenderDdlRequest {
769 region_id,
770 sender: sender.into(),
771 request: DdlRequest::ApplyStagingManifest(v),
772 }),
773 };
774
775 Ok((worker_request, receiver))
776 }
777
778 pub(crate) fn new_set_readonly_gracefully(
779 region_id: RegionId,
780 region_role_state: SettableRegionRoleState,
781 ) -> (WorkerRequest, Receiver<SetRegionRoleStateResponse>) {
782 let (sender, receiver) = oneshot::channel();
783
784 (
785 WorkerRequest::SetRegionRoleStateGracefully {
786 region_id,
787 region_role_state,
788 sender,
789 },
790 receiver,
791 )
792 }
793
794 pub(crate) fn new_sync_region_request(
795 region_id: RegionId,
796 manifest_version: ManifestVersion,
797 ) -> (WorkerRequest, Receiver<Result<(ManifestVersion, bool)>>) {
798 let (sender, receiver) = oneshot::channel();
799 (
800 WorkerRequest::SyncRegion(RegionSyncRequest {
801 region_id,
802 manifest_version,
803 sender,
804 }),
805 receiver,
806 )
807 }
808
809 #[allow(clippy::type_complexity)]
816 pub(crate) fn try_from_remap_manifests_request(
817 store_api::region_engine::RemapManifestsRequest {
818 region_id,
819 input_regions,
820 region_mapping,
821 new_partition_exprs,
822 }: store_api::region_engine::RemapManifestsRequest,
823 ) -> Result<(WorkerRequest, Receiver<Result<HashMap<RegionId, String>>>)> {
824 let (sender, receiver) = oneshot::channel();
825 let new_partition_exprs = new_partition_exprs
826 .into_iter()
827 .map(|(k, v)| {
828 Ok((
829 k,
830 PartitionExpr::from_json_str(&v)
831 .context(InvalidPartitionExprSnafu { expr: v })?
832 .context(MissingPartitionExprSnafu { region_id: k })?,
833 ))
834 })
835 .collect::<Result<HashMap<_, _>>>()?;
836
837 let request = RemapManifestsRequest {
838 region_id,
839 input_regions,
840 region_mapping,
841 new_partition_exprs,
842 sender,
843 };
844
845 Ok((WorkerRequest::RemapManifests(request), receiver))
846 }
847
848 pub(crate) fn try_from_copy_region_from_request(
850 region_id: RegionId,
851 store_api::region_engine::MitoCopyRegionFromRequest {
852 source_region_id,
853 parallelism,
854 }: store_api::region_engine::MitoCopyRegionFromRequest,
855 ) -> Result<(WorkerRequest, Receiver<Result<MitoCopyRegionFromResponse>>)> {
856 let (sender, receiver) = oneshot::channel();
857 let request = CopyRegionFromRequest {
858 region_id,
859 source_region_id,
860 parallelism,
861 sender,
862 };
863 Ok((WorkerRequest::CopyRegionFrom(request), receiver))
864 }
865}
866
867#[derive(Debug)]
869pub(crate) enum DdlRequest {
870 Create(RegionCreateRequest),
871 Drop(RegionDropRequest),
872 Open((RegionOpenRequest, Option<WalEntryReceiver>)),
873 OfflineCleanup(RegionCleanUpRequest),
874 Close(RegionCloseRequest),
875 Alter(RegionAlterRequest),
876 Flush(RegionFlushRequest),
877 Compact(RegionCompactRequest),
878 BuildIndex(RegionBuildIndexRequest),
879 Truncate(RegionTruncateRequest),
880 Catchup((RegionCatchupRequest, Option<WalEntryReceiver>)),
881 EnterStaging(EnterStagingRequest),
882 ApplyStagingManifest(ApplyStagingManifestRequest),
883}
884
885#[derive(Debug)]
887pub(crate) struct SenderDdlRequest {
888 pub(crate) region_id: RegionId,
890 pub(crate) sender: OptionOutputTx,
892 pub(crate) request: DdlRequest,
894}
895
896#[derive(Debug)]
898pub(crate) enum BackgroundNotify {
899 CompactionPickFinished(CompactionPickFinished),
901 FlushFinished(FlushFinished),
903 FlushFailed(FlushFailed),
905 IndexBuildFinished(IndexBuildFinished),
907 IndexBuildStopped(IndexBuildStopped),
909 IndexBuildFailed(IndexBuildFailed),
911 IndexBuildRetry(BuildIndexRequest),
913 CompactionFinished(CompactionFinished),
915 CompactionCancelled(CompactionCancelled),
917 CompactionFailed(CompactionFailed),
919 Truncate(TruncateResult),
921 DiscardUnflushed(DiscardUnflushedResult),
923 RegionChange(RegionChangeResult),
925 RegionEdit(RegionEditResult),
927 EnterStaging(EnterStagingResult),
929 CopyRegionFromFinished(CopyRegionFromFinished),
931}
932
933#[derive(Debug)]
935pub(crate) struct FlushFinished {
936 pub(crate) region_id: RegionId,
938 pub(crate) flush_reason: FlushReason,
940 pub(crate) flushed_entry_id: EntryId,
942 pub(crate) senders: Vec<OutputTx>,
944 pub(crate) _timer: HistogramTimer,
946 pub(crate) edit: RegionEdit,
948 pub(crate) memtables_to_remove: SmallVec<[MemtableId; 2]>,
950 pub(crate) is_staging: bool,
952}
953
954impl FlushFinished {
955 pub(crate) fn on_success(self) {
957 for sender in self.senders {
958 sender.send(Ok(0));
959 }
960 }
961}
962
963impl OnFailure for FlushFinished {
964 fn on_failure(&mut self, err: Error) {
965 let err = Arc::new(err);
966 for sender in self.senders.drain(..) {
967 sender.send(Err(err.clone()).context(FlushRegionSnafu {
968 region_id: self.region_id,
969 }));
970 }
971 }
972}
973
974#[derive(Debug)]
976pub(crate) struct FlushFailed {
977 pub(crate) err: Arc<Error>,
979}
980
981impl FlushFailed {
982 pub(crate) fn is_cancelled(&self) -> bool {
984 matches!(self.err.as_ref(), Error::FlushCancelled { .. })
985 }
986}
987
988#[derive(Debug)]
989pub(crate) struct IndexBuildFinished {
990 pub(crate) manifest_version: ManifestVersion,
991 pub(crate) file_meta: FileMeta,
992}
993
994#[derive(Debug)]
996pub(crate) struct IndexBuildStopped {
997 pub(crate) file_id: FileId,
998}
999
1000#[derive(Debug)]
1002pub(crate) struct IndexBuildFailed {
1003 pub(crate) err: Arc<Error>,
1004}
1005
1006#[derive(Debug)]
1008pub(crate) struct CompactionFinished {
1009 pub(crate) region_id: RegionId,
1011 pub(crate) execution: CompactionExecution,
1013 pub(crate) senders: Vec<OutputTx>,
1015 pub(crate) start_time: Instant,
1017 pub(crate) edit: RegionEdit,
1019}
1020
1021#[derive(Debug)]
1023pub(crate) struct CompactionCancelled {
1024 pub(crate) region_id: RegionId,
1026 pub(crate) execution: CompactionExecution,
1028 pub(crate) senders: Vec<OutputTx>,
1030}
1031
1032impl CompactionCancelled {
1033 pub(crate) fn on_success(self) {
1034 for sender in self.senders {
1035 sender.send(CompactionCancelledSnafu {}.fail());
1036 }
1037 info!("Compaction cancelled for region: {}", self.region_id);
1038 }
1039}
1040
1041impl CompactionFinished {
1042 pub fn on_success(self) {
1043 COMPACTION_ELAPSED_TOTAL.observe(self.start_time.elapsed().as_secs_f64());
1045
1046 for sender in self.senders {
1047 sender.send(Ok(0));
1048 }
1049 info!("Successfully compacted region: {}", self.region_id);
1050 }
1051}
1052
1053impl OnFailure for CompactionFinished {
1054 fn on_failure(&mut self, err: Error) {
1056 let err = Arc::new(err);
1057 for sender in self.senders.drain(..) {
1058 sender.send(Err(err.clone()).context(CompactRegionSnafu {
1059 region_id: self.region_id,
1060 }));
1061 }
1062 }
1063}
1064
1065#[derive(Debug)]
1067pub(crate) struct CompactionFailed {
1068 pub(crate) region_id: RegionId,
1069 pub(crate) execution: CompactionExecution,
1071 pub(crate) err: Arc<Error>,
1073}
1074
1075#[derive(Debug)]
1077pub(crate) struct TruncateResult {
1078 pub(crate) region_id: RegionId,
1080 pub(crate) sender: OptionOutputTx,
1082 pub(crate) result: Result<()>,
1084 pub(crate) kind: TruncateKind,
1085}
1086
1087#[derive(Debug)]
1089pub(crate) struct DiscardUnflushedResult {
1090 pub(crate) region_id: RegionId,
1092 pub(crate) sender: OptionOutputTx,
1094 pub(crate) result: Result<()>,
1096 pub(crate) discarded_entry_id: EntryId,
1098 pub(crate) discarded_sequence: SequenceNumber,
1100 pub(crate) discarded_rows: u64,
1102 pub(crate) discarded_bytes: u64,
1104}
1105
1106#[derive(Debug)]
1108pub(crate) struct RegionChangeResult {
1109 pub(crate) region_id: RegionId,
1111 pub(crate) new_meta: RegionMetadataRef,
1113 pub(crate) sender: OptionOutputTx,
1115 pub(crate) result: Result<()>,
1117 pub(crate) need_index: bool,
1119 pub(crate) new_options: Option<RegionOptions>,
1121}
1122
1123#[derive(Debug)]
1125pub(crate) struct EnterStagingResult {
1126 pub(crate) region_id: RegionId,
1128 pub(crate) partition_directive: StagingPartitionDirective,
1130 pub(crate) sender: OptionOutputTx,
1132 pub(crate) result: Result<()>,
1134}
1135
1136#[derive(Debug)]
1137pub(crate) struct CopyRegionFromFinished {
1138 pub(crate) region_id: RegionId,
1140 pub(crate) edit: RegionEdit,
1142 pub(crate) sender: Sender<Result<MitoCopyRegionFromResponse>>,
1144}
1145
1146#[derive(Debug, Default)]
1147pub(crate) struct Waiters(SmallVec<[Sender<Result<()>>; 1]>);
1148
1149impl Waiters {
1150 pub(crate) fn one(waiter: Sender<Result<()>>) -> Self {
1151 let mut waiters = SmallVec::new();
1152 waiters.push(waiter);
1153 Self(waiters)
1154 }
1155
1156 pub(crate) fn reply_with<F: Fn() -> Result<()>>(self, f: F) {
1157 for tx in self.0 {
1158 let _ = tx.send(f());
1159 }
1160 }
1161
1162 pub(crate) fn merge(&mut self, other: Self) {
1163 self.0.extend(other.0);
1164 }
1165}
1166
1167#[derive(Debug)]
1169pub(crate) struct RegionEditRequest {
1170 pub(crate) region_id: RegionId,
1171 pub(crate) edit: RegionEdit,
1172 pub(crate) preload_sst_cache: bool,
1174 pub(crate) waiters: Waiters,
1176}
1177
1178impl RegionEditRequest {
1179 pub(crate) fn new(
1180 region_id: RegionId,
1181 edit: RegionEdit,
1182 preload_sst_cache: bool,
1183 waiter: Sender<Result<()>>,
1184 ) -> Self {
1185 Self {
1186 region_id,
1187 edit,
1188 preload_sst_cache,
1189 waiters: Waiters::one(waiter),
1190 }
1191 }
1192}
1193
1194#[derive(Debug)]
1196pub(crate) struct RegionEditResult {
1197 pub(crate) region_id: RegionId,
1199 pub(crate) waiters: Waiters,
1201 pub(crate) edit: RegionEdit,
1203 pub(crate) result: std::result::Result<(), Arc<Error>>,
1205 pub(crate) update_region_state: bool,
1207 pub(crate) is_staging: bool,
1209}
1210
1211#[derive(Debug)]
1212pub(crate) struct BuildIndexRequest {
1213 pub(crate) region_id: RegionId,
1214 pub(crate) build_type: IndexBuildType,
1215 pub(crate) file_metas: Vec<FileMeta>,
1217}
1218
1219#[derive(Debug)]
1220pub(crate) struct RegionSyncRequest {
1221 pub(crate) region_id: RegionId,
1222 pub(crate) manifest_version: ManifestVersion,
1223 pub(crate) sender: Sender<Result<(ManifestVersion, bool)>>,
1225}
1226
1227#[derive(Debug)]
1228pub(crate) struct RemapManifestsRequest {
1229 pub(crate) region_id: RegionId,
1231 pub(crate) input_regions: Vec<RegionId>,
1233 pub(crate) region_mapping: HashMap<RegionId, Vec<RegionId>>,
1235 pub(crate) new_partition_exprs: HashMap<RegionId, PartitionExpr>,
1237 pub(crate) sender: Sender<Result<HashMap<RegionId, String>>>,
1241}
1242
1243#[derive(Debug)]
1244pub(crate) struct CopyRegionFromRequest {
1245 pub(crate) region_id: RegionId,
1247 pub(crate) source_region_id: RegionId,
1249 pub(crate) parallelism: usize,
1251 pub(crate) sender: Sender<Result<MitoCopyRegionFromResponse>>,
1253}
1254
1255#[cfg(test)]
1256mod tests {
1257 use api::v1::value::ValueData;
1258 use api::v1::{Row, SemanticType};
1259 use common_error::ext::ErrorExt;
1260 use common_error::status_code::StatusCode;
1261 use datatypes::prelude::ConcreteDataType;
1262 use datatypes::schema::ColumnDefaultConstraint;
1263 use mito_codec::test_util::i64_value;
1264 use store_api::metadata::RegionMetadataBuilder;
1265 use tokio::sync::oneshot;
1266
1267 use super::*;
1268 use crate::error::Error;
1269 use crate::test_util::ts_ms_value;
1270
1271 fn new_column_schema(
1272 name: &str,
1273 data_type: ColumnDataType,
1274 semantic_type: SemanticType,
1275 ) -> ColumnSchema {
1276 ColumnSchema {
1277 column_name: name.to_string(),
1278 datatype: data_type as i32,
1279 semantic_type: semantic_type as i32,
1280 ..Default::default()
1281 }
1282 }
1283
1284 fn check_invalid_request(err: &Error, expect: &str) {
1285 if let Error::InvalidRequest {
1286 region_id: _,
1287 reason,
1288 location: _,
1289 } = err
1290 {
1291 assert_eq!(reason, expect);
1292 } else {
1293 panic!("Unexpected error {err}")
1294 }
1295 }
1296
1297 fn waiter() -> (Sender<Result<()>>, Receiver<Result<()>>) {
1298 oneshot::channel()
1299 }
1300
1301 fn assert_waiter_ok(rx: &mut Receiver<Result<()>>) {
1302 rx.try_recv().unwrap().unwrap();
1303 }
1304
1305 #[test]
1306 fn test_waiters_reply_with_single_waiter() {
1307 let (tx, mut rx) = waiter();
1308 Waiters::one(tx).reply_with(|| Ok(()));
1309 assert_waiter_ok(&mut rx);
1310 }
1311
1312 #[test]
1313 fn test_waiters_reply_with_many_waiters() {
1314 let (tx1, mut rx1) = waiter();
1315 let (tx2, mut rx2) = waiter();
1316 let (tx3, mut rx3) = waiter();
1317
1318 let waiters = Waiters(vec![tx1, tx2, tx3].into());
1319 waiters.reply_with(|| Ok(()));
1320
1321 assert_waiter_ok(&mut rx1);
1322 assert_waiter_ok(&mut rx2);
1323 assert_waiter_ok(&mut rx3);
1324 }
1325
1326 #[test]
1327 fn test_waiters_merge() {
1328 let (tx1, mut rx1) = waiter();
1329 let (tx2, mut rx2) = waiter();
1330 let (tx3, mut rx3) = waiter();
1331 let (tx4, mut rx4) = waiter();
1332
1333 let mut waiters = Waiters::one(tx1);
1334 waiters.merge(Waiters::one(tx2));
1335 waiters.merge(Waiters(vec![tx3, tx4].into()));
1336 assert_eq!(4, waiters.0.len());
1337
1338 waiters.reply_with(|| Ok(()));
1339
1340 assert_waiter_ok(&mut rx1);
1341 assert_waiter_ok(&mut rx2);
1342 assert_waiter_ok(&mut rx3);
1343 assert_waiter_ok(&mut rx4);
1344 }
1345
1346 #[test]
1347 fn test_write_request_duplicate_column() {
1348 let rows = Rows {
1349 schema: vec![
1350 new_column_schema("c0", ColumnDataType::Int64, SemanticType::Tag),
1351 new_column_schema("c0", ColumnDataType::Int64, SemanticType::Tag),
1352 ],
1353 rows: vec![],
1354 };
1355
1356 let err = WriteRequest::new(RegionId::new(1, 1), OpType::Put, rows, None).unwrap_err();
1357 check_invalid_request(&err, "duplicate column c0");
1358 }
1359
1360 #[test]
1361 fn test_valid_write_request() {
1362 let rows = Rows {
1363 schema: vec![
1364 new_column_schema("c0", ColumnDataType::Int64, SemanticType::Tag),
1365 new_column_schema("c1", ColumnDataType::Int64, SemanticType::Tag),
1366 ],
1367 rows: vec![Row {
1368 values: vec![i64_value(1), i64_value(2)],
1369 }],
1370 };
1371
1372 let request = WriteRequest::new(RegionId::new(1, 1), OpType::Put, rows, None).unwrap();
1373 assert_eq!(0, request.column_index_by_name("c0").unwrap());
1374 assert_eq!(1, request.column_index_by_name("c1").unwrap());
1375 assert_eq!(None, request.column_index_by_name("c2"));
1376 }
1377
1378 #[test]
1379 fn test_compaction_cancelled_sends_cancelled_error() {
1380 let (tx, rx) = oneshot::channel();
1381 let request = CompactionCancelled {
1382 region_id: RegionId::new(1, 1),
1383 execution: crate::compaction::CompactionExecution::for_test(0),
1384 senders: vec![OutputTx::new(tx)],
1385 };
1386
1387 request.on_success();
1388
1389 let err = rx.blocking_recv().unwrap().unwrap_err();
1390 assert!(matches!(err, Error::CompactionCancelled { .. }));
1391 assert_eq!(err.status_code(), StatusCode::Cancelled);
1392 }
1393
1394 #[test]
1395 fn test_write_request_column_num() {
1396 let rows = Rows {
1397 schema: vec![
1398 new_column_schema("c0", ColumnDataType::Int64, SemanticType::Tag),
1399 new_column_schema("c1", ColumnDataType::Int64, SemanticType::Tag),
1400 ],
1401 rows: vec![Row {
1402 values: vec![i64_value(1), i64_value(2), i64_value(3)],
1403 }],
1404 };
1405
1406 let err = WriteRequest::new(RegionId::new(1, 1), OpType::Put, rows, None).unwrap_err();
1407 check_invalid_request(&err, "row has 3 columns but schema has 2");
1408 }
1409
1410 fn new_region_metadata() -> RegionMetadata {
1411 let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 1));
1412 builder
1413 .push_column_metadata(ColumnMetadata {
1414 column_schema: datatypes::schema::ColumnSchema::new(
1415 "ts",
1416 ConcreteDataType::timestamp_millisecond_datatype(),
1417 false,
1418 ),
1419 semantic_type: SemanticType::Timestamp,
1420 column_id: 1,
1421 })
1422 .push_column_metadata(ColumnMetadata {
1423 column_schema: datatypes::schema::ColumnSchema::new(
1424 "k0",
1425 ConcreteDataType::int64_datatype(),
1426 true,
1427 ),
1428 semantic_type: SemanticType::Tag,
1429 column_id: 2,
1430 })
1431 .primary_key(vec![2]);
1432 builder.build().unwrap()
1433 }
1434
1435 #[test]
1436 fn test_check_schema() {
1437 let rows = Rows {
1438 schema: vec![
1439 new_column_schema(
1440 "ts",
1441 ColumnDataType::TimestampMillisecond,
1442 SemanticType::Timestamp,
1443 ),
1444 new_column_schema("k0", ColumnDataType::Int64, SemanticType::Tag),
1445 ],
1446 rows: vec![Row {
1447 values: vec![ts_ms_value(1), i64_value(2)],
1448 }],
1449 };
1450 let metadata = new_region_metadata();
1451
1452 let request = WriteRequest::new(RegionId::new(1, 1), OpType::Put, rows, None).unwrap();
1453 request.check_schema(&metadata).unwrap();
1454 }
1455
1456 #[test]
1457 fn test_column_type() {
1458 let rows = Rows {
1459 schema: vec![
1460 new_column_schema("ts", ColumnDataType::Int64, SemanticType::Timestamp),
1461 new_column_schema("k0", ColumnDataType::Int64, SemanticType::Tag),
1462 ],
1463 rows: vec![Row {
1464 values: vec![i64_value(1), i64_value(2)],
1465 }],
1466 };
1467 let metadata = new_region_metadata();
1468
1469 let request = WriteRequest::new(RegionId::new(1, 1), OpType::Put, rows, None).unwrap();
1470 let err = request.check_schema(&metadata).unwrap_err();
1471 check_invalid_request(
1472 &err,
1473 "column ts expect type Timestamp(Millisecond(TimestampMillisecondType)), given: INT64(4)",
1474 );
1475 }
1476
1477 #[test]
1478 fn test_semantic_type() {
1479 let rows = Rows {
1480 schema: vec![
1481 new_column_schema(
1482 "ts",
1483 ColumnDataType::TimestampMillisecond,
1484 SemanticType::Tag,
1485 ),
1486 new_column_schema("k0", ColumnDataType::Int64, SemanticType::Tag),
1487 ],
1488 rows: vec![Row {
1489 values: vec![ts_ms_value(1), i64_value(2)],
1490 }],
1491 };
1492 let metadata = new_region_metadata();
1493
1494 let request = WriteRequest::new(RegionId::new(1, 1), OpType::Put, rows, None).unwrap();
1495 let err = request.check_schema(&metadata).unwrap_err();
1496 check_invalid_request(&err, "column ts has semantic type Timestamp, given: TAG(0)");
1497 }
1498
1499 #[test]
1500 fn test_column_nullable() {
1501 let rows = Rows {
1502 schema: vec![
1503 new_column_schema(
1504 "ts",
1505 ColumnDataType::TimestampMillisecond,
1506 SemanticType::Timestamp,
1507 ),
1508 new_column_schema("k0", ColumnDataType::Int64, SemanticType::Tag),
1509 ],
1510 rows: vec![Row {
1511 values: vec![Value { value_data: None }, i64_value(2)],
1512 }],
1513 };
1514 let metadata = new_region_metadata();
1515
1516 let request = WriteRequest::new(RegionId::new(1, 1), OpType::Put, rows, None).unwrap();
1517 let err = request.check_schema(&metadata).unwrap_err();
1518 check_invalid_request(&err, "column ts is not null but input has null");
1519 }
1520
1521 #[test]
1522 fn test_column_default() {
1523 let rows = Rows {
1524 schema: vec![new_column_schema(
1525 "k0",
1526 ColumnDataType::Int64,
1527 SemanticType::Tag,
1528 )],
1529 rows: vec![Row {
1530 values: vec![i64_value(1)],
1531 }],
1532 };
1533 let metadata = new_region_metadata();
1534
1535 let request = WriteRequest::new(RegionId::new(1, 1), OpType::Put, rows, None).unwrap();
1536 let err = request.check_schema(&metadata).unwrap_err();
1537 check_invalid_request(&err, "missing column ts");
1538 }
1539
1540 #[test]
1541 fn test_unknown_column() {
1542 let rows = Rows {
1543 schema: vec![
1544 new_column_schema(
1545 "ts",
1546 ColumnDataType::TimestampMillisecond,
1547 SemanticType::Timestamp,
1548 ),
1549 new_column_schema("k0", ColumnDataType::Int64, SemanticType::Tag),
1550 new_column_schema("k1", ColumnDataType::Int64, SemanticType::Tag),
1551 ],
1552 rows: vec![Row {
1553 values: vec![ts_ms_value(1), i64_value(2), i64_value(3)],
1554 }],
1555 };
1556 let metadata = new_region_metadata();
1557
1558 let request = WriteRequest::new(RegionId::new(1, 1), OpType::Put, rows, None).unwrap();
1559 let err = request.check_schema(&metadata).unwrap_err();
1560 check_invalid_request(&err, r#"unknown columns: ["k1"]"#);
1561 }
1562
1563 #[test]
1564 fn test_fill_impure_columns_err() {
1565 let rows = Rows {
1566 schema: vec![new_column_schema(
1567 "k0",
1568 ColumnDataType::Int64,
1569 SemanticType::Tag,
1570 )],
1571 rows: vec![Row {
1572 values: vec![i64_value(1)],
1573 }],
1574 };
1575 let metadata = {
1576 let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 1));
1577 builder
1578 .push_column_metadata(ColumnMetadata {
1579 column_schema: datatypes::schema::ColumnSchema::new(
1580 "ts",
1581 ConcreteDataType::timestamp_millisecond_datatype(),
1582 false,
1583 )
1584 .with_default_constraint(Some(ColumnDefaultConstraint::Function(
1585 "now()".to_string(),
1586 )))
1587 .unwrap(),
1588 semantic_type: SemanticType::Timestamp,
1589 column_id: 1,
1590 })
1591 .push_column_metadata(ColumnMetadata {
1592 column_schema: datatypes::schema::ColumnSchema::new(
1593 "k0",
1594 ConcreteDataType::int64_datatype(),
1595 true,
1596 ),
1597 semantic_type: SemanticType::Tag,
1598 column_id: 2,
1599 })
1600 .primary_key(vec![2]);
1601 builder.build().unwrap()
1602 };
1603
1604 let mut request = WriteRequest::new(RegionId::new(1, 1), OpType::Put, rows, None).unwrap();
1605 let err = request.check_schema(&metadata).unwrap_err();
1606 assert!(err.is_fill_default());
1607 assert!(
1608 request
1609 .fill_missing_columns(&metadata)
1610 .unwrap_err()
1611 .to_string()
1612 .contains("unexpected impure default value with region_id")
1613 );
1614 }
1615
1616 #[test]
1617 fn test_fill_missing_columns() {
1618 let rows = Rows {
1619 schema: vec![new_column_schema(
1620 "ts",
1621 ColumnDataType::TimestampMillisecond,
1622 SemanticType::Timestamp,
1623 )],
1624 rows: vec![Row {
1625 values: vec![ts_ms_value(1)],
1626 }],
1627 };
1628 let metadata = new_region_metadata();
1629
1630 let mut request = WriteRequest::new(RegionId::new(1, 1), OpType::Put, rows, None).unwrap();
1631 let err = request.check_schema(&metadata).unwrap_err();
1632 assert!(err.is_fill_default());
1633 request.fill_missing_columns(&metadata).unwrap();
1634
1635 let expect_rows = Rows {
1636 schema: vec![new_column_schema(
1637 "ts",
1638 ColumnDataType::TimestampMillisecond,
1639 SemanticType::Timestamp,
1640 )],
1641 rows: vec![Row {
1642 values: vec![ts_ms_value(1)],
1643 }],
1644 };
1645 assert_eq!(expect_rows, request.rows);
1646 }
1647
1648 fn builder_with_ts_tag() -> RegionMetadataBuilder {
1649 let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 1));
1650 builder
1651 .push_column_metadata(ColumnMetadata {
1652 column_schema: datatypes::schema::ColumnSchema::new(
1653 "ts",
1654 ConcreteDataType::timestamp_millisecond_datatype(),
1655 false,
1656 ),
1657 semantic_type: SemanticType::Timestamp,
1658 column_id: 1,
1659 })
1660 .push_column_metadata(ColumnMetadata {
1661 column_schema: datatypes::schema::ColumnSchema::new(
1662 "k0",
1663 ConcreteDataType::int64_datatype(),
1664 true,
1665 ),
1666 semantic_type: SemanticType::Tag,
1667 column_id: 2,
1668 })
1669 .primary_key(vec![2]);
1670 builder
1671 }
1672
1673 fn region_metadata_two_fields() -> RegionMetadata {
1674 let mut builder = builder_with_ts_tag();
1675 builder
1676 .push_column_metadata(ColumnMetadata {
1677 column_schema: datatypes::schema::ColumnSchema::new(
1678 "f0",
1679 ConcreteDataType::int64_datatype(),
1680 true,
1681 ),
1682 semantic_type: SemanticType::Field,
1683 column_id: 3,
1684 })
1685 .push_column_metadata(ColumnMetadata {
1687 column_schema: datatypes::schema::ColumnSchema::new(
1688 "f1",
1689 ConcreteDataType::int64_datatype(),
1690 false,
1691 )
1692 .with_default_constraint(Some(ColumnDefaultConstraint::Value(
1693 datatypes::value::Value::Int64(100),
1694 )))
1695 .unwrap(),
1696 semantic_type: SemanticType::Field,
1697 column_id: 4,
1698 });
1699 builder.build().unwrap()
1700 }
1701
1702 #[test]
1703 fn test_fill_missing_for_delete() {
1704 let rows = Rows {
1705 schema: vec![new_column_schema(
1706 "ts",
1707 ColumnDataType::TimestampMillisecond,
1708 SemanticType::Timestamp,
1709 )],
1710 rows: vec![Row {
1711 values: vec![ts_ms_value(1)],
1712 }],
1713 };
1714 let metadata = region_metadata_two_fields();
1715
1716 let mut request =
1717 WriteRequest::new(RegionId::new(1, 1), OpType::Delete, rows, None).unwrap();
1718 let err = request.check_schema(&metadata).unwrap_err();
1719 check_invalid_request(&err, "delete requests need column k0");
1720 let err = request.fill_missing_columns(&metadata).unwrap_err();
1721 check_invalid_request(&err, "delete requests need column k0");
1722
1723 let rows = Rows {
1724 schema: vec![
1725 new_column_schema("k0", ColumnDataType::Int64, SemanticType::Tag),
1726 new_column_schema(
1727 "ts",
1728 ColumnDataType::TimestampMillisecond,
1729 SemanticType::Timestamp,
1730 ),
1731 ],
1732 rows: vec![Row {
1733 values: vec![i64_value(100), ts_ms_value(1)],
1734 }],
1735 };
1736 let mut request =
1737 WriteRequest::new(RegionId::new(1, 1), OpType::Delete, rows, None).unwrap();
1738 let err = request.check_schema(&metadata).unwrap_err();
1739 assert!(err.is_fill_default());
1740 request.fill_missing_columns(&metadata).unwrap();
1741
1742 let expect_rows = Rows {
1743 schema: vec![
1744 new_column_schema("k0", ColumnDataType::Int64, SemanticType::Tag),
1745 new_column_schema(
1746 "ts",
1747 ColumnDataType::TimestampMillisecond,
1748 SemanticType::Timestamp,
1749 ),
1750 new_column_schema("f1", ColumnDataType::Int64, SemanticType::Field),
1751 ],
1752 rows: vec![Row {
1754 values: vec![i64_value(100), ts_ms_value(1), i64_value(0)],
1755 }],
1756 };
1757 assert_eq!(expect_rows, request.rows);
1758 }
1759
1760 #[test]
1761 fn test_fill_missing_without_default_in_delete() {
1762 let mut builder = builder_with_ts_tag();
1763 builder
1764 .push_column_metadata(ColumnMetadata {
1766 column_schema: datatypes::schema::ColumnSchema::new(
1767 "f0",
1768 ConcreteDataType::int64_datatype(),
1769 true,
1770 ),
1771 semantic_type: SemanticType::Field,
1772 column_id: 3,
1773 })
1774 .push_column_metadata(ColumnMetadata {
1776 column_schema: datatypes::schema::ColumnSchema::new(
1777 "f1",
1778 ConcreteDataType::int64_datatype(),
1779 false,
1780 ),
1781 semantic_type: SemanticType::Field,
1782 column_id: 4,
1783 });
1784 let metadata = builder.build().unwrap();
1785
1786 let rows = Rows {
1787 schema: vec![
1788 new_column_schema("k0", ColumnDataType::Int64, SemanticType::Tag),
1789 new_column_schema(
1790 "ts",
1791 ColumnDataType::TimestampMillisecond,
1792 SemanticType::Timestamp,
1793 ),
1794 ],
1795 rows: vec![Row {
1797 values: vec![i64_value(100), ts_ms_value(1)],
1798 }],
1799 };
1800 let mut request =
1801 WriteRequest::new(RegionId::new(1, 1), OpType::Delete, rows, None).unwrap();
1802 let err = request.check_schema(&metadata).unwrap_err();
1803 assert!(err.is_fill_default());
1804 request.fill_missing_columns(&metadata).unwrap();
1805
1806 let expect_rows = Rows {
1807 schema: vec![
1808 new_column_schema("k0", ColumnDataType::Int64, SemanticType::Tag),
1809 new_column_schema(
1810 "ts",
1811 ColumnDataType::TimestampMillisecond,
1812 SemanticType::Timestamp,
1813 ),
1814 new_column_schema("f1", ColumnDataType::Int64, SemanticType::Field),
1815 ],
1816 rows: vec![Row {
1818 values: vec![i64_value(100), ts_ms_value(1), i64_value(0)],
1819 }],
1820 };
1821 assert_eq!(expect_rows, request.rows);
1822 }
1823
1824 #[test]
1825 fn test_no_default() {
1826 let rows = Rows {
1827 schema: vec![new_column_schema(
1828 "k0",
1829 ColumnDataType::Int64,
1830 SemanticType::Tag,
1831 )],
1832 rows: vec![Row {
1833 values: vec![i64_value(1)],
1834 }],
1835 };
1836 let metadata = new_region_metadata();
1837
1838 let mut request = WriteRequest::new(RegionId::new(1, 1), OpType::Put, rows, None).unwrap();
1839 let err = request.fill_missing_columns(&metadata).unwrap_err();
1840 check_invalid_request(&err, "column ts does not have default value");
1841 }
1842
1843 #[test]
1844 fn test_missing_and_invalid() {
1845 let rows = Rows {
1847 schema: vec![
1848 new_column_schema("k0", ColumnDataType::Int64, SemanticType::Tag),
1849 new_column_schema(
1850 "ts",
1851 ColumnDataType::TimestampMillisecond,
1852 SemanticType::Timestamp,
1853 ),
1854 new_column_schema("f1", ColumnDataType::String, SemanticType::Field),
1855 ],
1856 rows: vec![Row {
1857 values: vec![
1858 i64_value(100),
1859 ts_ms_value(1),
1860 Value {
1861 value_data: Some(ValueData::StringValue("xxxxx".to_string())),
1862 },
1863 ],
1864 }],
1865 };
1866 let metadata = region_metadata_two_fields();
1867
1868 let request = WriteRequest::new(RegionId::new(1, 1), OpType::Put, rows, None).unwrap();
1869 let err = request.check_schema(&metadata).unwrap_err();
1870 check_invalid_request(
1871 &err,
1872 "column f1 expect type Int64(Int64Type), given: STRING(12)",
1873 );
1874 }
1875
1876 #[test]
1877 fn test_write_request_metadata() {
1878 let rows = Rows {
1879 schema: vec![
1880 new_column_schema("c0", ColumnDataType::Int64, SemanticType::Tag),
1881 new_column_schema("c1", ColumnDataType::Int64, SemanticType::Tag),
1882 ],
1883 rows: vec![Row {
1884 values: vec![i64_value(1), i64_value(2)],
1885 }],
1886 };
1887
1888 let metadata = Arc::new(new_region_metadata());
1889 let request = WriteRequest::new(
1890 RegionId::new(1, 1),
1891 OpType::Put,
1892 rows,
1893 Some(metadata.clone()),
1894 )
1895 .unwrap();
1896
1897 assert!(request.region_metadata.is_some());
1898 assert_eq!(
1899 request.region_metadata.unwrap().region_id,
1900 RegionId::new(1, 1)
1901 );
1902 }
1903}