1use 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 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 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 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 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 self.validate_batch_requests(physical_region_id, &mut requests)
149 .await?;
150
151 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 self.data_region
159 .write_data(data_region_id, RegionRequest::Put(merged_request))
160 .await?;
161
162 Ok(total_affected_rows)
163 }
164
165 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 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 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 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 let merged_schema =
282 Self::build_union_schema(requests.iter().map(|(_, req)| req.rows.schema.as_slice()));
283
284 let (merged_rows, table_ids, merged_version) =
286 Self::align_requests_to_schema(requests, &merged_schema)?;
287
288 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 let schema = test_util::row_schema_with_tags(&["job"]);
1657
1658 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 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 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 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 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 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 let result = engine.inner.put_regions_batch(requests.into_iter()).await;
1800 assert!(result.is_err());
1801
1802 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 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 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 #[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 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 #[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 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 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, ),
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}