Skip to main content

metric_engine/engine/
put.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
15use std::collections::HashMap;
16
17use api::helper::ColumnDataTypeWrapper;
18use api::v1::{
19    ColumnSchema, PrimaryKeyEncoding as PrimaryKeyEncodingProto, Row, Rows, SemanticType, Value,
20    WriteHint,
21};
22use common_telemetry::{error, info};
23use fxhash::FxHashMap;
24use snafu::{OptionExt, ResultExt, ensure};
25use store_api::codec::PrimaryKeyEncoding;
26use store_api::metadata::ColumnMetadata;
27use store_api::region_request::{
28    AffectedRows, RegionDeleteRequest, RegionPutRequest, RegionRequest,
29};
30use store_api::storage::{RegionId, TableId};
31
32use crate::engine::MetricEngineInner;
33use crate::error::{
34    ColumnNotFoundSnafu, CreateDefaultSnafu, ForbiddenPhysicalWriteSnafu, InvalidRequestSnafu,
35    LogicalRegionNotFoundSnafu, PhysicalRegionNotFoundSnafu, Result, UnexpectedRequestSnafu,
36    UnsupportedRegionRequestSnafu,
37};
38use crate::metrics::{FORBIDDEN_OPERATION_COUNT, MITO_OPERATION_ELAPSED};
39use crate::row_modifier::{RowsIter, TableIdInput};
40use crate::utils::to_data_region_id;
41
42impl MetricEngineInner {
43    /// Dispatch region put request
44    pub async fn put_region(
45        &self,
46        region_id: RegionId,
47        request: RegionPutRequest,
48    ) -> Result<AffectedRows> {
49        let is_putting_physical_region =
50            self.state.read().unwrap().exist_physical_region(region_id);
51
52        if is_putting_physical_region {
53            info!(
54                "Metric region received put request {request:?} on physical region {region_id:?}"
55            );
56            FORBIDDEN_OPERATION_COUNT.inc();
57
58            ForbiddenPhysicalWriteSnafu.fail()
59        } else {
60            self.put_logical_region(region_id, request).await
61        }
62    }
63
64    /// Batch write multiple logical regions to the same physical region.
65    ///
66    /// Dispatch region put requests in batch.
67    ///
68    /// Requests may span multiple physical regions. We group them by physical
69    /// region and write sequentially. This method fails fast on validation or
70    /// preparation errors within a group and stops at the first failure.
71    /// Writes in earlier physical-region groups are not rolled back if a later
72    /// group fails.
73    pub async fn put_regions_batch(
74        &self,
75        requests: impl ExactSizeIterator<Item = (RegionId, RegionPutRequest)>,
76    ) -> Result<AffectedRows> {
77        let len = requests.len();
78
79        if len == 0 {
80            return Ok(0);
81        }
82
83        let _timer = MITO_OPERATION_ELAPSED
84            .with_label_values(&["put_batch"])
85            .start_timer();
86
87        // Fast path: single request, no batching overhead
88        if len == 1 {
89            let (region_id, req) = requests.into_iter().next().unwrap();
90            let is_putting_physical_region =
91                self.state.read().unwrap().exist_physical_region(region_id);
92            if is_putting_physical_region {
93                FORBIDDEN_OPERATION_COUNT.inc();
94                return ForbiddenPhysicalWriteSnafu.fail();
95            }
96
97            return self.put_logical_region(region_id, req).await;
98        }
99
100        let mut requests_per_physical: HashMap<RegionId, Vec<(RegionId, RegionPutRequest)>> =
101            HashMap::new();
102        for (region_id, request) in requests {
103            let is_putting_physical_region =
104                self.state.read().unwrap().exist_physical_region(region_id);
105            if is_putting_physical_region {
106                FORBIDDEN_OPERATION_COUNT.inc();
107                return ForbiddenPhysicalWriteSnafu.fail();
108            }
109            let physical_region_id = self.find_physical_region_id(region_id)?;
110            requests_per_physical
111                .entry(physical_region_id)
112                .or_default()
113                .push((region_id, request));
114        }
115
116        let mut total_affected_rows: AffectedRows = 0;
117        for (physical_region_id, requests) in requests_per_physical {
118            let affected_rows = self
119                .put_regions_batch_single_physical(physical_region_id, requests)
120                .await?;
121            total_affected_rows += affected_rows;
122        }
123
124        Ok(total_affected_rows)
125    }
126
127    /// Write a batch of requests that all belong to the same physical region.
128    ///
129    /// This function orchestrates the batch write process:
130    /// 1. Validates all requests
131    /// 2. Merges requests according to the encoding strategy (sparse or dense)
132    /// 3. Writes the merged batch to the physical region
133    async fn put_regions_batch_single_physical(
134        &self,
135        physical_region_id: RegionId,
136        mut requests: Vec<(RegionId, RegionPutRequest)>,
137    ) -> Result<AffectedRows> {
138        if requests.is_empty() {
139            return Ok(0);
140        }
141
142        let data_region_id = to_data_region_id(physical_region_id);
143        let primary_key_encoding = self.get_primary_key_encoding(data_region_id)?;
144
145        // TODO(weny): Consolidate validation and merging to avoid redundant request traversals,
146        // while ensuring the entire batch is validated before writing.
147        // Validate all requests
148        self.validate_batch_requests(physical_region_id, &mut requests)
149            .await?;
150
151        // Merge requests according to encoding strategy
152        let (merged_request, total_affected_rows) = match primary_key_encoding {
153            PrimaryKeyEncoding::Sparse => self.merge_sparse_batch(physical_region_id, requests)?,
154            PrimaryKeyEncoding::Dense => self.merge_dense_batch(data_region_id, requests)?,
155        };
156
157        // Write once to the physical region
158        self.data_region
159            .write_data(data_region_id, RegionRequest::Put(merged_request))
160            .await?;
161
162        Ok(total_affected_rows)
163    }
164
165    /// Get primary key encoding for a data region.
166    fn get_primary_key_encoding(&self, data_region_id: RegionId) -> Result<PrimaryKeyEncoding> {
167        let state = self.state.read().unwrap();
168        state
169            .get_primary_key_encoding(data_region_id)
170            .context(PhysicalRegionNotFoundSnafu {
171                region_id: data_region_id,
172            })
173    }
174
175    /// Validates all requests in a batch.
176    async fn validate_batch_requests(
177        &self,
178        physical_region_id: RegionId,
179        requests: &mut [(RegionId, RegionPutRequest)],
180    ) -> Result<()> {
181        let skip_wal = requests
182            .first()
183            .is_some_and(|(_, request)| request.skip_wal);
184        ensure!(
185            requests
186                .iter()
187                .all(|(_, request)| request.skip_wal == skip_wal),
188            InvalidRequestSnafu {
189                region_id: physical_region_id,
190                reason: "inconsistent WAL policy in batch"
191            }
192        );
193        for (logical_region_id, request) in requests {
194            self.verify_rows(
195                *logical_region_id,
196                physical_region_id,
197                &mut request.rows,
198                true,
199            )
200            .await?;
201        }
202        Ok(())
203    }
204
205    /// Merges multiple requests using sparse primary key encoding.
206    fn merge_sparse_batch(
207        &self,
208        physical_region_id: RegionId,
209        requests: Vec<(RegionId, RegionPutRequest)>,
210    ) -> Result<(RegionPutRequest, AffectedRows)> {
211        let skip_wal = requests
212            .first()
213            .is_some_and(|(_, request)| request.skip_wal);
214        let total_rows: usize = requests.iter().map(|(_, req)| req.rows.rows.len()).sum();
215        let mut modified_requests = Vec::with_capacity(requests.len());
216        let mut total_affected_rows: AffectedRows = 0;
217        let mut merged_version: Option<u64> = None;
218
219        for (logical_region_id, mut request) in requests {
220            if let Some(request_version) = request.partition_expr_version {
221                if let Some(merged_version) = merged_version {
222                    ensure!(
223                        merged_version == request_version,
224                        InvalidRequestSnafu {
225                            region_id: physical_region_id,
226                            reason: "inconsistent partition expr version in batch"
227                        }
228                    );
229                } else {
230                    merged_version = Some(request_version);
231                }
232            }
233            self.modify_rows(
234                physical_region_id,
235                logical_region_id.table_id(),
236                &mut request.rows,
237                PrimaryKeyEncoding::Sparse,
238            )?;
239
240            let row_count = request.rows.rows.len();
241            total_affected_rows += row_count as AffectedRows;
242            modified_requests.push(request.rows);
243        }
244
245        let schema =
246            Self::build_union_schema(modified_requests.iter().map(|rows| rows.schema.as_slice()));
247        let mut merged_rows = Vec::with_capacity(total_rows);
248        for rows in modified_requests {
249            merged_rows.extend(Self::align_rows_to_schema(rows, &schema));
250        }
251
252        let merged_request = RegionPutRequest {
253            skip_wal,
254            rows: Rows {
255                schema,
256                rows: merged_rows,
257            },
258            hint: Some(WriteHint {
259                primary_key_encoding: PrimaryKeyEncodingProto::Sparse.into(),
260            }),
261            partition_expr_version: merged_version,
262        };
263
264        Ok((merged_request, total_affected_rows))
265    }
266
267    /// Merges multiple requests using dense primary key encoding.
268    ///
269    /// In dense mode, different requests can have different columns.
270    /// We merge all schemas into a union schema, align each row to this schema,
271    /// then batch-modify all rows together (adding __table_id and __tsid).
272    fn merge_dense_batch(
273        &self,
274        data_region_id: RegionId,
275        requests: Vec<(RegionId, RegionPutRequest)>,
276    ) -> Result<(RegionPutRequest, AffectedRows)> {
277        let skip_wal = requests
278            .first()
279            .is_some_and(|(_, request)| request.skip_wal);
280        // Build union schema from all requests
281        let merged_schema =
282            Self::build_union_schema(requests.iter().map(|(_, req)| req.rows.schema.as_slice()));
283
284        // Align all rows to the merged schema and collect table_ids
285        let (merged_rows, table_ids, merged_version) =
286            Self::align_requests_to_schema(requests, &merged_schema)?;
287
288        // Batch-modify all rows (add __table_id and __tsid columns)
289        let final_rows = {
290            let state = self.state.read().unwrap();
291            let physical_columns = state
292                .physical_region_states()
293                .get(&data_region_id)
294                .with_context(|| PhysicalRegionNotFoundSnafu {
295                    region_id: data_region_id,
296                })?
297                .physical_columns();
298
299            let iter = RowsIter::new(
300                Rows {
301                    schema: merged_schema,
302                    rows: merged_rows,
303                },
304                physical_columns,
305            );
306
307            self.row_modifier.modify_rows(
308                iter,
309                TableIdInput::Batch(&table_ids),
310                PrimaryKeyEncoding::Dense,
311            )?
312        };
313
314        let merged_request = RegionPutRequest {
315            skip_wal,
316            rows: final_rows,
317            hint: None,
318            partition_expr_version: merged_version,
319        };
320
321        Ok((merged_request, table_ids.len() as AffectedRows))
322    }
323
324    fn build_union_schema<'a>(
325        schemas: impl IntoIterator<Item = &'a [ColumnSchema]>,
326    ) -> Vec<ColumnSchema> {
327        let mut schema = Vec::new();
328        for columns in schemas {
329            for col in columns {
330                if !schema
331                    .iter()
332                    .any(|existing: &ColumnSchema| existing.column_name == col.column_name)
333                {
334                    schema.push(col.clone());
335                }
336            }
337        }
338        schema
339    }
340
341    fn align_requests_to_schema(
342        requests: Vec<(RegionId, RegionPutRequest)>,
343        merged_schema: &[ColumnSchema],
344    ) -> Result<(Vec<Row>, Vec<TableId>, Option<u64>)> {
345        // Pre-calculate total capacity
346        let total_rows: usize = requests.iter().map(|(_, req)| req.rows.rows.len()).sum();
347        let mut merged_rows = Vec::with_capacity(total_rows);
348        let mut table_ids = Vec::with_capacity(total_rows);
349        let mut merged_version: Option<u64> = None;
350
351        for (logical_region_id, request) in requests {
352            if let Some(request_version) = request.partition_expr_version {
353                if let Some(merged_version) = merged_version {
354                    ensure!(
355                        merged_version == request_version,
356                        InvalidRequestSnafu {
357                            region_id: logical_region_id,
358                            reason: "inconsistent partition expr version in batch"
359                        }
360                    );
361                } else {
362                    merged_version = Some(request_version);
363                }
364            }
365            let table_id = logical_region_id.table_id();
366            let row_count = request.rows.rows.len();
367            merged_rows.extend(Self::align_rows_to_schema(request.rows, merged_schema));
368            table_ids.extend(std::iter::repeat_n(table_id, row_count));
369        }
370
371        Ok((merged_rows, table_ids, merged_version))
372    }
373
374    fn align_rows_to_schema(rows: Rows, merged_schema: &[ColumnSchema]) -> Vec<Row> {
375        let Rows { schema, rows } = rows;
376        if schema.len() == merged_schema.len()
377            && schema
378                .iter()
379                .zip(merged_schema)
380                .all(|(left, right)| left.column_name == right.column_name)
381        {
382            return rows;
383        }
384
385        let col_name_to_idx: FxHashMap<&str, usize> = schema
386            .iter()
387            .enumerate()
388            .map(|(idx, col)| (col.column_name.as_str(), idx))
389            .collect();
390        let col_mapping: Vec<Option<usize>> = merged_schema
391            .iter()
392            .map(|merged_col| {
393                col_name_to_idx
394                    .get(merged_col.column_name.as_str())
395                    .copied()
396            })
397            .collect();
398        let null_value = Value { value_data: None };
399
400        rows.into_iter()
401            .map(|mut row| {
402                let values = col_mapping
403                    .iter()
404                    .map(|opt_idx| match opt_idx {
405                        Some(idx) => std::mem::take(&mut row.values[*idx]),
406                        None => null_value.clone(),
407                    })
408                    .collect();
409                Row { values }
410            })
411            .collect()
412    }
413
414    /// Find the physical region id for a logical region.
415    fn find_physical_region_id(&self, logical_region_id: RegionId) -> Result<RegionId> {
416        let state = self.state.read().unwrap();
417        state
418            .logical_regions()
419            .get(&logical_region_id)
420            .copied()
421            .context(LogicalRegionNotFoundSnafu {
422                region_id: logical_region_id,
423            })
424    }
425
426    /// Dispatch region delete request
427    pub async fn delete_region(
428        &self,
429        region_id: RegionId,
430        request: RegionDeleteRequest,
431    ) -> Result<AffectedRows> {
432        if self.is_physical_region(region_id) {
433            info!(
434                "Metric region received delete request {request:?} on physical region {region_id:?}"
435            );
436            FORBIDDEN_OPERATION_COUNT.inc();
437
438            UnsupportedRegionRequestSnafu {
439                request: RegionRequest::Delete(request),
440            }
441            .fail()
442        } else {
443            self.delete_logical_region(region_id, request).await
444        }
445    }
446
447    async fn put_logical_region(
448        &self,
449        logical_region_id: RegionId,
450        mut request: RegionPutRequest,
451    ) -> Result<AffectedRows> {
452        let _timer = MITO_OPERATION_ELAPSED
453            .with_label_values(&["put"])
454            .start_timer();
455
456        let (physical_region_id, data_region_id, primary_key_encoding) =
457            self.find_data_region_meta(logical_region_id)?;
458
459        self.verify_rows(
460            logical_region_id,
461            physical_region_id,
462            &mut request.rows,
463            true,
464        )
465        .await?;
466
467        // write to data region
468        // TODO: retrieve table name
469        self.modify_rows(
470            physical_region_id,
471            logical_region_id.table_id(),
472            &mut request.rows,
473            primary_key_encoding,
474        )?;
475        if primary_key_encoding == PrimaryKeyEncoding::Sparse {
476            request.hint = Some(WriteHint {
477                primary_key_encoding: PrimaryKeyEncodingProto::Sparse.into(),
478            });
479        }
480        self.data_region
481            .write_data(data_region_id, RegionRequest::Put(request))
482            .await
483    }
484
485    async fn delete_logical_region(
486        &self,
487        logical_region_id: RegionId,
488        mut request: RegionDeleteRequest,
489    ) -> Result<AffectedRows> {
490        let _timer = MITO_OPERATION_ELAPSED
491            .with_label_values(&["delete"])
492            .start_timer();
493
494        let (physical_region_id, data_region_id, primary_key_encoding) =
495            self.find_data_region_meta(logical_region_id)?;
496
497        self.verify_rows(
498            logical_region_id,
499            physical_region_id,
500            &mut request.rows,
501            false,
502        )
503        .await?;
504
505        // write to data region
506        // TODO: retrieve table name
507        self.modify_rows(
508            physical_region_id,
509            logical_region_id.table_id(),
510            &mut request.rows,
511            primary_key_encoding,
512        )?;
513        if primary_key_encoding == PrimaryKeyEncoding::Sparse {
514            request.hint = Some(WriteHint {
515                primary_key_encoding: PrimaryKeyEncodingProto::Sparse.into(),
516            });
517        }
518        self.data_region
519            .write_data(data_region_id, RegionRequest::Delete(request))
520            .await
521    }
522
523    pub(crate) fn find_data_region_meta(
524        &self,
525        logical_region_id: RegionId,
526    ) -> Result<(RegionId, RegionId, PrimaryKeyEncoding)> {
527        let state = self.state.read().unwrap();
528        let physical_region_id = *state
529            .logical_regions()
530            .get(&logical_region_id)
531            .with_context(|| LogicalRegionNotFoundSnafu {
532                region_id: logical_region_id,
533            })?;
534        let data_region_id = to_data_region_id(physical_region_id);
535        let primary_key_encoding = state.get_primary_key_encoding(data_region_id).context(
536            PhysicalRegionNotFoundSnafu {
537                region_id: data_region_id,
538            },
539        )?;
540        Ok((physical_region_id, data_region_id, primary_key_encoding))
541    }
542
543    /// Verifies a request for a logical region against its corresponding metadata region.
544    ///
545    /// Includes:
546    /// - Check if the logical region exists
547    /// - Check if every column in the request exists in the physical region
548    /// - Check each column's datatype and semantic type match the physical region's schema
549    /// - Check the time index column is present
550    /// - When `check_fields` is true, check every logical field column is present.
551    ///   Set this to `false` for delete requests, which legitimately carry only
552    ///   the primary key + timestamp.
553    async fn verify_rows(
554        &self,
555        logical_region_id: RegionId,
556        physical_region_id: RegionId,
557        rows: &mut Rows,
558        check_fields: bool,
559    ) -> Result<()> {
560        // Check if the region exists
561        let data_region_id = to_data_region_id(physical_region_id);
562        let (physical_columns, ts_name) = {
563            let state = self.state.read().unwrap();
564            if !state.is_logical_region_exist(logical_region_id) {
565                error!("Trying to write to an nonexistent region {logical_region_id}");
566                return LogicalRegionNotFoundSnafu {
567                    region_id: logical_region_id,
568                }
569                .fail();
570            }
571
572            let physical_state = state
573                .physical_region_states()
574                .get(&data_region_id)
575                .context(PhysicalRegionNotFoundSnafu {
576                    region_id: data_region_id,
577                })?;
578            (
579                physical_state.physical_columns_snapshot(),
580                physical_state.time_index_column_name().to_string(),
581            )
582        };
583
584        // Type + semantic check on every column in the request schema.
585        for col in &rows.schema {
586            let info = physical_columns
587                .get(&col.column_name)
588                .context(ColumnNotFoundSnafu {
589                    name: &col.column_name,
590                    region_id: logical_region_id,
591                })?;
592
593            ensure!(
594                api::helper::is_column_type_value_eq(
595                    col.datatype,
596                    col.datatype_extension.clone(),
597                    &info.column_schema.data_type
598                ),
599                InvalidRequestSnafu {
600                    region_id: logical_region_id,
601                    reason: format!(
602                        "column {} expect type {:?}, given: {}({})",
603                        col.column_name,
604                        info.column_schema.data_type,
605                        api::v1::ColumnDataType::try_from(col.datatype)
606                            .map(|v| v.as_str_name())
607                            .unwrap_or("Unknown"),
608                        col.datatype,
609                    ),
610                }
611            );
612
613            ensure!(
614                api::helper::is_semantic_type_eq(col.semantic_type, info.semantic_type),
615                InvalidRequestSnafu {
616                    region_id: logical_region_id,
617                    reason: format!(
618                        "column {} expect semantic type {:?}, given: {}({})",
619                        col.column_name,
620                        info.semantic_type,
621                        api::v1::SemanticType::try_from(col.semantic_type)
622                            .map(|v| v.as_str_name())
623                            .unwrap_or("Unknown"),
624                        col.semantic_type,
625                    ),
626                }
627            );
628        }
629
630        ensure!(
631            rows.schema.iter().any(|col| col.column_name == ts_name),
632            InvalidRequestSnafu {
633                region_id: logical_region_id,
634                reason: format!("missing required time index column {ts_name}"),
635            }
636        );
637
638        let logical_columns = self
639            .load_logical_columns(physical_region_id, logical_region_id)
640            .await?;
641        let logical_fields = logical_columns
642            .iter()
643            .filter(|col| col.semantic_type == SemanticType::Field)
644            .map(|col| (col.column_schema.name.as_str(), col))
645            .collect::<HashMap<_, _>>();
646
647        for col in &rows.schema {
648            if api::helper::is_semantic_type_eq(col.semantic_type, SemanticType::Field) {
649                ensure!(
650                    logical_fields.contains_key(col.column_name.as_str()),
651                    InvalidRequestSnafu {
652                        region_id: logical_region_id,
653                        reason: format!(
654                            "field column {} does not belong to logical region {logical_region_id}",
655                            col.column_name,
656                        ),
657                    }
658                );
659            }
660        }
661
662        if check_fields {
663            // Sparse logical writes may omit nullable field columns. Fill them
664            // before the rows are rewritten for the shared physical table.
665            for (field_name, field_meta) in logical_fields {
666                if !rows.schema.iter().any(|col| col.column_name == field_name) {
667                    Self::fill_missing_field_column(
668                        logical_region_id,
669                        field_name,
670                        field_meta,
671                        rows,
672                    )?;
673                }
674            }
675
676            for (field_name, field_meta) in physical_columns
677                .iter()
678                .filter(|(_, col)| col.semantic_type == SemanticType::Field)
679            {
680                if !rows.schema.iter().any(|col| col.column_name == *field_name) {
681                    Self::fill_missing_field_column(
682                        logical_region_id,
683                        field_name,
684                        field_meta,
685                        rows,
686                    )?;
687                }
688            }
689        }
690
691        Ok(())
692    }
693
694    fn fill_missing_field_column(
695        logical_region_id: RegionId,
696        field_name: &str,
697        field_meta: &ColumnMetadata,
698        rows: &mut Rows,
699    ) -> Result<()> {
700        // This is only for schema columns with a concrete default, usually NULL
701        // for field columns from other logical tables sharing this physical table.
702        ensure!(
703            !field_meta.column_schema.is_default_impure(),
704            UnexpectedRequestSnafu {
705                reason: format!(
706                    "unexpected impure default value with region_id: {logical_region_id}, column: {field_name}, default_value: {:?}",
707                    field_meta.column_schema.default_constraint(),
708                ),
709            }
710        );
711
712        let default_value = field_meta
713            .column_schema
714            .create_default()
715            .context(CreateDefaultSnafu {
716                region_id: logical_region_id,
717                column: field_name,
718            })?
719            .with_context(|| InvalidRequestSnafu {
720                region_id: logical_region_id,
721                reason: format!("missing required field column {field_name}"),
722            })?;
723        let default_value = api::helper::to_grpc_value(default_value);
724        let (datatype, datatype_extension) =
725            ColumnDataTypeWrapper::try_from(field_meta.column_schema.data_type.clone())
726                .map_err(|e| {
727                    InvalidRequestSnafu {
728                        region_id: logical_region_id,
729                        reason: format!(
730                            "no protobuf type for field column {field_name} ({:?}): {e}",
731                            field_meta.column_schema.data_type
732                        ),
733                    }
734                    .build()
735                })?
736                .to_parts();
737
738        rows.schema.push(ColumnSchema {
739            column_name: field_name.to_string(),
740            datatype: datatype as i32,
741            semantic_type: SemanticType::Field as i32,
742            datatype_extension,
743            options: None,
744        });
745
746        for row in &mut rows.rows {
747            row.values.push(default_value.clone());
748        }
749
750        Ok(())
751    }
752
753    /// Perform metric engine specific logic to incoming rows.
754    /// - Add table_id column
755    /// - Generate tsid
756    fn modify_rows(
757        &self,
758        physical_region_id: RegionId,
759        table_id: TableId,
760        rows: &mut Rows,
761        encoding: PrimaryKeyEncoding,
762    ) -> Result<()> {
763        let input = std::mem::take(rows);
764        let iter = {
765            let state = self.state.read().unwrap();
766            let physical_columns = state
767                .physical_region_states()
768                .get(&physical_region_id)
769                .with_context(|| PhysicalRegionNotFoundSnafu {
770                    region_id: physical_region_id,
771                })?
772                .physical_columns();
773            RowsIter::new(input, physical_columns)
774        };
775        let output =
776            self.row_modifier
777                .modify_rows(iter, TableIdInput::Single(table_id), encoding)?;
778        *rows = output;
779        Ok(())
780    }
781}
782
783#[cfg(test)]
784mod tests {
785    use std::collections::HashSet;
786
787    use api::v1::value::ValueData;
788    use api::v1::{ColumnDataType, ColumnSchema as PbColumnSchema};
789    use common_error::ext::ErrorExt;
790    use common_error::status_code::StatusCode;
791    use common_function::utils::partition_expr_version;
792    use common_query::prelude::{greptime_native_histogram, greptime_timestamp, greptime_value};
793    use common_recordbatch::RecordBatches;
794    use datatypes::arrow::array::{
795        Float64Array, TimestampMicrosecondArray, TimestampMillisecondArray,
796    };
797    use datatypes::prelude::ConcreteDataType;
798    use datatypes::schema::{ColumnDefaultConstraint, ColumnSchema};
799    use datatypes::value::Value as PartitionValue;
800    use partition::expr::col;
801    use store_api::metadata::ColumnMetadata;
802    use store_api::metric_engine_consts::{
803        DATA_SCHEMA_TABLE_ID_COLUMN_NAME, DATA_SCHEMA_TSID_COLUMN_NAME, METRIC_ENGINE_NAME,
804        PHYSICAL_TABLE_METADATA_KEY, PRIMARY_KEY_ENCODING,
805    };
806    use store_api::path_utils::table_dir;
807    use store_api::region_engine::RegionEngine;
808    use store_api::region_request::{
809        EnterStagingRequest, PathType, RegionCloseRequest, RegionOpenRequest, RegionRequest,
810        StagingPartitionDirective,
811    };
812    use store_api::storage::ScanRequest;
813    use store_api::storage::consts::PRIMARY_KEY_COLUMN_NAME;
814
815    use super::*;
816    use crate::engine::MetricEngine;
817    use crate::test_util::{self, TestEnv};
818
819    async fn scan_timestamp_values(engine: &MetricEngine, region_id: RegionId) -> Vec<(i64, f64)> {
820        let stream = engine
821            .scan_to_stream(region_id, ScanRequest::default())
822            .await
823            .unwrap();
824        let batches = RecordBatches::try_collect(stream).await.unwrap();
825        let mut rows = Vec::new();
826        for batch in batches.iter() {
827            let batch = batch.df_record_batch();
828            let timestamp_index = batch.schema().index_of(greptime_timestamp()).unwrap();
829            let value_index = batch.schema().index_of(greptime_value()).unwrap();
830            let timestamps = batch
831                .column(timestamp_index)
832                .as_any()
833                .downcast_ref::<TimestampMillisecondArray>()
834                .unwrap();
835            let values = batch
836                .column(value_index)
837                .as_any()
838                .downcast_ref::<Float64Array>()
839                .unwrap();
840            rows.extend(
841                timestamps
842                    .values()
843                    .iter()
844                    .copied()
845                    .zip(values.values().iter().copied()),
846            );
847        }
848        rows.sort_unstable_by_key(|(timestamp, _)| *timestamp);
849        rows
850    }
851
852    #[tokio::test]
853    async fn test_batch_partition_versions() {
854        check_batch_partition_versions("sparse").await;
855        check_batch_partition_versions("dense").await;
856    }
857
858    async fn check_batch_partition_versions(encoding: &str) {
859        let env = TestEnv::new().await;
860        let physical_region_id = env.default_physical_region_id();
861        let logical_region_id = env.default_logical_region_id();
862        env.create_physical_region(
863            physical_region_id,
864            &TestEnv::default_table_dir(),
865            vec![(PRIMARY_KEY_ENCODING.to_string(), encoding.to_string())],
866        )
867        .await;
868        create_logical_region_with_tags(&env, physical_region_id, logical_region_id, &["job"])
869            .await;
870        let build_requests = |versions: [Option<u64>; 3]| {
871            versions
872                .into_iter()
873                .map(|partition_expr_version| {
874                    (
875                        logical_region_id,
876                        RegionPutRequest {
877                            skip_wal: false,
878                            rows: Rows {
879                                schema: test_util::row_schema_with_tags(&["job"]),
880                                rows: test_util::build_rows(1, 1),
881                            },
882                            hint: None,
883                            partition_expr_version,
884                        },
885                    )
886                })
887                .collect::<Vec<_>>()
888        };
889        // Conflicting explicit versions must fail before any data is written.
890        let err = env
891            .metric()
892            .inner
893            .put_regions_batch_single_physical(
894                physical_region_id,
895                build_requests([None, Some(10), Some(11)]),
896            )
897            .await
898            .unwrap_err();
899        assert!(
900            err.to_string()
901                .contains("inconsistent partition expr version")
902        );
903        assert!(
904            scan_timestamp_values(&env.metric(), logical_region_id)
905                .await
906                .is_empty()
907        );
908
909        for (versions, expected) in [
910            ([None, None, None], None),
911            ([None, Some(7), None], Some(7)),
912            ([Some(7), None, Some(7)], Some(7)),
913        ] {
914            let mut requests = build_requests(versions);
915            let engine = env.metric();
916            engine
917                .inner
918                .validate_batch_requests(physical_region_id, &mut requests)
919                .await
920                .unwrap();
921            let (merged, _) = match encoding {
922                "sparse" => engine
923                    .inner
924                    .merge_sparse_batch(physical_region_id, requests),
925                "dense" => engine
926                    .inner
927                    .merge_dense_batch(to_data_region_id(physical_region_id), requests),
928                _ => unreachable!(),
929            }
930            .unwrap();
931            assert_eq!(merged.partition_expr_version, expected);
932        }
933    }
934
935    #[tokio::test]
936    async fn test_put_skip_wal_batch_recovery() {
937        check_put_skip_wal_batch_recovery("sparse", false).await;
938        check_put_skip_wal_batch_recovery("sparse", true).await;
939        check_put_skip_wal_batch_recovery("dense", false).await;
940        check_put_skip_wal_batch_recovery("dense", true).await;
941    }
942
943    async fn check_put_skip_wal_batch_recovery(encoding: &str, skip_wal: bool) {
944        let env = TestEnv::new().await;
945        let engine = env.metric();
946        engine.inner.flush_task.stop().await.unwrap();
947        let physical_region_id = env.default_physical_region_id();
948        let logical_region_id = env.default_logical_region_id();
949        env.create_physical_region(
950            physical_region_id,
951            &TestEnv::default_table_dir(),
952            vec![(PRIMARY_KEY_ENCODING.to_string(), encoding.to_string())],
953        )
954        .await;
955        create_logical_region_with_tags(&env, physical_region_id, logical_region_id, &["job"])
956            .await;
957        let metadata_before = engine.get_metadata(logical_region_id).await.unwrap();
958
959        let requests = [skip_wal; 3]
960            .into_iter()
961            .enumerate()
962            .map(|(index, skip_wal)| {
963                let timestamp = index as i64 + 1;
964                let value = timestamp as f64 * 10.0;
965                // Every request updates the same key at timestamp zero and
966                // also inserts a distinct key to verify merge order.
967                let rows = [0, timestamp]
968                    .into_iter()
969                    .map(|timestamp| Row {
970                        values: vec![
971                            Value {
972                                value_data: Some(ValueData::TimestampMillisecondValue(timestamp)),
973                            },
974                            Value {
975                                value_data: Some(ValueData::F64Value(value)),
976                            },
977                            Value {
978                                value_data: Some(ValueData::StringValue("tag_0".to_string())),
979                            },
980                        ],
981                    })
982                    .collect();
983                (
984                    logical_region_id,
985                    RegionPutRequest {
986                        rows: Rows {
987                            schema: test_util::row_schema_with_tags(&["job"]),
988                            rows,
989                        },
990                        hint: None,
991                        partition_expr_version: None,
992                        skip_wal,
993                    },
994                )
995            });
996        let affected_rows = engine.inner.put_regions_batch(requests).await.unwrap();
997        assert_eq!(affected_rows, 6);
998        assert_eq!(
999            scan_timestamp_values(&engine, logical_region_id).await,
1000            vec![(0, 30.0), (1, 10.0), (2, 20.0), (3, 30.0)]
1001        );
1002
1003        // Neither data nor metadata has an SST to hide missing WAL.
1004        for region_id in [
1005            to_data_region_id(physical_region_id),
1006            crate::utils::to_metadata_region_id(physical_region_id),
1007        ] {
1008            let stat = env.mito().region_statistic(region_id).unwrap();
1009            assert!(stat.memtable_size > 0);
1010            assert_eq!(stat.sst_num, 0);
1011        }
1012        engine
1013            .handle_request(
1014                physical_region_id,
1015                RegionRequest::Close(RegionCloseRequest {
1016                    flush_on_close: false,
1017                }),
1018            )
1019            .await
1020            .unwrap();
1021
1022        // Recreate the wrapper as well, discarding its metadata cache.
1023        let reopened = MetricEngine::try_new(env.mito(), Default::default()).unwrap();
1024        reopened.inner.flush_task.stop().await.unwrap();
1025        reopened
1026            .handle_request(
1027                physical_region_id,
1028                RegionRequest::Open(RegionOpenRequest {
1029                    engine: METRIC_ENGINE_NAME.to_string(),
1030                    table_dir: TestEnv::default_table_dir(),
1031                    path_type: PathType::Bare,
1032                    options: [
1033                        (PHYSICAL_TABLE_METADATA_KEY.to_string(), String::new()),
1034                        (PRIMARY_KEY_ENCODING.to_string(), encoding.to_string()),
1035                    ]
1036                    .into_iter()
1037                    .collect(),
1038                    skip_wal_replay: false,
1039                    checkpoint: None,
1040                    requirements: Default::default(),
1041                }),
1042            )
1043            .await
1044            .unwrap();
1045        let recovered_metadata = reopened.get_metadata(logical_region_id).await.unwrap();
1046        assert_eq!(
1047            metadata_before.column_metadatas,
1048            recovered_metadata.column_metadatas
1049        );
1050        let expected = if skip_wal {
1051            vec![]
1052        } else {
1053            vec![(0, 30.0), (1, 10.0), (2, 20.0), (3, 30.0)]
1054        };
1055        assert_eq!(
1056            scan_timestamp_values(&reopened, logical_region_id).await,
1057            expected,
1058            "encoding={encoding}, skip_wal={skip_wal}"
1059        );
1060    }
1061
1062    fn assert_merged_schema(rows: &Rows, expect_sparse: bool) {
1063        let column_names: HashSet<String> = rows
1064            .schema
1065            .iter()
1066            .map(|col| col.column_name.clone())
1067            .collect();
1068
1069        if expect_sparse {
1070            assert!(
1071                column_names.contains(PRIMARY_KEY_COLUMN_NAME),
1072                "sparse encoding should include primary key column"
1073            );
1074            assert!(
1075                !column_names.contains(DATA_SCHEMA_TABLE_ID_COLUMN_NAME),
1076                "sparse encoding should not include table id column"
1077            );
1078            assert!(
1079                !column_names.contains(DATA_SCHEMA_TSID_COLUMN_NAME),
1080                "sparse encoding should not include tsid column"
1081            );
1082            assert!(
1083                !column_names.contains("job"),
1084                "sparse encoding should not include tag columns"
1085            );
1086            assert!(
1087                !column_names.contains("instance"),
1088                "sparse encoding should not include tag columns"
1089            );
1090        } else {
1091            assert!(
1092                !column_names.contains(PRIMARY_KEY_COLUMN_NAME),
1093                "dense encoding should not include primary key column"
1094            );
1095            assert!(
1096                column_names.contains(DATA_SCHEMA_TABLE_ID_COLUMN_NAME),
1097                "dense encoding should include table id column"
1098            );
1099            assert!(
1100                column_names.contains(DATA_SCHEMA_TSID_COLUMN_NAME),
1101                "dense encoding should include tsid column"
1102            );
1103            assert!(
1104                column_names.contains("job"),
1105                "dense encoding should keep tag columns"
1106            );
1107            assert!(
1108                column_names.contains("instance"),
1109                "dense encoding should keep tag columns"
1110            );
1111        }
1112    }
1113
1114    fn job_partition_expr_json() -> String {
1115        let expr = col("job")
1116            .gt_eq(PartitionValue::String("job-0".into()))
1117            .and(col("job").lt(PartitionValue::String("job-9".into())));
1118        expr.as_json_str().unwrap()
1119    }
1120
1121    #[tokio::test]
1122    async fn test_put_and_scan_microsecond_physical_region() {
1123        let env = TestEnv::new().await;
1124        let engine = env.metric();
1125        let physical_region_id = env.default_physical_region_id();
1126        let logical_region_id = env.default_logical_region_id();
1127        env.create_physical_region_with_ts_type(
1128            physical_region_id,
1129            &TestEnv::default_table_dir(),
1130            vec![],
1131            ConcreteDataType::timestamp_microsecond_datatype(),
1132        )
1133        .await;
1134
1135        // A logical region with a matching microsecond time index is accepted.
1136        let region_create_request = test_util::create_logical_region_request_with_ts_type(
1137            &["job"],
1138            physical_region_id,
1139            &table_dir("test", logical_region_id.table_id()),
1140            ConcreteDataType::timestamp_microsecond_datatype(),
1141        );
1142        engine
1143            .handle_request(
1144                logical_region_id,
1145                RegionRequest::Create(region_create_request),
1146            )
1147            .await
1148            .unwrap();
1149
1150        // Writing microsecond rows works.
1151        let affected_rows = engine
1152            .handle_request(
1153                logical_region_id,
1154                RegionRequest::Put(RegionPutRequest {
1155                    rows: Rows {
1156                        schema: test_util::row_schema_with_tags_and_ts_datatype(
1157                            &["job"],
1158                            ColumnDataType::TimestampMicrosecond,
1159                        ),
1160                        rows: test_util::build_rows_with_ts_datatype(
1161                            1,
1162                            2,
1163                            ColumnDataType::TimestampMicrosecond,
1164                        ),
1165                    },
1166                    hint: None,
1167                    partition_expr_version: None,
1168                    skip_wal: false,
1169                }),
1170            )
1171            .await
1172            .unwrap();
1173        assert_eq!(affected_rows.affected_rows, 2);
1174
1175        // The scan returns timestamps in the physical region's unit.
1176        let stream = engine
1177            .scan_to_stream(logical_region_id, ScanRequest::default())
1178            .await
1179            .unwrap();
1180        let batches = RecordBatches::try_collect(stream).await.unwrap();
1181        let mut rows = Vec::new();
1182        for batch in batches.iter() {
1183            let batch = batch.df_record_batch();
1184            let timestamps = batch
1185                .column(batch.schema().index_of(greptime_timestamp()).unwrap())
1186                .as_any()
1187                .downcast_ref::<TimestampMicrosecondArray>()
1188                .unwrap();
1189            let values = batch
1190                .column(batch.schema().index_of(greptime_value()).unwrap())
1191                .as_any()
1192                .downcast_ref::<Float64Array>()
1193                .unwrap();
1194            rows.extend(
1195                timestamps
1196                    .values()
1197                    .iter()
1198                    .copied()
1199                    .zip(values.values().iter().copied()),
1200            );
1201        }
1202        rows.sort_unstable_by_key(|(timestamp, _)| *timestamp);
1203        assert_eq!(rows, vec![(0, 0.0), (1, 1.0)]);
1204    }
1205
1206    async fn create_logical_region_with_tags(
1207        env: &TestEnv,
1208        physical_region_id: RegionId,
1209        logical_region_id: RegionId,
1210        tags: &[&str],
1211    ) {
1212        let region_create_request = test_util::create_logical_region_request(
1213            tags,
1214            physical_region_id,
1215            &table_dir("test", logical_region_id.table_id()),
1216        );
1217        env.metric()
1218            .handle_request(
1219                logical_region_id,
1220                RegionRequest::Create(region_create_request),
1221            )
1222            .await
1223            .unwrap();
1224    }
1225
1226    fn column_index(rows: &Rows, name: &str) -> usize {
1227        rows.schema
1228            .iter()
1229            .position(|col| col.column_name == name)
1230            .unwrap()
1231    }
1232
1233    fn check_batch_merge_wal_policy(
1234        env: &TestEnv,
1235        physical_region_id: RegionId,
1236        mut requests: Vec<(RegionId, RegionPutRequest)>,
1237        expect_sparse: bool,
1238        skip_wal: bool,
1239    ) {
1240        for (_, request) in &mut requests {
1241            request.skip_wal = skip_wal;
1242        }
1243        let (merged_request, affected_rows) = if expect_sparse {
1244            let (merged_request, affected_rows) = env
1245                .metric()
1246                .inner
1247                .merge_sparse_batch(physical_region_id, requests)
1248                .unwrap();
1249            let hint = merged_request
1250                .hint
1251                .as_ref()
1252                .expect("missing sparse write hint");
1253            assert_eq!(
1254                hint.primary_key_encoding,
1255                PrimaryKeyEncodingProto::Sparse as i32
1256            );
1257            (merged_request, affected_rows)
1258        } else {
1259            let (merged_request, affected_rows) = env
1260                .metric()
1261                .inner
1262                .merge_dense_batch(to_data_region_id(physical_region_id), requests)
1263                .unwrap();
1264            assert!(merged_request.hint.is_none());
1265            (merged_request, affected_rows)
1266        };
1267        assert_merged_schema(&merged_request.rows, expect_sparse);
1268        assert_eq!(merged_request.skip_wal, skip_wal);
1269        assert_eq!(affected_rows, 5);
1270    }
1271
1272    async fn run_batch_write_with_schema_variants(
1273        env: &TestEnv,
1274        physical_region_id: RegionId,
1275        options: Vec<(String, String)>,
1276        expect_sparse: bool,
1277    ) {
1278        env.create_physical_region(physical_region_id, &TestEnv::default_table_dir(), options)
1279            .await;
1280
1281        let logical_region_1 = env.default_logical_region_id();
1282        let logical_region_2 = RegionId::new(1024, 1);
1283
1284        create_logical_region_with_tags(env, physical_region_id, logical_region_1, &["job"]).await;
1285        create_logical_region_with_tags(
1286            env,
1287            physical_region_id,
1288            logical_region_2,
1289            &["job", "instance"],
1290        )
1291        .await;
1292
1293        let schema_1 = test_util::row_schema_with_tags(&["job"]);
1294        let schema_2 = test_util::row_schema_with_tags(&["job", "instance"]);
1295
1296        let data_region_id = RegionId::new(physical_region_id.table_id(), 2);
1297        let primary_key_encoding = env
1298            .metric()
1299            .inner
1300            .get_primary_key_encoding(data_region_id)
1301            .unwrap();
1302        assert_eq!(
1303            primary_key_encoding,
1304            if expect_sparse {
1305                PrimaryKeyEncoding::Sparse
1306            } else {
1307                PrimaryKeyEncoding::Dense
1308            }
1309        );
1310
1311        let build_requests = || {
1312            let rows_1 = test_util::build_rows(1, 3);
1313            let rows_2 = test_util::build_rows(2, 2);
1314
1315            vec![
1316                (
1317                    logical_region_1,
1318                    RegionPutRequest {
1319                        skip_wal: false,
1320                        rows: Rows {
1321                            schema: schema_1.clone(),
1322                            rows: rows_1,
1323                        },
1324                        hint: None,
1325                        partition_expr_version: None,
1326                    },
1327                ),
1328                (
1329                    logical_region_2,
1330                    RegionPutRequest {
1331                        skip_wal: false,
1332                        rows: Rows {
1333                            schema: schema_2.clone(),
1334                            rows: rows_2,
1335                        },
1336                        hint: None,
1337                        partition_expr_version: None,
1338                    },
1339                ),
1340            ]
1341        };
1342
1343        check_batch_merge_wal_policy(
1344            env,
1345            physical_region_id,
1346            build_requests(),
1347            expect_sparse,
1348            false,
1349        );
1350        check_batch_merge_wal_policy(
1351            env,
1352            physical_region_id,
1353            build_requests(),
1354            expect_sparse,
1355            true,
1356        );
1357
1358        for policies in [[false, true], [true, false]] {
1359            let mut mixed_requests = build_requests();
1360            for ((_, request), skip_wal) in mixed_requests.iter_mut().zip(policies) {
1361                request.skip_wal = skip_wal;
1362            }
1363            let err = env
1364                .metric()
1365                .inner
1366                .put_regions_batch(mixed_requests.into_iter())
1367                .await
1368                .unwrap_err();
1369            assert!(err.to_string().contains("inconsistent WAL policy in batch"));
1370            for logical_region_id in [logical_region_1, logical_region_2] {
1371                assert!(
1372                    scan_timestamp_values(&env.metric(), logical_region_id)
1373                        .await
1374                        .is_empty()
1375                );
1376            }
1377        }
1378
1379        let affected_rows = env
1380            .metric()
1381            .inner
1382            .put_regions_batch(build_requests().into_iter())
1383            .await
1384            .unwrap();
1385        assert_eq!(affected_rows, 5);
1386
1387        let request = ScanRequest::default();
1388        let stream = env
1389            .mito()
1390            .scan_to_stream(data_region_id, request)
1391            .await
1392            .unwrap();
1393        let batches = RecordBatches::try_collect(stream).await.unwrap();
1394
1395        assert_eq!(batches.iter().map(|b| b.num_rows()).sum::<usize>(), 5);
1396    }
1397
1398    #[test]
1399    fn test_sparse_batch_aligns_mixed_field_order() {
1400        let primary_key = PbColumnSchema {
1401            column_name: PRIMARY_KEY_COLUMN_NAME.to_string(),
1402            datatype: ColumnDataType::Binary as i32,
1403            semantic_type: SemanticType::Tag as _,
1404            datatype_extension: None,
1405            options: None,
1406        };
1407        let timestamp = PbColumnSchema {
1408            column_name: greptime_timestamp().to_string(),
1409            datatype: ColumnDataType::TimestampMillisecond as i32,
1410            semantic_type: SemanticType::Timestamp as _,
1411            datatype_extension: None,
1412            options: None,
1413        };
1414        let value = PbColumnSchema {
1415            column_name: greptime_value().to_string(),
1416            datatype: ColumnDataType::Float64 as i32,
1417            semantic_type: SemanticType::Field as _,
1418            datatype_extension: None,
1419            options: None,
1420        };
1421        let histogram = PbColumnSchema {
1422            column_name: greptime_native_histogram().to_string(),
1423            datatype: ColumnDataType::Struct as i32,
1424            semantic_type: SemanticType::Field as _,
1425            datatype_extension: None,
1426            options: None,
1427        };
1428
1429        let sample_rows = Rows {
1430            schema: vec![
1431                primary_key.clone(),
1432                timestamp.clone(),
1433                value.clone(),
1434                histogram.clone(),
1435            ],
1436            rows: vec![Row {
1437                values: vec![
1438                    ValueData::BinaryValue(vec![1]).into(),
1439                    ValueData::TimestampMillisecondValue(0).into(),
1440                    ValueData::F64Value(1.0).into(),
1441                    Value { value_data: None },
1442                ],
1443            }],
1444        };
1445        let histogram_rows = Rows {
1446            schema: vec![primary_key, timestamp, histogram, value],
1447            rows: vec![Row {
1448                values: vec![
1449                    ValueData::BinaryValue(vec![2]).into(),
1450                    ValueData::TimestampMillisecondValue(0).into(),
1451                    ValueData::StructValue(api::v1::StructValue { items: vec![] }).into(),
1452                    Value { value_data: None },
1453                ],
1454            }],
1455        };
1456
1457        let schema = MetricEngineInner::build_union_schema([
1458            sample_rows.schema.as_slice(),
1459            histogram_rows.schema.as_slice(),
1460        ]);
1461        let merged_rows = MetricEngineInner::align_rows_to_schema(sample_rows, &schema)
1462            .into_iter()
1463            .chain(MetricEngineInner::align_rows_to_schema(
1464                histogram_rows,
1465                &schema,
1466            ))
1467            .collect();
1468        let merged_request = Rows {
1469            schema,
1470            rows: merged_rows,
1471        };
1472
1473        let value_idx = column_index(&merged_request, greptime_value());
1474        let histogram_idx = column_index(&merged_request, greptime_native_histogram());
1475        assert!(matches!(
1476            merged_request.rows[0].values[value_idx].value_data,
1477            Some(ValueData::F64Value(_))
1478        ));
1479        assert!(
1480            merged_request.rows[0].values[histogram_idx]
1481                .value_data
1482                .is_none()
1483        );
1484        assert!(
1485            merged_request.rows[1].values[value_idx]
1486                .value_data
1487                .is_none()
1488        );
1489        assert!(matches!(
1490            merged_request.rows[1].values[histogram_idx].value_data,
1491            Some(ValueData::StructValue(_))
1492        ));
1493    }
1494
1495    #[tokio::test]
1496    async fn test_write_logical_region() {
1497        let env = TestEnv::new().await;
1498        env.init_metric_region().await;
1499
1500        // prepare data
1501        let schema = test_util::row_schema_with_tags(&["job"]);
1502        let rows = test_util::build_rows(1, 5);
1503        let request = RegionRequest::Put(RegionPutRequest {
1504            skip_wal: false,
1505            rows: Rows { schema, rows },
1506            hint: None,
1507            partition_expr_version: None,
1508        });
1509
1510        // write data
1511        let logical_region_id = env.default_logical_region_id();
1512        let result = env
1513            .metric()
1514            .handle_request(logical_region_id, request)
1515            .await
1516            .unwrap();
1517        assert_eq!(result.affected_rows, 5);
1518
1519        // read data from physical region
1520        let physical_region_id = env.default_physical_region_id();
1521        let request = ScanRequest::default();
1522        let stream = env
1523            .metric()
1524            .scan_to_stream(physical_region_id, request)
1525            .await
1526            .unwrap();
1527        let batches = RecordBatches::try_collect(stream).await.unwrap();
1528        let expected = "\
1529+-------------------------+----------------+------------+---------------------+-------+
1530| greptime_timestamp      | greptime_value | __table_id | __tsid              | job   |
1531+-------------------------+----------------+------------+---------------------+-------+
1532| 1970-01-01T00:00:00     | 0.0            | 3          | 2955007454552897459 | tag_0 |
1533| 1970-01-01T00:00:00.001 | 1.0            | 3          | 2955007454552897459 | tag_0 |
1534| 1970-01-01T00:00:00.002 | 2.0            | 3          | 2955007454552897459 | tag_0 |
1535| 1970-01-01T00:00:00.003 | 3.0            | 3          | 2955007454552897459 | tag_0 |
1536| 1970-01-01T00:00:00.004 | 4.0            | 3          | 2955007454552897459 | tag_0 |
1537+-------------------------+----------------+------------+---------------------+-------+";
1538        assert_eq!(expected, batches.pretty_print().unwrap(), "physical region");
1539
1540        // read data from logical region
1541        let request = ScanRequest::default();
1542        let stream = env
1543            .metric()
1544            .scan_to_stream(logical_region_id, request)
1545            .await
1546            .unwrap();
1547        let batches = RecordBatches::try_collect(stream).await.unwrap();
1548        let expected = "\
1549+-------------------------+----------------+-------+
1550| greptime_timestamp      | greptime_value | job   |
1551+-------------------------+----------------+-------+
1552| 1970-01-01T00:00:00     | 0.0            | tag_0 |
1553| 1970-01-01T00:00:00.001 | 1.0            | tag_0 |
1554| 1970-01-01T00:00:00.002 | 2.0            | tag_0 |
1555| 1970-01-01T00:00:00.003 | 3.0            | tag_0 |
1556| 1970-01-01T00:00:00.004 | 4.0            | tag_0 |
1557+-------------------------+----------------+-------+";
1558        assert_eq!(expected, batches.pretty_print().unwrap(), "logical region");
1559    }
1560
1561    #[tokio::test]
1562    async fn test_write_logical_region_row_count() {
1563        let env = TestEnv::new().await;
1564        env.init_metric_region().await;
1565        let engine = env.metric();
1566
1567        // add columns
1568        let logical_region_id = env.default_logical_region_id();
1569        let columns = &["odd", "even", "Ev_En"];
1570        let alter_request = test_util::alter_logical_region_add_tag_columns(123456, columns);
1571        engine
1572            .handle_request(logical_region_id, RegionRequest::Alter(alter_request))
1573            .await
1574            .unwrap();
1575
1576        // prepare data
1577        let schema = test_util::row_schema_with_tags(columns);
1578        let rows = test_util::build_rows(3, 100);
1579        let request = RegionRequest::Put(RegionPutRequest {
1580            skip_wal: false,
1581            rows: Rows { schema, rows },
1582            hint: None,
1583            partition_expr_version: None,
1584        });
1585
1586        // write data
1587        let result = engine
1588            .handle_request(logical_region_id, request)
1589            .await
1590            .unwrap();
1591        assert_eq!(100, result.affected_rows);
1592    }
1593
1594    #[tokio::test]
1595    async fn test_write_physical_region() {
1596        let env = TestEnv::new().await;
1597        env.init_metric_region().await;
1598        let engine = env.metric();
1599
1600        let physical_region_id = env.default_physical_region_id();
1601        let schema = test_util::row_schema_with_tags(&["abc"]);
1602        let rows = test_util::build_rows(1, 100);
1603        let request = RegionRequest::Put(RegionPutRequest {
1604            skip_wal: false,
1605            rows: Rows { schema, rows },
1606            hint: None,
1607            partition_expr_version: None,
1608        });
1609
1610        engine
1611            .handle_request(physical_region_id, request)
1612            .await
1613            .unwrap_err();
1614    }
1615
1616    #[tokio::test]
1617    async fn test_write_nonexist_logical_region() {
1618        let env = TestEnv::new().await;
1619        env.init_metric_region().await;
1620        let engine = env.metric();
1621
1622        let logical_region_id = RegionId::new(175, 8345);
1623        let schema = test_util::row_schema_with_tags(&["def"]);
1624        let rows = test_util::build_rows(1, 100);
1625        let request = RegionRequest::Put(RegionPutRequest {
1626            skip_wal: false,
1627            rows: Rows { schema, rows },
1628            hint: None,
1629            partition_expr_version: None,
1630        });
1631
1632        engine
1633            .handle_request(logical_region_id, request)
1634            .await
1635            .unwrap_err();
1636    }
1637
1638    #[tokio::test]
1639    async fn test_batch_write_multiple_logical_regions() {
1640        let env = TestEnv::new().await;
1641        env.init_metric_region().await;
1642        let engine = env.metric();
1643
1644        // Create two additional logical regions
1645        let physical_region_id = env.default_physical_region_id();
1646        let logical_region_1 = env.default_logical_region_id();
1647        let logical_region_2 = RegionId::new(1024, 1);
1648        let logical_region_3 = RegionId::new(1024, 2);
1649
1650        env.create_logical_region(physical_region_id, logical_region_2)
1651            .await;
1652        env.create_logical_region(physical_region_id, logical_region_3)
1653            .await;
1654
1655        // Prepare batch requests with non-overlapping timestamps
1656        let schema = test_util::row_schema_with_tags(&["job"]);
1657
1658        // Use build_rows_with_ts to create non-overlapping timestamps
1659        // logical_region_1: ts 0, 1, 2
1660        // logical_region_2: ts 10, 11  (offset to avoid overlap)
1661        // logical_region_3: ts 20, 21, 22, 23, 24  (offset to avoid overlap)
1662        let rows1 = test_util::build_rows(1, 3);
1663        let mut rows2 = test_util::build_rows(1, 2);
1664        let mut rows3 = test_util::build_rows(1, 5);
1665
1666        // Adjust timestamps to avoid conflicts
1667        use api::v1::value::ValueData;
1668        for (i, row) in rows2.iter_mut().enumerate() {
1669            if let Some(ValueData::TimestampMillisecondValue(ts)) =
1670                row.values.get_mut(0).and_then(|v| v.value_data.as_mut())
1671            {
1672                *ts = (10 + i) as i64;
1673            }
1674        }
1675        for (i, row) in rows3.iter_mut().enumerate() {
1676            if let Some(ValueData::TimestampMillisecondValue(ts)) =
1677                row.values.get_mut(0).and_then(|v| v.value_data.as_mut())
1678            {
1679                *ts = (20 + i) as i64;
1680            }
1681        }
1682
1683        let requests = vec![
1684            (
1685                logical_region_1,
1686                RegionPutRequest {
1687                    skip_wal: false,
1688                    rows: Rows {
1689                        schema: schema.clone(),
1690                        rows: rows1,
1691                    },
1692                    hint: None,
1693                    partition_expr_version: None,
1694                },
1695            ),
1696            (
1697                logical_region_2,
1698                RegionPutRequest {
1699                    skip_wal: false,
1700                    rows: Rows {
1701                        schema: schema.clone(),
1702                        rows: rows2,
1703                    },
1704                    hint: None,
1705                    partition_expr_version: None,
1706                },
1707            ),
1708            (
1709                logical_region_3,
1710                RegionPutRequest {
1711                    skip_wal: false,
1712                    rows: Rows {
1713                        schema: schema.clone(),
1714                        rows: rows3,
1715                    },
1716                    hint: None,
1717                    partition_expr_version: None,
1718                },
1719            ),
1720        ];
1721
1722        // Batch write
1723        let affected_rows = engine
1724            .inner
1725            .put_regions_batch(requests.into_iter())
1726            .await
1727            .unwrap();
1728        assert_eq!(affected_rows, 10);
1729
1730        // Verify physical region contains data from all logical regions
1731        let request = ScanRequest::default();
1732        let stream = env
1733            .metric()
1734            .scan_to_stream(physical_region_id, request)
1735            .await
1736            .unwrap();
1737        let batches = RecordBatches::try_collect(stream).await.unwrap();
1738
1739        // Should have 3 + 2 + 5 = 10 rows total
1740        assert_eq!(batches.iter().map(|b| b.num_rows()).sum::<usize>(), 10);
1741    }
1742
1743    #[tokio::test]
1744    async fn test_batch_write_with_partial_failure() {
1745        let env = TestEnv::new().await;
1746        env.init_metric_region().await;
1747        let engine = env.metric();
1748
1749        let physical_region_id = env.default_physical_region_id();
1750        let logical_region_1 = env.default_logical_region_id();
1751        let logical_region_2 = RegionId::new(1024, 1);
1752        let nonexistent_region = RegionId::new(9999, 9999);
1753
1754        env.create_logical_region(physical_region_id, logical_region_2)
1755            .await;
1756
1757        // Prepare batch with one invalid region
1758        let schema = test_util::row_schema_with_tags(&["job"]);
1759        let requests = vec![
1760            (
1761                logical_region_1,
1762                RegionPutRequest {
1763                    skip_wal: false,
1764                    rows: Rows {
1765                        schema: schema.clone(),
1766                        rows: test_util::build_rows(1, 3),
1767                    },
1768                    hint: None,
1769                    partition_expr_version: None,
1770                },
1771            ),
1772            (
1773                nonexistent_region,
1774                RegionPutRequest {
1775                    skip_wal: false,
1776                    rows: Rows {
1777                        schema: schema.clone(),
1778                        rows: test_util::build_rows(1, 2),
1779                    },
1780                    hint: None,
1781                    partition_expr_version: None,
1782                },
1783            ),
1784            (
1785                logical_region_2,
1786                RegionPutRequest {
1787                    skip_wal: false,
1788                    rows: Rows {
1789                        schema: schema.clone(),
1790                        rows: test_util::build_rows(1, 5),
1791                    },
1792                    hint: None,
1793                    partition_expr_version: None,
1794                },
1795            ),
1796        ];
1797
1798        // Batch write
1799        let result = engine.inner.put_regions_batch(requests.into_iter()).await;
1800        assert!(result.is_err());
1801
1802        // Invalid region is detected before any write, so the physical region remains empty.
1803        // Fail-fast is per physical-region group; cross-group partial success is possible.
1804        let request = ScanRequest::default();
1805        let stream = env
1806            .metric()
1807            .scan_to_stream(physical_region_id, request)
1808            .await
1809            .unwrap();
1810        let batches = RecordBatches::try_collect(stream).await.unwrap();
1811
1812        assert_eq!(batches.iter().map(|b| b.num_rows()).sum::<usize>(), 0);
1813    }
1814
1815    #[tokio::test]
1816    async fn test_batch_write_single_physical_region_forbidden() {
1817        let env = TestEnv::new().await;
1818        env.init_metric_region().await;
1819        let engine = env.metric();
1820
1821        let physical_region_id = env.default_physical_region_id();
1822        let schema = test_util::row_schema_with_tags(&["job"]);
1823        let requests = vec![(
1824            physical_region_id,
1825            RegionPutRequest {
1826                skip_wal: false,
1827                rows: Rows {
1828                    schema,
1829                    rows: test_util::build_rows(1, 1),
1830                },
1831                hint: None,
1832                partition_expr_version: None,
1833            },
1834        )];
1835
1836        let err = engine
1837            .inner
1838            .put_regions_batch(requests.into_iter())
1839            .await
1840            .unwrap_err();
1841
1842        assert!(matches!(
1843            err,
1844            crate::error::Error::ForbiddenPhysicalWrite { .. }
1845        ));
1846    }
1847
1848    #[tokio::test]
1849    async fn test_batch_write_physical_region_forbidden() {
1850        let env = TestEnv::new().await;
1851        env.init_metric_region().await;
1852        let engine = env.metric();
1853
1854        let physical_region_id = env.default_physical_region_id();
1855        let logical_region_id = env.default_logical_region_id();
1856        let schema = test_util::row_schema_with_tags(&["job"]);
1857        let requests = vec![
1858            (
1859                logical_region_id,
1860                RegionPutRequest {
1861                    skip_wal: false,
1862                    rows: Rows {
1863                        schema: schema.clone(),
1864                        rows: test_util::build_rows(1, 1),
1865                    },
1866                    hint: None,
1867                    partition_expr_version: None,
1868                },
1869            ),
1870            (
1871                physical_region_id,
1872                RegionPutRequest {
1873                    skip_wal: false,
1874                    rows: Rows {
1875                        schema,
1876                        rows: test_util::build_rows(1, 1),
1877                    },
1878                    hint: None,
1879                    partition_expr_version: None,
1880                },
1881            ),
1882        ];
1883
1884        let err = engine
1885            .inner
1886            .put_regions_batch(requests.into_iter())
1887            .await
1888            .unwrap_err();
1889
1890        assert!(matches!(
1891            err,
1892            crate::error::Error::ForbiddenPhysicalWrite { .. }
1893        ));
1894    }
1895
1896    #[tokio::test]
1897    async fn test_batch_write_single_request_fast_path() {
1898        let env = TestEnv::new().await;
1899        env.init_metric_region().await;
1900        let engine = env.metric();
1901
1902        let logical_region_id = env.default_logical_region_id();
1903        let schema = test_util::row_schema_with_tags(&["job"]);
1904
1905        // Single request should use fast path
1906        let requests = vec![(
1907            logical_region_id,
1908            RegionPutRequest {
1909                skip_wal: false,
1910                rows: Rows {
1911                    schema,
1912                    rows: test_util::build_rows(1, 5),
1913                },
1914                hint: None,
1915                partition_expr_version: None,
1916            },
1917        )];
1918
1919        let affected_rows = engine
1920            .inner
1921            .put_regions_batch(requests.into_iter())
1922            .await
1923            .unwrap();
1924        assert_eq!(affected_rows, 5);
1925    }
1926
1927    #[tokio::test]
1928    async fn test_batch_write_empty_requests() {
1929        let env = TestEnv::new().await;
1930        env.init_metric_region().await;
1931        let engine = env.metric();
1932
1933        // Empty batch should return zero affected rows
1934        let requests = vec![];
1935        let affected_rows = engine
1936            .inner
1937            .put_regions_batch(requests.into_iter())
1938            .await
1939            .unwrap();
1940
1941        assert_eq!(affected_rows, 0);
1942    }
1943
1944    #[tokio::test]
1945    async fn test_batch_write_sparse_encoding() {
1946        let env = TestEnv::new().await;
1947        let physical_region_id = env.default_physical_region_id();
1948
1949        run_batch_write_with_schema_variants(
1950            &env,
1951            physical_region_id,
1952            vec![(PRIMARY_KEY_ENCODING.to_string(), "sparse".to_string())],
1953            true,
1954        )
1955        .await;
1956    }
1957
1958    #[tokio::test]
1959    async fn test_batch_write_dense_encoding() {
1960        let env = TestEnv::new().await;
1961        let physical_region_id = env.default_physical_region_id();
1962
1963        run_batch_write_with_schema_variants(
1964            &env,
1965            physical_region_id,
1966            vec![(PRIMARY_KEY_ENCODING.to_string(), "dense".to_string())],
1967            false,
1968        )
1969        .await;
1970    }
1971
1972    #[tokio::test]
1973    async fn test_metric_put_rejects_bad_partition_expr_version() {
1974        let env = TestEnv::new().await;
1975        env.init_metric_region().await;
1976
1977        let logical_region_id = env.default_logical_region_id();
1978        let rows = Rows {
1979            schema: test_util::row_schema_with_tags(&["job"]),
1980            rows: test_util::build_rows(1, 3),
1981        };
1982
1983        let err = env
1984            .metric()
1985            .handle_request(
1986                logical_region_id,
1987                RegionRequest::Put(RegionPutRequest {
1988                    skip_wal: false,
1989                    rows,
1990                    hint: None,
1991                    partition_expr_version: Some(1),
1992                }),
1993            )
1994            .await
1995            .unwrap_err();
1996
1997        assert_eq!(err.status_code(), StatusCode::InvalidArguments);
1998    }
1999
2000    #[tokio::test]
2001    async fn test_metric_put_respects_staging_partition_expr_version() {
2002        let env = TestEnv::new().await;
2003        env.init_metric_region().await;
2004
2005        let logical_region_id = env.default_logical_region_id();
2006        let physical_region_id = env.default_physical_region_id();
2007        let partition_expr = job_partition_expr_json();
2008        env.metric()
2009            .handle_request(
2010                physical_region_id,
2011                RegionRequest::EnterStaging(EnterStagingRequest {
2012                    partition_directive: StagingPartitionDirective::UpdatePartitionExpr(
2013                        partition_expr.clone(),
2014                    ),
2015                }),
2016            )
2017            .await
2018            .unwrap();
2019
2020        let expected_version = partition_expr_version(Some(&partition_expr));
2021        let rows = Rows {
2022            schema: test_util::row_schema_with_tags(&["job"]),
2023            rows: test_util::build_rows(1, 3),
2024        };
2025
2026        let err = env
2027            .metric()
2028            .handle_request(
2029                logical_region_id,
2030                RegionRequest::Put(RegionPutRequest {
2031                    skip_wal: false,
2032                    rows: rows.clone(),
2033                    hint: None,
2034                    partition_expr_version: Some(expected_version.wrapping_add(1)),
2035                }),
2036            )
2037            .await
2038            .unwrap_err();
2039        assert_eq!(err.status_code(), StatusCode::InvalidArguments);
2040
2041        let response = env
2042            .metric()
2043            .handle_request(
2044                logical_region_id,
2045                RegionRequest::Put(RegionPutRequest {
2046                    skip_wal: false,
2047                    rows: rows.clone(),
2048                    hint: None,
2049                    partition_expr_version: None,
2050                }),
2051            )
2052            .await
2053            .unwrap();
2054        assert_eq!(response.affected_rows, 3);
2055
2056        let response = env
2057            .metric()
2058            .handle_request(
2059                logical_region_id,
2060                RegionRequest::Put(RegionPutRequest {
2061                    skip_wal: false,
2062                    rows,
2063                    hint: None,
2064                    partition_expr_version: Some(expected_version),
2065                }),
2066            )
2067            .await
2068            .unwrap();
2069        assert_eq!(response.affected_rows, 3);
2070    }
2071
2072    /// Regression test for issue #7990: the metric engine must reject a row
2073    /// whose timestamp column carries a non-timestamp datatype, rather than
2074    /// letting it panic inside mito's `ValueBuilder::push`.
2075    #[tokio::test]
2076    async fn test_verify_rows_rejects_wrong_type() {
2077        use api::v1::value::ValueData;
2078        use api::v1::{ColumnDataType, ColumnSchema as PbColumnSchema, SemanticType};
2079        use common_query::prelude::{greptime_timestamp, greptime_value};
2080
2081        let env = TestEnv::new().await;
2082        env.init_metric_region().await;
2083
2084        let logical_region_id = env.default_logical_region_id();
2085
2086        // Timestamp column is declared as String — the very payload that
2087        // caused #7990. It should surface a typed error rather than panic.
2088        let schema = vec![
2089            PbColumnSchema {
2090                column_name: greptime_timestamp().to_string(),
2091                datatype: ColumnDataType::String as i32,
2092                semantic_type: SemanticType::Timestamp as _,
2093                datatype_extension: None,
2094                options: None,
2095            },
2096            PbColumnSchema {
2097                column_name: greptime_value().to_string(),
2098                datatype: ColumnDataType::Float64 as i32,
2099                semantic_type: SemanticType::Field as _,
2100                datatype_extension: None,
2101                options: None,
2102            },
2103            PbColumnSchema {
2104                column_name: "job".to_string(),
2105                datatype: ColumnDataType::String as i32,
2106                semantic_type: SemanticType::Tag as _,
2107                datatype_extension: None,
2108                options: None,
2109            },
2110        ];
2111        let rows = vec![Row {
2112            values: vec![
2113                Value {
2114                    value_data: Some(ValueData::StringValue("not-a-timestamp".to_string())),
2115                },
2116                Value {
2117                    value_data: Some(ValueData::F64Value(1.0)),
2118                },
2119                Value {
2120                    value_data: Some(ValueData::StringValue("tag_0".to_string())),
2121                },
2122            ],
2123        }];
2124
2125        let err = env
2126            .metric()
2127            .handle_request(
2128                logical_region_id,
2129                RegionRequest::Put(RegionPutRequest {
2130                    skip_wal: false,
2131                    rows: Rows { schema, rows },
2132                    hint: None,
2133                    partition_expr_version: None,
2134                }),
2135            )
2136            .await
2137            .unwrap_err();
2138        assert_eq!(err.status_code(), StatusCode::InvalidArguments);
2139    }
2140
2141    /// The completeness check must reject requests that omit the time index
2142    /// column, since mito cannot default-fill a `TimeIndex` column and would
2143    /// previously panic on the empty builder.
2144    #[tokio::test]
2145    async fn test_verify_rows_rejects_missing_time_index() {
2146        use api::v1::{ColumnDataType, ColumnSchema as PbColumnSchema, SemanticType};
2147        use common_query::prelude::greptime_value;
2148
2149        let env = TestEnv::new().await;
2150        env.init_metric_region().await;
2151
2152        let logical_region_id = env.default_logical_region_id();
2153
2154        // Payload only carries the field and a tag — no timestamp column.
2155        let schema = vec![
2156            PbColumnSchema {
2157                column_name: greptime_value().to_string(),
2158                datatype: ColumnDataType::Float64 as i32,
2159                semantic_type: SemanticType::Field as _,
2160                datatype_extension: None,
2161                options: None,
2162            },
2163            PbColumnSchema {
2164                column_name: "job".to_string(),
2165                datatype: ColumnDataType::String as i32,
2166                semantic_type: SemanticType::Tag as _,
2167                datatype_extension: None,
2168                options: None,
2169            },
2170        ];
2171        let rows = vec![Row {
2172            values: vec![
2173                Value {
2174                    value_data: Some(api::v1::value::ValueData::F64Value(1.0)),
2175                },
2176                Value {
2177                    value_data: Some(api::v1::value::ValueData::StringValue("tag_0".to_string())),
2178                },
2179            ],
2180        }];
2181
2182        let err = env
2183            .metric()
2184            .handle_request(
2185                logical_region_id,
2186                RegionRequest::Put(RegionPutRequest {
2187                    skip_wal: false,
2188                    rows: Rows { schema, rows },
2189                    hint: None,
2190                    partition_expr_version: None,
2191                }),
2192            )
2193            .await
2194            .unwrap_err();
2195        assert_eq!(err.status_code(), StatusCode::InvalidArguments);
2196    }
2197
2198    #[tokio::test]
2199    async fn test_verify_rows_rejects_missing_field() {
2200        use api::v1::value::ValueData;
2201        use api::v1::{ColumnDataType, ColumnSchema as PbColumnSchema, SemanticType};
2202        use common_query::prelude::greptime_timestamp;
2203
2204        let env = TestEnv::new().await;
2205        env.init_metric_region().await;
2206
2207        let logical_region_id = env.default_logical_region_id();
2208
2209        // Schema has timestamp + tag but no field column.
2210        let schema = vec![
2211            PbColumnSchema {
2212                column_name: greptime_timestamp().to_string(),
2213                datatype: ColumnDataType::TimestampMillisecond as i32,
2214                semantic_type: SemanticType::Timestamp as _,
2215                datatype_extension: None,
2216                options: None,
2217            },
2218            PbColumnSchema {
2219                column_name: "job".to_string(),
2220                datatype: ColumnDataType::String as i32,
2221                semantic_type: SemanticType::Tag as _,
2222                datatype_extension: None,
2223                options: None,
2224            },
2225        ];
2226        let rows = vec![Row {
2227            values: vec![
2228                Value {
2229                    value_data: Some(ValueData::TimestampMillisecondValue(0)),
2230                },
2231                Value {
2232                    value_data: Some(ValueData::StringValue("tag_0".to_string())),
2233                },
2234            ],
2235        }];
2236
2237        let err = env
2238            .metric()
2239            .handle_request(
2240                logical_region_id,
2241                RegionRequest::Put(RegionPutRequest {
2242                    skip_wal: false,
2243                    rows: Rows { schema, rows },
2244                    hint: None,
2245                    partition_expr_version: None,
2246                }),
2247            )
2248            .await
2249            .unwrap_err();
2250        let message = err.to_string();
2251        assert!(
2252            message.contains("missing required field column"),
2253            "expected field-completeness rejection, got: {message}"
2254        );
2255        assert_eq!(err.status_code(), StatusCode::InvalidArguments);
2256    }
2257
2258    #[test]
2259    fn test_fill_missing_field_column_nullable_no_default() {
2260        let field_meta = ColumnMetadata {
2261            column_id: 1,
2262            semantic_type: SemanticType::Field,
2263            column_schema: ColumnSchema::new(
2264                "greptime_value".to_string(),
2265                ConcreteDataType::float64_datatype(),
2266                true, // nullable, no default
2267            ),
2268        };
2269        let mut rows = Rows {
2270            schema: vec![PbColumnSchema {
2271                column_name: "ts".to_string(),
2272                datatype: ColumnDataType::TimestampMillisecond as i32,
2273                semantic_type: SemanticType::Timestamp as _,
2274                datatype_extension: None,
2275                options: None,
2276            }],
2277            rows: vec![Row {
2278                values: vec![Value {
2279                    value_data: Some(ValueData::TimestampMillisecondValue(0)),
2280                }],
2281            }],
2282        };
2283
2284        MetricEngineInner::fill_missing_field_column(
2285            RegionId::new(1, 1),
2286            "greptime_value",
2287            &field_meta,
2288            &mut rows,
2289        )
2290        .unwrap();
2291
2292        assert_eq!(rows.schema.len(), 2);
2293        assert_eq!(rows.schema[1].column_name, "greptime_value");
2294        assert_eq!(rows.rows[0].values.len(), 2);
2295        assert!(
2296            rows.rows[0].values[1].value_data.is_none(),
2297            "missing nullable field should be filled with null"
2298        );
2299    }
2300
2301    #[test]
2302    fn test_fill_missing_field_column_rejects_impure_default() {
2303        let field_meta = ColumnMetadata {
2304            column_id: 1,
2305            semantic_type: SemanticType::Field,
2306            column_schema: ColumnSchema::new(
2307                "greptime_value".to_string(),
2308                ConcreteDataType::timestamp_millisecond_datatype(),
2309                false,
2310            )
2311            .with_default_constraint(Some(ColumnDefaultConstraint::Function("now()".to_string())))
2312            .unwrap(),
2313        };
2314        let mut rows = Rows {
2315            schema: vec![PbColumnSchema {
2316                column_name: "ts".to_string(),
2317                datatype: api::v1::ColumnDataType::TimestampMillisecond as i32,
2318                semantic_type: SemanticType::Timestamp as _,
2319                datatype_extension: None,
2320                options: None,
2321            }],
2322            rows: vec![Row {
2323                values: vec![Value {
2324                    value_data: Some(ValueData::TimestampMillisecondValue(0)),
2325                }],
2326            }],
2327        };
2328
2329        let err = MetricEngineInner::fill_missing_field_column(
2330            RegionId::new(1, 1),
2331            "greptime_value",
2332            &field_meta,
2333            &mut rows,
2334        )
2335        .unwrap_err();
2336        assert!(
2337            err.to_string().contains("impure default value"),
2338            "expected impure-default rejection, got: {err}"
2339        );
2340    }
2341}