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