Skip to main content

mito2/
request.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Worker requests.
16
17use 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/// Request to write a region.
67#[derive(Debug)]
68pub struct WriteRequest {
69    /// Region to write.
70    pub region_id: RegionId,
71    /// Type of the write request.
72    pub op_type: OpType,
73    /// Rows to write.
74    pub rows: Rows,
75    /// Map column name to column index in `rows`.
76    pub name_to_index: HashMap<String, usize>,
77    /// Whether each column has null.
78    pub has_null: Vec<bool>,
79    /// Write hint.
80    pub hint: Option<WriteHint>,
81    /// Region metadata on the time of this request is created.
82    pub(crate) region_metadata: Option<RegionMetadataRef>,
83    /// Partition expression version for the region.
84    pub partition_expr_version: Option<u64>,
85}
86
87impl WriteRequest {
88    /// Creates a new request.
89    ///
90    /// Returns `Err` if `rows` are invalid.
91    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    /// Sets the write hint.
146    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    /// Returns the encoding hint.
157    pub fn primary_key_encoding(&self) -> PrimaryKeyEncoding {
158        infer_primary_key_encoding_from_hint(self.hint.as_ref())
159    }
160
161    /// Returns estimated size of the request.
162    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    /// Gets column index by name.
173    pub fn column_index_by_name(&self, name: &str) -> Option<usize> {
174        self.name_to_index.get(name).copied()
175    }
176
177    /// Checks schema of rows is compatible with schema of the region.
178    ///
179    /// If column with default value is missing, it returns a special [FillDefault](crate::error::Error::FillDefault)
180    /// error.
181    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        // Index all columns in rows.
186        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        // Checks all columns in this region.
195        for column in &metadata.column_metadatas {
196            if let Some(input_col) = rows_columns.remove(&column.column_schema.name) {
197                // Check data type.
198                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                // Check semantic type.
219                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                // Check nullable.
236                // Safety: `rows_columns` ensures this column exists.
237                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                // Rows don't have this column.
250                self.check_missing_column(column)?;
251
252                need_fill_default = true;
253            }
254        }
255
256        // Checks all columns in rows exist in the region.
257        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        // If we need to fill default values, return a special error.
267        ensure!(!need_fill_default, FillDefaultSnafu { region_id });
268
269        Ok(())
270    }
271
272    /// Tries to fill missing columns.
273    ///
274    /// Currently, our protobuf format might be inefficient when we need to fill lots of null
275    /// values.
276    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    /// Checks the schema and fill missing columns.
291    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                // TODO(yingwen): Add metrics for this case.
295                // We need to fill default value. The write request may be a request
296                // sent before changing the schema.
297                self.fill_missing_columns(metadata)?;
298            } else {
299                return Err(e);
300            }
301        }
302
303        Ok(())
304    }
305
306    /// Fills default value for specific `columns`.
307    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    /// Checks whether we should allow a row doesn't provide this column.
345    fn check_missing_column(&self, column: &ColumnMetadata) -> Result<()> {
346        if self.op_type == OpType::Delete {
347            if column.semantic_type == SemanticType::Field {
348                // For delete request, all tags and timestamp is required. We don't fill default
349                // tag or timestamp while deleting rows.
350                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        // Not a delete request. Checks whether they have default value.
361        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    /// Returns the default value for specific column.
374    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                // For delete request, we need a default value for padding so we
389                // can delete a row even a field doesn't have a default value. So the
390                // value doesn't need to following the default value constraint of the
391                // column.
392                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                // For put requests, we use the default value from column schema.
400                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                    // This column doesn't have default value.
419                    .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        // Convert default value into proto's value.
430        Ok(api::helper::to_grpc_value(default_value))
431    }
432}
433
434/// Validate proto value schema.
435pub(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/// Oneshot output result sender.
476#[derive(Debug)]
477pub struct OutputTx(Sender<Result<AffectedRows>>);
478
479impl OutputTx {
480    /// Creates a new output sender.
481    pub(crate) fn new(sender: Sender<Result<AffectedRows>>) -> OutputTx {
482        OutputTx(sender)
483    }
484
485    /// Sends the `result`.
486    pub(crate) fn send(self, result: Result<AffectedRows>) {
487        // Ignores send result.
488        let _ = self.0.send(result);
489    }
490}
491
492/// Optional output result sender.
493#[derive(Debug)]
494pub(crate) struct OptionOutputTx(Option<OutputTx>);
495
496impl OptionOutputTx {
497    /// Creates a sender.
498    pub(crate) fn new(sender: Option<OutputTx>) -> OptionOutputTx {
499        OptionOutputTx(sender)
500    }
501
502    /// Creates an empty sender.
503    pub(crate) fn none() -> OptionOutputTx {
504        OptionOutputTx(None)
505    }
506
507    /// Sends the `result` and consumes the inner sender.
508    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    /// Sends the `result` and consumes the sender.
515    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    /// Takes the inner sender.
522    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
539/// Callback on failure.
540pub(crate) trait OnFailure {
541    /// Handles `err` on failure.
542    fn on_failure(&mut self, err: Error);
543}
544
545/// Sender and write request.
546#[derive(Debug)]
547pub(crate) struct SenderWriteRequest {
548    /// Result sender.
549    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/// Request sent to a worker with timestamp
569#[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/// Request sent to a worker
585#[derive(Debug)]
586pub(crate) enum WorkerRequest {
587    /// Write to a region.
588    Write(SenderWriteRequest),
589
590    /// Ddl request to a region.
591    Ddl(SenderDdlRequest),
592
593    /// Notifications from internal background jobs.
594    Background {
595        /// Id of the region to send.
596        region_id: RegionId,
597        /// Internal notification.
598        notify: BackgroundNotify,
599    },
600
601    /// The internal commands.
602    SetRegionRoleStateGracefully {
603        /// Id of the region to send.
604        region_id: RegionId,
605        /// The [SettableRegionRoleState].
606        region_role_state: SettableRegionRoleState,
607        /// The sender of [SetReadonlyResponse].
608        sender: Sender<SetRegionRoleStateResponse>,
609    },
610
611    /// Notify a worker to stop.
612    Stop,
613
614    /// Use [RegionEdit] to edit a region directly.
615    EditRegion(RegionEditRequest),
616
617    /// Keep the manifest of a region up to date.
618    SyncRegion(RegionSyncRequest),
619
620    /// Bulk inserts request and region metadata.
621    BulkInserts(BulkInsertRequest),
622
623    /// Remap manifests request.
624    RemapManifests(RemapManifestsRequest),
625
626    /// Copy region from request.
627    CopyRegionFrom(CopyRegionFromRequest),
628}
629
630impl WorkerRequest {
631    /// Creates a new open region request.
632    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    /// Creates a new catchup region request.
649    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    /// Converts request from a [RegionRequest].
664    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) = &region_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) = &region_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    /// Converts [RemapManifestsRequest] from a [RemapManifestsRequest](store_api::region_engine::RemapManifestsRequest).
810    ///
811    /// # Errors
812    ///
813    /// Returns an error if the partition expression is invalid or missing.
814    /// Returns an error if the new partition expressions are not found for some regions.
815    #[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    /// Converts [CopyRegionFromRequest] from a [MitoCopyRegionFromRequest](store_api::region_engine::MitoCopyRegionFromRequest).
849    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/// DDL request to a region.
868#[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/// Sender and Ddl request.
886#[derive(Debug)]
887pub(crate) struct SenderDdlRequest {
888    /// Region id of the request.
889    pub(crate) region_id: RegionId,
890    /// Result sender.
891    pub(crate) sender: OptionOutputTx,
892    /// Ddl request.
893    pub(crate) request: DdlRequest,
894}
895
896/// Notification from a background job.
897#[derive(Debug)]
898pub(crate) enum BackgroundNotify {
899    /// Compaction planning has finished.
900    CompactionPickFinished(CompactionPickFinished),
901    /// Flush has finished.
902    FlushFinished(FlushFinished),
903    /// Flush has failed.
904    FlushFailed(FlushFailed),
905    /// Index build has finished.
906    IndexBuildFinished(IndexBuildFinished),
907    /// Index build has been stopped (aborted or succeeded).
908    IndexBuildStopped(IndexBuildStopped),
909    /// Index build has failed.
910    IndexBuildFailed(IndexBuildFailed),
911    /// An index build must be retried against the latest schema generation.
912    IndexBuildRetry(BuildIndexRequest),
913    /// Compaction has finished.
914    CompactionFinished(CompactionFinished),
915    /// Compaction has been cancelled cooperatively.
916    CompactionCancelled(CompactionCancelled),
917    /// Compaction has failed.
918    CompactionFailed(CompactionFailed),
919    /// Truncate result.
920    Truncate(TruncateResult),
921    /// Discard unflushed data result.
922    DiscardUnflushed(DiscardUnflushedResult),
923    /// Region change result.
924    RegionChange(RegionChangeResult),
925    /// Region edit result.
926    RegionEdit(RegionEditResult),
927    /// Enter staging result.
928    EnterStaging(EnterStagingResult),
929    /// Copy region result.
930    CopyRegionFromFinished(CopyRegionFromFinished),
931}
932
933/// Notifies a flush job is finished.
934#[derive(Debug)]
935pub(crate) struct FlushFinished {
936    /// Region id.
937    pub(crate) region_id: RegionId,
938    /// Reason to flush.
939    pub(crate) flush_reason: FlushReason,
940    /// Entry id of flushed data.
941    pub(crate) flushed_entry_id: EntryId,
942    /// Flush result senders.
943    pub(crate) senders: Vec<OutputTx>,
944    /// Flush timer.
945    pub(crate) _timer: HistogramTimer,
946    /// Region edit to apply.
947    pub(crate) edit: RegionEdit,
948    /// Memtables to remove.
949    pub(crate) memtables_to_remove: SmallVec<[MemtableId; 2]>,
950    /// Whether the region is in staging mode.
951    pub(crate) is_staging: bool,
952}
953
954impl FlushFinished {
955    /// Marks the flush job as successful and observes the timer.
956    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/// Notifies a flush job is failed.
975#[derive(Debug)]
976pub(crate) struct FlushFailed {
977    /// The error source of the failure.
978    pub(crate) err: Arc<Error>,
979}
980
981impl FlushFailed {
982    /// Returns whether the flush was cancelled cooperatively.
983    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/// Notifies an index build job has been stopped.
995#[derive(Debug)]
996pub(crate) struct IndexBuildStopped {
997    pub(crate) file_id: FileId,
998}
999
1000/// Notifies an index build job has failed.
1001#[derive(Debug)]
1002pub(crate) struct IndexBuildFailed {
1003    pub(crate) err: Arc<Error>,
1004}
1005
1006/// Notifies a compaction job has finished.
1007#[derive(Debug)]
1008pub(crate) struct CompactionFinished {
1009    /// Region id.
1010    pub(crate) region_id: RegionId,
1011    /// Identity and reservation lease of the accepted execution.
1012    pub(crate) execution: CompactionExecution,
1013    /// Compaction result senders.
1014    pub(crate) senders: Vec<OutputTx>,
1015    /// Start time of compaction task.
1016    pub(crate) start_time: Instant,
1017    /// Region edit to apply.
1018    pub(crate) edit: RegionEdit,
1019}
1020
1021/// Notifies a compaction job has been cancelled cooperatively.
1022#[derive(Debug)]
1023pub(crate) struct CompactionCancelled {
1024    /// Region id.
1025    pub(crate) region_id: RegionId,
1026    /// Identity and reservation lease of the accepted execution.
1027    pub(crate) execution: CompactionExecution,
1028    /// Waiters to wake once the cancellation has been observed by the worker.
1029    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        // only update compaction time on success
1044        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    /// Compaction succeeded but failed to update manifest or region's already been dropped.
1055    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/// A failing compaction result.
1066#[derive(Debug)]
1067pub(crate) struct CompactionFailed {
1068    pub(crate) region_id: RegionId,
1069    /// Identity and reservation lease of the accepted execution.
1070    pub(crate) execution: CompactionExecution,
1071    /// The error source of the failure.
1072    pub(crate) err: Arc<Error>,
1073}
1074
1075/// Notifies the truncate result of a region.
1076#[derive(Debug)]
1077pub(crate) struct TruncateResult {
1078    /// Region id.
1079    pub(crate) region_id: RegionId,
1080    /// Result sender.
1081    pub(crate) sender: OptionOutputTx,
1082    /// Truncate result.
1083    pub(crate) result: Result<()>,
1084    pub(crate) kind: TruncateKind,
1085}
1086
1087/// Notifies the result of discarding unflushed data from a region.
1088#[derive(Debug)]
1089pub(crate) struct DiscardUnflushedResult {
1090    /// Region id.
1091    pub(crate) region_id: RegionId,
1092    /// Result sender.
1093    pub(crate) sender: OptionOutputTx,
1094    /// Manifest update result.
1095    pub(crate) result: Result<()>,
1096    /// Last WAL entry covered by the discard operation.
1097    pub(crate) discarded_entry_id: EntryId,
1098    /// Last sequence covered by the discard operation.
1099    pub(crate) discarded_sequence: SequenceNumber,
1100    /// Estimated number of discarded rows.
1101    pub(crate) discarded_rows: u64,
1102    /// Estimated number of discarded bytes.
1103    pub(crate) discarded_bytes: u64,
1104}
1105
1106/// Notifies the region the result of writing region change action.
1107#[derive(Debug)]
1108pub(crate) struct RegionChangeResult {
1109    /// Region id.
1110    pub(crate) region_id: RegionId,
1111    /// The new region metadata to apply.
1112    pub(crate) new_meta: RegionMetadataRef,
1113    /// Result sender.
1114    pub(crate) sender: OptionOutputTx,
1115    /// Result from the manifest manager.
1116    pub(crate) result: Result<()>,
1117    /// Used for index build in schema change.
1118    pub(crate) need_index: bool,
1119    /// New options for the region.
1120    pub(crate) new_options: Option<RegionOptions>,
1121}
1122
1123/// Notifies the region the result of entering staging.
1124#[derive(Debug)]
1125pub(crate) struct EnterStagingResult {
1126    /// Region id.
1127    pub(crate) region_id: RegionId,
1128    /// The new staging partition directive to apply.
1129    pub(crate) partition_directive: StagingPartitionDirective,
1130    /// Result sender.
1131    pub(crate) sender: OptionOutputTx,
1132    /// Result from the manifest manager.
1133    pub(crate) result: Result<()>,
1134}
1135
1136#[derive(Debug)]
1137pub(crate) struct CopyRegionFromFinished {
1138    /// Region id.
1139    pub(crate) region_id: RegionId,
1140    /// Region edit to apply.
1141    pub(crate) edit: RegionEdit,
1142    /// Result sender.
1143    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/// Request to edit a region directly.
1168#[derive(Debug)]
1169pub(crate) struct RegionEditRequest {
1170    pub(crate) region_id: RegionId,
1171    pub(crate) edit: RegionEdit,
1172    /// Whether to preload SST files into the write cache.
1173    pub(crate) preload_sst_cache: bool,
1174    /// The waiters that are waiting for this region edit's result.
1175    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/// Notifies the regin the result of editing region.
1195#[derive(Debug)]
1196pub(crate) struct RegionEditResult {
1197    /// Region id.
1198    pub(crate) region_id: RegionId,
1199    /// Result waiters.
1200    pub(crate) waiters: Waiters,
1201    /// Region edit to apply.
1202    pub(crate) edit: RegionEdit,
1203    /// Result from the manifest manager.
1204    pub(crate) result: std::result::Result<(), Arc<Error>>,
1205    /// Whether region state need to be set to Writable after handling this request.
1206    pub(crate) update_region_state: bool,
1207    /// The region is in staging mode before handling this request.
1208    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    /// files need to build index, empty means all.
1216    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    /// Returns the latest manifest version and a boolean indicating whether new maniefst is installed.
1224    pub(crate) sender: Sender<Result<(ManifestVersion, bool)>>,
1225}
1226
1227#[derive(Debug)]
1228pub(crate) struct RemapManifestsRequest {
1229    /// The [`RegionId`] of a staging region used to obtain table directory and storage configuration for the remap operation.
1230    pub(crate) region_id: RegionId,
1231    /// Regions to remap manifests from.
1232    pub(crate) input_regions: Vec<RegionId>,
1233    /// For each old region, which new regions should receive its files
1234    pub(crate) region_mapping: HashMap<RegionId, Vec<RegionId>>,
1235    /// New partition expressions for the new regions.
1236    pub(crate) new_partition_exprs: HashMap<RegionId, PartitionExpr>,
1237    /// Sender for the result of the remap operation.
1238    ///
1239    /// The result is a map from region IDs to their corresponding staging manifest paths.
1240    pub(crate) sender: Sender<Result<HashMap<RegionId, String>>>,
1241}
1242
1243#[derive(Debug)]
1244pub(crate) struct CopyRegionFromRequest {
1245    /// The [`RegionId`] of the target region.
1246    pub(crate) region_id: RegionId,
1247    /// The [`RegionId`] of the source region.
1248    pub(crate) source_region_id: RegionId,
1249    /// The parallelism of the copy operation.
1250    pub(crate) parallelism: usize,
1251    /// Result sender.
1252    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            // Column is not nullable.
1686            .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            // Column f1 is not nullable and we use 0 for padding.
1753            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            // f0 is nullable.
1765            .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            // f1 is not nullable and don't has default.
1775            .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            // Missing f0 (nullable), f1 (not nullable).
1796            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            // Column f1 is not nullable and we use 0 for padding.
1817            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        // Missing f0 and f1 has invalid type (string).
1846        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}