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    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/// Request to write a region.
68#[derive(Debug)]
69pub struct WriteRequest {
70    /// Region to write.
71    pub region_id: RegionId,
72    /// Type of the write request.
73    pub op_type: OpType,
74    /// Rows to write.
75    pub rows: Rows,
76    /// Map column name to column index in `rows`.
77    pub name_to_index: HashMap<String, usize>,
78    /// Whether each column has null.
79    pub has_null: Vec<bool>,
80    /// Whether this insert should skip WAL. Never applies to deletes.
81    pub skip_wal: bool,
82    /// Write hint.
83    pub hint: Option<WriteHint>,
84    /// Region metadata on the time of this request is created.
85    pub(crate) region_metadata: Option<RegionMetadataRef>,
86    /// Partition expression version for the region.
87    pub partition_expr_version: Option<u64>,
88}
89
90impl WriteRequest {
91    /// Creates a new request.
92    ///
93    /// Returns `Err` if `rows` are invalid.
94    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    /// Sets the request-level WAL policy.
150    pub fn with_skip_wal(mut self, skip_wal: bool) -> Self {
151        self.skip_wal = skip_wal;
152        self
153    }
154
155    /// Sets the write hint.
156    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    /// Returns the encoding hint.
167    pub fn primary_key_encoding(&self) -> PrimaryKeyEncoding {
168        infer_primary_key_encoding_from_hint(self.hint.as_ref())
169    }
170
171    /// Returns estimated size of the request.
172    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    /// Gets column index by name.
183    pub fn column_index_by_name(&self, name: &str) -> Option<usize> {
184        self.name_to_index.get(name).copied()
185    }
186
187    /// Checks schema of rows is compatible with schema of the region.
188    ///
189    /// If column with default value is missing, it returns a special [FillDefault](crate::error::Error::FillDefault)
190    /// error.
191    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        // Index all columns in rows.
196        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        // Checks all columns in this region.
205        for column in &metadata.column_metadatas {
206            if let Some(input_col) = rows_columns.remove(&column.column_schema.name) {
207                // Check data type.
208                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                // Check semantic type.
229                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                // Check nullable.
246                // Safety: `rows_columns` ensures this column exists.
247                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                // Rows don't have this column.
260                self.check_missing_column(column)?;
261
262                need_fill_default = true;
263            }
264        }
265
266        // Checks all columns in rows exist in the region.
267        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        // If we need to fill default values, return a special error.
277        ensure!(!need_fill_default, FillDefaultSnafu { region_id });
278
279        Ok(())
280    }
281
282    /// Tries to fill missing columns.
283    ///
284    /// Currently, our protobuf format might be inefficient when we need to fill lots of null
285    /// values.
286    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    /// Checks the schema and fill missing columns.
301    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                // TODO(yingwen): Add metrics for this case.
305                // We need to fill default value. The write request may be a request
306                // sent before changing the schema.
307                self.fill_missing_columns(metadata)?;
308            } else {
309                return Err(e);
310            }
311        }
312
313        Ok(())
314    }
315
316    /// Fills default value for specific `columns`.
317    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    /// Checks whether we should allow a row doesn't provide this column.
355    fn check_missing_column(&self, column: &ColumnMetadata) -> Result<()> {
356        if self.op_type == OpType::Delete {
357            if column.semantic_type == SemanticType::Field {
358                // For delete request, all tags and timestamp is required. We don't fill default
359                // tag or timestamp while deleting rows.
360                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        // Not a delete request. Checks whether they have default value.
371        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    /// Returns the default value for specific column.
384    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                // For delete request, we need a default value for padding so we
399                // can delete a row even a field doesn't have a default value. So the
400                // value doesn't need to following the default value constraint of the
401                // column.
402                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                // For put requests, we use the default value from column schema.
410                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                    // This column doesn't have default value.
429                    .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        // Convert default value into proto's value.
440        Ok(api::helper::to_grpc_value(default_value))
441    }
442}
443
444/// Validate proto value schema.
445pub(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/// Oneshot output result sender.
477#[derive(Debug)]
478pub struct OutputTx(Sender<Result<AffectedRows>>);
479
480impl OutputTx {
481    /// Creates a new output sender.
482    pub(crate) fn new(sender: Sender<Result<AffectedRows>>) -> OutputTx {
483        OutputTx(sender)
484    }
485
486    /// Sends the `result`.
487    pub(crate) fn send(self, result: Result<AffectedRows>) {
488        // Ignores send result.
489        let _ = self.0.send(result);
490    }
491}
492
493/// Optional output result sender.
494#[derive(Debug)]
495pub(crate) struct OptionOutputTx(Option<OutputTx>);
496
497impl OptionOutputTx {
498    /// Creates a sender.
499    pub(crate) fn new(sender: Option<OutputTx>) -> OptionOutputTx {
500        OptionOutputTx(sender)
501    }
502
503    /// Creates an empty sender.
504    pub(crate) fn none() -> OptionOutputTx {
505        OptionOutputTx(None)
506    }
507
508    /// Sends the `result` and consumes the inner sender.
509    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    /// Sends the `result` and consumes the sender.
516    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    /// Takes the inner sender.
523    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
540/// Callback on failure.
541pub(crate) trait OnFailure {
542    /// Handles `err` on failure.
543    fn on_failure(&mut self, err: Error);
544}
545
546/// Sender and write request.
547#[derive(Debug)]
548pub(crate) struct SenderWriteRequest {
549    /// Result sender.
550    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/// Request sent to a worker with timestamp
571#[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/// Request sent to a worker
587#[derive(Debug)]
588pub(crate) enum WorkerRequest {
589    /// Write to a region.
590    Write(SenderWriteRequest),
591
592    /// Ddl request to a region.
593    Ddl(SenderDdlRequest),
594
595    /// Notifications from internal background jobs.
596    Background {
597        /// Id of the region to send.
598        region_id: RegionId,
599        /// Internal notification.
600        notify: BackgroundNotify,
601    },
602
603    /// The internal commands.
604    SetRegionRoleStateGracefully {
605        /// Id of the region to send.
606        region_id: RegionId,
607        /// The [SettableRegionRoleState].
608        region_role_state: SettableRegionRoleState,
609        /// The sender of [SetReadonlyResponse].
610        sender: Sender<SetRegionRoleStateResponse>,
611    },
612
613    /// Notify a worker to stop.
614    Stop,
615
616    /// Use [RegionEdit] to edit a region directly.
617    EditRegion(RegionEditRequest),
618
619    /// Keep the manifest of a region up to date.
620    SyncRegion(RegionSyncRequest),
621
622    /// Bulk inserts request and region metadata.
623    BulkInserts(BulkInsertRequest),
624
625    /// Remap manifests request.
626    RemapManifests(RemapManifestsRequest),
627
628    /// Copy region from request.
629    CopyRegionFrom(CopyRegionFromRequest),
630}
631
632impl WorkerRequest {
633    /// Creates a new open region request.
634    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    /// Creates a new catchup region request.
651    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    /// Converts request from a [RegionRequest].
666    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) = &region_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) = &region_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    /// Converts [RemapManifestsRequest] from a [RemapManifestsRequest](store_api::region_engine::RemapManifestsRequest).
813    ///
814    /// # Errors
815    ///
816    /// Returns an error if the partition expression is invalid or missing.
817    /// Returns an error if the new partition expressions are not found for some regions.
818    #[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    /// Converts [CopyRegionFromRequest] from a [MitoCopyRegionFromRequest](store_api::region_engine::MitoCopyRegionFromRequest).
852    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/// DDL request to a region.
871#[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/// Sender and Ddl request.
889#[derive(Debug)]
890pub(crate) struct SenderDdlRequest {
891    /// Region id of the request.
892    pub(crate) region_id: RegionId,
893    /// Result sender.
894    pub(crate) sender: OptionOutputTx,
895    /// Ddl request.
896    pub(crate) request: DdlRequest,
897}
898
899/// Notification from a background job.
900#[derive(Debug)]
901pub(crate) enum BackgroundNotify {
902    /// Compaction planning has finished.
903    CompactionPickFinished(CompactionPickFinished),
904    /// Flush has finished.
905    FlushFinished(FlushFinished),
906    /// Flush has failed.
907    FlushFailed(FlushFailed),
908    /// Index build has finished.
909    IndexBuildFinished(IndexBuildFinished),
910    /// Index build has been stopped (aborted or succeeded).
911    IndexBuildStopped(IndexBuildStopped),
912    /// Index build has failed.
913    IndexBuildFailed(IndexBuildFailed),
914    /// An index build must be retried against the latest schema generation.
915    IndexBuildRetry(BuildIndexRequest),
916    /// Compaction has finished.
917    CompactionFinished(CompactionFinished),
918    /// Compaction has been cancelled cooperatively.
919    CompactionCancelled(CompactionCancelled),
920    /// Compaction has failed.
921    CompactionFailed(CompactionFailed),
922    /// Truncate result.
923    Truncate(TruncateResult),
924    /// Discard unflushed data result.
925    DiscardUnflushed(DiscardUnflushedResult),
926    /// Region change result.
927    RegionChange(RegionChangeResult),
928    /// Region edit result.
929    RegionEdit(RegionEditResult),
930    /// Enter staging result.
931    EnterStaging(EnterStagingResult),
932    /// Copy region result.
933    CopyRegionFromFinished(CopyRegionFromFinished),
934}
935
936/// Notifies a flush job is finished.
937#[derive(Debug)]
938pub(crate) struct FlushFinished {
939    /// Region id.
940    pub(crate) region_id: RegionId,
941    /// Reason to flush.
942    pub(crate) flush_reason: FlushReason,
943    /// Entry id of flushed data.
944    pub(crate) flushed_entry_id: EntryId,
945    /// Flush result senders.
946    pub(crate) senders: Vec<OutputTx>,
947    /// Flush timer.
948    pub(crate) _timer: HistogramTimer,
949    /// Region edit to apply.
950    pub(crate) edit: RegionEdit,
951    /// Memtables to remove.
952    pub(crate) memtables_to_remove: SmallVec<[MemtableId; 2]>,
953    /// Whether the region is in staging mode.
954    pub(crate) is_staging: bool,
955}
956
957impl FlushFinished {
958    /// Marks the flush job as successful and observes the timer.
959    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/// Notifies a flush job is failed.
978#[derive(Debug)]
979pub(crate) struct FlushFailed {
980    /// The error source of the failure.
981    pub(crate) err: Arc<Error>,
982}
983
984impl FlushFailed {
985    /// Returns whether the flush was cancelled cooperatively.
986    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/// Notifies an index build job has been stopped.
998#[derive(Debug)]
999pub(crate) struct IndexBuildStopped {
1000    pub(crate) file_id: FileId,
1001}
1002
1003/// Notifies an index build job has failed.
1004#[derive(Debug)]
1005pub(crate) struct IndexBuildFailed {
1006    pub(crate) err: Arc<Error>,
1007}
1008
1009/// Notifies a compaction job has finished.
1010#[derive(Debug)]
1011pub(crate) struct CompactionFinished {
1012    /// Region id.
1013    pub(crate) region_id: RegionId,
1014    /// Identity and reservation lease of the accepted execution.
1015    pub(crate) execution: CompactionExecution,
1016    /// Compaction result senders.
1017    pub(crate) senders: Vec<OutputTx>,
1018    /// Start time of compaction task.
1019    pub(crate) start_time: Instant,
1020    /// Region edit to apply.
1021    pub(crate) edit: RegionEdit,
1022}
1023
1024/// Notifies a compaction job has been cancelled cooperatively.
1025#[derive(Debug)]
1026pub(crate) struct CompactionCancelled {
1027    /// Region id.
1028    pub(crate) region_id: RegionId,
1029    /// Identity and reservation lease of the accepted execution.
1030    pub(crate) execution: CompactionExecution,
1031    /// Waiters to wake once the cancellation has been observed by the worker.
1032    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        // only update compaction time on success
1047        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    /// Compaction succeeded but failed to update manifest or region's already been dropped.
1058    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/// A failing compaction result.
1069#[derive(Debug)]
1070pub(crate) struct CompactionFailed {
1071    pub(crate) region_id: RegionId,
1072    /// Identity and reservation lease of the accepted execution.
1073    pub(crate) execution: CompactionExecution,
1074    /// The error source of the failure.
1075    pub(crate) err: Arc<Error>,
1076}
1077
1078/// Notifies the truncate result of a region.
1079#[derive(Debug)]
1080pub(crate) struct TruncateResult {
1081    /// Region id.
1082    pub(crate) region_id: RegionId,
1083    /// Result sender.
1084    pub(crate) sender: OptionOutputTx,
1085    /// Truncate result.
1086    pub(crate) result: Result<()>,
1087    pub(crate) kind: TruncateKind,
1088}
1089
1090/// Notifies the result of discarding unflushed data from a region.
1091#[derive(Debug)]
1092pub(crate) struct DiscardUnflushedResult {
1093    /// Region id.
1094    pub(crate) region_id: RegionId,
1095    /// Result sender.
1096    pub(crate) sender: OptionOutputTx,
1097    /// Manifest update result.
1098    pub(crate) result: Result<()>,
1099    /// Last WAL entry covered by the discard operation.
1100    pub(crate) discarded_entry_id: EntryId,
1101    /// Last sequence covered by the discard operation.
1102    pub(crate) discarded_sequence: SequenceNumber,
1103    /// Estimated number of discarded rows.
1104    pub(crate) discarded_rows: u64,
1105    /// Estimated number of discarded bytes.
1106    pub(crate) discarded_bytes: u64,
1107}
1108
1109/// Notifies the region the result of writing region change action.
1110#[derive(Debug)]
1111pub(crate) struct RegionChangeResult {
1112    /// Region id.
1113    pub(crate) region_id: RegionId,
1114    /// The new region metadata to apply.
1115    pub(crate) new_meta: RegionMetadataRef,
1116    /// Result sender.
1117    pub(crate) sender: OptionOutputTx,
1118    /// Result from the manifest manager.
1119    pub(crate) result: Result<()>,
1120    /// Used for index build in schema change.
1121    pub(crate) need_index: bool,
1122    /// New options for the region.
1123    pub(crate) new_options: Option<RegionOptions>,
1124}
1125
1126/// Notifies the region the result of entering staging.
1127#[derive(Debug)]
1128pub(crate) struct EnterStagingResult {
1129    /// Region id.
1130    pub(crate) region_id: RegionId,
1131    /// The new staging partition directive to apply.
1132    pub(crate) partition_directive: StagingPartitionDirective,
1133    /// Result sender.
1134    pub(crate) sender: OptionOutputTx,
1135    /// Result from the manifest manager.
1136    pub(crate) result: Result<()>,
1137}
1138
1139#[derive(Debug)]
1140pub(crate) struct CopyRegionFromFinished {
1141    /// Region id.
1142    pub(crate) region_id: RegionId,
1143    /// Region edit to apply.
1144    pub(crate) edit: RegionEdit,
1145    /// Result sender.
1146    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/// Request to edit a region directly.
1171#[derive(Debug)]
1172pub(crate) struct RegionEditRequest {
1173    pub(crate) region_id: RegionId,
1174    pub(crate) edit: RegionEdit,
1175    /// Whether to preload SST files into the write cache.
1176    pub(crate) preload_sst_cache: bool,
1177    /// The waiters that are waiting for this region edit's result.
1178    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/// Notifies the regin the result of editing region.
1198#[derive(Debug)]
1199pub(crate) struct RegionEditResult {
1200    /// Region id.
1201    pub(crate) region_id: RegionId,
1202    /// Result waiters.
1203    pub(crate) waiters: Waiters,
1204    /// Region edit to apply.
1205    pub(crate) edit: RegionEdit,
1206    /// Result from the manifest manager.
1207    pub(crate) result: std::result::Result<(), Arc<Error>>,
1208    /// Whether region state need to be set to Writable after handling this request.
1209    pub(crate) update_region_state: bool,
1210    /// The region is in staging mode before handling this request.
1211    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    /// files need to build index, empty means all.
1219    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    /// Returns the latest manifest version and a boolean indicating whether new maniefst is installed.
1227    pub(crate) sender: Sender<Result<(ManifestVersion, bool)>>,
1228}
1229
1230#[derive(Debug)]
1231pub(crate) struct RemapManifestsRequest {
1232    /// The [`RegionId`] of a staging region used to obtain table directory and storage configuration for the remap operation.
1233    pub(crate) region_id: RegionId,
1234    /// Regions to remap manifests from.
1235    pub(crate) input_regions: Vec<RegionId>,
1236    /// For each old region, which new regions should receive its files
1237    pub(crate) region_mapping: HashMap<RegionId, Vec<RegionId>>,
1238    /// New partition expressions for the new regions.
1239    pub(crate) new_partition_exprs: HashMap<RegionId, PartitionExpr>,
1240    /// Sender for the result of the remap operation.
1241    ///
1242    /// The result is a map from region IDs to their corresponding staging manifest paths.
1243    pub(crate) sender: Sender<Result<HashMap<RegionId, String>>>,
1244}
1245
1246#[derive(Debug)]
1247pub(crate) struct CopyRegionFromRequest {
1248    /// The [`RegionId`] of the target region.
1249    pub(crate) region_id: RegionId,
1250    /// The [`RegionId`] of the source region.
1251    pub(crate) source_region_id: RegionId,
1252    /// The parallelism of the copy operation.
1253    pub(crate) parallelism: usize,
1254    /// Result sender.
1255    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            // Column is not nullable.
1689            .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            // Column f1 is not nullable and we use 0 for padding.
1756            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            // f0 is nullable.
1768            .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            // f1 is not nullable and don't has default.
1778            .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            // Missing f0 (nullable), f1 (not nullable).
1799            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            // Column f1 is not nullable and we use 0 for padding.
1820            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        // Missing f0 and f1 has invalid type (string).
1849        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}