Skip to main content

mito2/
sst.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Sorted strings tables.
16
17use std::collections::HashMap;
18use std::sync::Arc;
19
20use api::v1::SemanticType;
21use arrow_schema::DataType;
22use arrow_schema::extension::{EXTENSION_TYPE_NAME_KEY, ExtensionType};
23use common_base::readable_size::ReadableSize;
24use common_query::native_histogram::{
25    is_native_histogram_value_type, native_histogram_list_element_id, native_histogram_subfield_id,
26};
27use datatypes::arrow::datatypes::{
28    DataType as ArrowDataType, Field, FieldRef, Fields, Schema, SchemaRef,
29};
30use datatypes::arrow::record_batch::RecordBatch;
31use datatypes::extension::histogram::HistogramExtensionType;
32use datatypes::prelude::ConcreteDataType;
33use datatypes::timestamp::timestamp_array_to_primitive;
34use serde::{Deserialize, Serialize};
35use store_api::codec::PrimaryKeyEncoding;
36use store_api::metadata::RegionMetadata;
37use store_api::storage::consts::{
38    OP_TYPE_COLUMN_NAME, PRIMARY_KEY_COLUMN_NAME, SEQUENCE_COLUMN_NAME,
39};
40
41use crate::error::{InvalidNativeHistogramFieldIdSnafu, InvalidNativeHistogramSubfieldSnafu};
42use crate::sst::parquet::flat_format::time_index_column_index;
43
44pub mod file;
45pub mod file_purger;
46pub mod file_ref;
47pub mod index;
48pub mod location;
49pub mod parquet;
50pub(crate) mod primary_key;
51pub mod range_index;
52pub(crate) mod version;
53
54/// Default write buffer size, it should be greater than the default minimum upload part of S3 (5mb).
55pub const DEFAULT_WRITE_BUFFER_SIZE: ReadableSize = ReadableSize::mb(8);
56
57/// Default number of concurrent write, it only works on object store backend(e.g., S3).
58pub const DEFAULT_WRITE_CONCURRENCY: usize = 8;
59
60/// Format type of the SST file.
61#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, strum::EnumString)]
62#[serde(rename_all = "snake_case")]
63#[strum(serialize_all = "snake_case")]
64pub enum FormatType {
65    /// Parquet with primary key encoded.
66    #[default]
67    PrimaryKey,
68    /// Flat Parquet format.
69    Flat,
70}
71
72/// Iceberg-compatible column field ID key stored in Parquet column metadata.
73pub const PARQUET_FIELD_ID_KEY: &str = "PARQUET:field_id";
74
75/// Adds `PARQUET:field_id` metadata to a top-level Arrow field.
76///
77/// Native-histogram sub-field ids are stamped separately at parquet-write
78/// time by [`stamp_native_histogram_subfield_ids`], not here.
79pub fn with_field_id(mut field: Field, column_id: u32) -> Field {
80    field
81        .metadata_mut()
82        .insert(PARQUET_FIELD_ID_KEY.to_string(), column_id.to_string());
83    field
84}
85
86/// Stamps the `greptime.histogram` extension and reserved `PARQUET:field_id`s
87/// onto a native-histogram struct field (and its sub-fields / list element
88/// fields), so external readers can identify it by extension and resolve
89/// nested fields by id.
90///
91/// Detection is by the exact native-histogram struct type
92/// (`is_native_histogram_value_type`); other struct columns are left untouched.
93/// mito2 reads SST columns by schema position, never by field metadata, so this
94/// only affects external readers.
95///
96/// Returns an error if the parent column's `PARQUET:field_id` is missing,
97/// malformed, or exceeds `i32::MAX`, or if a sub-field id cannot be derived
98/// — because the sub-field name is not a known native-histogram field, or
99/// the derived id overflows a positive `i32` (an absurdly large parent
100/// `column_id`); see [`native_histogram_subfield_id`].
101fn stamp_native_histogram_subfield_ids(field: &mut Field) -> crate::error::Result<()> {
102    if !is_native_histogram_value_type(&ConcreteDataType::from_arrow_type(field.data_type())) {
103        return Ok(());
104    }
105    // Namespace sub-field ids by the parent column's field id (its
106    // `PARQUET:field_id`, stamped earlier by `with_field_id`) so several
107    // histogram columns in one table get disjoint ids. Fail loudly if the id
108    // is absent, malformed, or too large to fit a positive `i32`.
109    let column_id = field
110        .metadata()
111        .get(PARQUET_FIELD_ID_KEY)
112        .and_then(|s| s.parse::<i32>().ok())
113        .ok_or_else(|| {
114            InvalidNativeHistogramFieldIdSnafu {
115                field_name: field.name().clone(),
116            }
117            .build()
118        })?;
119    // Tag the field with the greptime.histogram extension.
120    field.metadata_mut().insert(
121        EXTENSION_TYPE_NAME_KEY.to_string(),
122        HistogramExtensionType::NAME.to_string(),
123    );
124    let ArrowDataType::Struct(children) = field.data_type() else {
125        return Ok(());
126    };
127    let new_children: crate::error::Result<Fields> = children
128        .iter()
129        .map(|child| {
130            let mut c = (**child).clone();
131            // `None` here means either the sub-field name is not a known
132            // native-histogram field, or the derived id overflowed i32.
133            // Surface it as an error rather than silently leaving the field
134            // without an id.
135            let id = native_histogram_subfield_id(column_id, c.name()).ok_or_else(|| {
136                InvalidNativeHistogramSubfieldSnafu {
137                    column_id,
138                    field_name: c.name().clone(),
139                }
140                .build()
141            })?;
142            // Stamp the sub-field's own id.
143            c.metadata_mut()
144                .insert(PARQUET_FIELD_ID_KEY.to_string(), id.to_string());
145            // If the sub-field is a list, stamp its element field's id.
146            if let ArrowDataType::List(elem) = c.data_type() {
147                let elem_id =
148                    native_histogram_list_element_id(column_id, c.name()).ok_or_else(|| {
149                        InvalidNativeHistogramSubfieldSnafu {
150                            column_id,
151                            field_name: c.name().clone(),
152                        }
153                        .build()
154                    })?;
155                let mut new_elem = (**elem).clone();
156                new_elem
157                    .metadata_mut()
158                    .insert(PARQUET_FIELD_ID_KEY.to_string(), elem_id.to_string());
159                c.set_data_type(ArrowDataType::List(Arc::new(new_elem)));
160            }
161            Ok(Arc::new(c))
162        })
163        .collect();
164    field.set_data_type(ArrowDataType::Struct(new_children?));
165    Ok(())
166}
167
168/// Returns a copy of `schema` with native-histogram sub-field ids stamped,
169/// for the parquet writer.
170///
171/// This is called on the schema handed to `AsyncArrowWriter`, not in
172/// [`with_field_id`], because the SST arrow schema is also the memtable's
173/// in-memory schema, whose `Struct` equality (`PartialEq`) is
174/// metadata-sensitive — stamping there would break writes. The parquet writer
175/// compares types with `DataType::equals_datatype`, which ignores field
176/// metadata, so a stamped schema accepts an unstamped batch.
177pub fn maybe_wrap_schema(schema: &SchemaRef) -> crate::error::Result<SchemaRef> {
178    // Fast path: only a struct column can be a native histogram; if there are
179    // none, skip the rebuild.
180    if !schema
181        .fields()
182        .iter()
183        .any(|f| matches!(f.data_type(), ArrowDataType::Struct(_)))
184    {
185        return Ok(schema.clone());
186    }
187    let new_fields: crate::error::Result<Vec<FieldRef>> = schema
188        .fields()
189        .iter()
190        .map(|f| {
191            let mut field = (**f).clone();
192            stamp_native_histogram_subfield_ids(&mut field)?;
193            Ok(Arc::new(field))
194        })
195        .collect();
196    Ok(Arc::new(Schema::new_with_metadata(
197        Fields::from(new_fields?),
198        schema.metadata().clone(),
199    )))
200}
201
202/// Parquet field ID base for internal columns (__primary_key, __sequence, __op_type).
203/// Uses bit 30 to distinguish from user column IDs and fit in positive i32 range.
204pub(crate) const INTERNAL_PARQUET_FIELD_ID_BASE: u32 = 1 << 30;
205
206/// Parquet field ID for the __primary_key column.
207pub(crate) const PRIMARY_KEY_PARQUET_FIELD_ID: u32 = INTERNAL_PARQUET_FIELD_ID_BASE;
208/// Parquet field ID for the __sequence column.
209pub(crate) const SEQUENCE_PARQUET_FIELD_ID: u32 = INTERNAL_PARQUET_FIELD_ID_BASE + 1;
210/// Parquet field ID for the __op_type column.
211pub(crate) const OP_TYPE_PARQUET_FIELD_ID: u32 = INTERNAL_PARQUET_FIELD_ID_BASE + 2;
212
213/// Gets the arrow schema to store in parquet.
214pub fn to_sst_arrow_schema(metadata: &RegionMetadata) -> SchemaRef {
215    let fields = Fields::from_iter(
216        metadata
217            .schema
218            .arrow_schema()
219            .fields()
220            .iter()
221            .zip(&metadata.column_metadatas)
222            .filter_map(|(field, column_meta)| {
223                if column_meta.semantic_type == SemanticType::Field {
224                    Some(Arc::new(with_field_id(
225                        (**field).clone(),
226                        column_meta.column_id,
227                    )))
228                } else {
229                    // We have fixed positions for tags (primary key) and time index.
230                    None
231                }
232            })
233            .chain([Arc::new(with_field_id(
234                (*metadata.time_index_field()).clone(),
235                metadata.time_index_column().column_id,
236            ))])
237            .chain(internal_fields()),
238    );
239
240    Arc::new(Schema::new(fields))
241}
242
243/// Options of flat schema.
244pub struct FlatSchemaOptions {
245    /// Whether to store primary key columns additionally instead of an encoded column.
246    pub raw_pk_columns: bool,
247    /// Whether to use dictionary encoding for string primary key columns
248    /// when storing primary key columns.
249    /// Only takes effect when `raw_pk_columns` is true.
250    pub string_pk_use_dict: bool,
251    /// The column's concretized JSON types, to be set into Arrow schema.
252    /// Otherwise it's empty struct in the Arrow schema.
253    pub concretized_json_types: HashMap<String, DataType>,
254}
255
256impl Default for FlatSchemaOptions {
257    fn default() -> Self {
258        Self {
259            raw_pk_columns: true,
260            string_pk_use_dict: true,
261            concretized_json_types: HashMap::new(),
262        }
263    }
264}
265
266impl FlatSchemaOptions {
267    /// Creates a options according to the primary key encoding.
268    pub fn from_encoding(encoding: PrimaryKeyEncoding) -> Self {
269        if encoding == PrimaryKeyEncoding::Dense {
270            Self::default()
271        } else {
272            Self {
273                raw_pk_columns: false,
274                string_pk_use_dict: false,
275                concretized_json_types: HashMap::new(),
276            }
277        }
278    }
279}
280
281/// Gets the arrow schema to store in parquet.
282///
283/// The schema is:
284/// ```text
285/// primary key columns, field columns, time index, __primary_key, __sequence, __op_type
286/// ```
287///
288/// # Panics
289/// Panics if the metadata is invalid.
290pub fn to_flat_sst_arrow_schema(
291    metadata: &RegionMetadata,
292    options: &FlatSchemaOptions,
293) -> SchemaRef {
294    let num_fields = flat_sst_arrow_schema_column_num(metadata, options);
295    let mut fields = Vec::with_capacity(num_fields);
296    let schema = metadata.schema.arrow_schema();
297    if options.raw_pk_columns {
298        for pk_id in &metadata.primary_key {
299            let pk_index = metadata.column_index_by_id(*pk_id).unwrap();
300            let column_id = metadata.column_metadatas[pk_index].column_id;
301            if options.string_pk_use_dict {
302                let old_field = &schema.fields[pk_index];
303                let new_field = tag_maybe_to_dictionary_field(
304                    &metadata.column_metadatas[pk_index].column_schema.data_type,
305                    old_field,
306                );
307                let new_field = concretize_json_type(new_field, options);
308                fields.push(Arc::new(with_field_id((*new_field).clone(), column_id)));
309            }
310        }
311    }
312    let remaining_fields = schema
313        .fields()
314        .iter()
315        .zip(&metadata.column_metadatas)
316        .filter_map(|(field, column_meta)| {
317            if column_meta.semantic_type == SemanticType::Field {
318                let field = concretize_json_type(field.clone(), options);
319                Some(Arc::new(with_field_id(
320                    Arc::unwrap_or_clone(field),
321                    column_meta.column_id,
322                )))
323            } else {
324                None
325            }
326        })
327        .chain([Arc::new(with_field_id(
328            (*metadata.time_index_field()).clone(),
329            metadata.time_index_column().column_id,
330        ))])
331        .chain(internal_fields());
332    for field in remaining_fields {
333        fields.push(field);
334    }
335
336    Arc::new(Schema::new(fields))
337}
338
339fn concretize_json_type(field: Arc<Field>, options: &FlatSchemaOptions) -> Arc<Field> {
340    if let Some(data_type) = options.concretized_json_types.get(field.name()) {
341        let mut field = Arc::unwrap_or_clone(field);
342        field.set_data_type(data_type.clone());
343        Arc::new(field)
344    } else {
345        field
346    }
347}
348
349/// Returns the number of columns in the flat format.
350pub fn flat_sst_arrow_schema_column_num(
351    metadata: &RegionMetadata,
352    options: &FlatSchemaOptions,
353) -> usize {
354    if options.raw_pk_columns {
355        metadata.column_metadatas.len() + 3
356    } else {
357        metadata.column_metadatas.len() + 3 - metadata.primary_key.len()
358    }
359}
360
361/// Helper function to create a dictionary field from a field.
362fn to_dictionary_field(field: &Field) -> Field {
363    let mut new_field = Field::new_dictionary(
364        field.name(),
365        datatypes::arrow::datatypes::DataType::UInt32,
366        field.data_type().clone(),
367        field.is_nullable(),
368    );
369
370    // retain field_id metadata
371    if let Some(field_id) = field.metadata().get(PARQUET_FIELD_ID_KEY) {
372        new_field
373            .metadata_mut()
374            .insert(PARQUET_FIELD_ID_KEY.to_string(), field_id.clone());
375    }
376
377    new_field
378}
379
380/// Helper function to create a dictionary field from a field if it is a string column.
381pub(crate) fn tag_maybe_to_dictionary_field(
382    data_type: &ConcreteDataType,
383    field: &Arc<Field>,
384) -> Arc<Field> {
385    if data_type.is_string() {
386        Arc::new(to_dictionary_field(field))
387    } else {
388        field.clone()
389    }
390}
391
392/// Fields for internal columns.
393pub(crate) fn internal_fields() -> [FieldRef; 3] {
394    // Internal columns are always not null.
395    [
396        Arc::new(with_field_id(
397            Field::new_dictionary(
398                PRIMARY_KEY_COLUMN_NAME,
399                ArrowDataType::UInt32,
400                ArrowDataType::Binary,
401                false,
402            ),
403            PRIMARY_KEY_PARQUET_FIELD_ID,
404        )),
405        Arc::new(with_field_id(
406            Field::new(SEQUENCE_COLUMN_NAME, ArrowDataType::UInt64, false),
407            SEQUENCE_PARQUET_FIELD_ID,
408        )),
409        Arc::new(with_field_id(
410            Field::new(OP_TYPE_COLUMN_NAME, ArrowDataType::UInt8, false),
411            OP_TYPE_PARQUET_FIELD_ID,
412        )),
413    ]
414}
415
416/// Returns a copy of `schema` with the `__primary_key` field replaced by a plain `Binary` field.
417pub(crate) fn override_pk_field_to_binary(schema: &SchemaRef) -> SchemaRef {
418    let new_fields = schema
419        .fields()
420        .iter()
421        .map(|field| {
422            if field.name() == PRIMARY_KEY_COLUMN_NAME {
423                let mut new_field = Field::new(
424                    PRIMARY_KEY_COLUMN_NAME,
425                    ArrowDataType::Binary,
426                    field.is_nullable(),
427                );
428                // Preserve the field_id metadata so parquet readers that require
429                // all columns to carry a field_id don't fail.
430                if let Some(field_id) = field.metadata().get(PARQUET_FIELD_ID_KEY) {
431                    new_field
432                        .metadata_mut()
433                        .insert(PARQUET_FIELD_ID_KEY.to_string(), field_id.clone());
434                }
435                Arc::new(new_field)
436            } else {
437                field.clone()
438            }
439        })
440        .collect::<Vec<_>>();
441    Arc::new(Schema::new(new_fields))
442}
443
444/// Gets the estimated number of series from record batches.
445///
446/// This struct tracks the last timestamp value to detect series boundaries
447/// by observing when timestamps decrease (indicating a new series).
448#[derive(Default)]
449pub(crate) struct SeriesEstimator {
450    /// The last timestamp value seen
451    last_timestamp: Option<i64>,
452    /// The estimated number of series
453    series_count: u64,
454}
455
456impl SeriesEstimator {
457    /// Updates the estimator with a new record batch in flat format.
458    ///
459    /// This method examines the time index column to detect series boundaries.
460    pub(crate) fn update_flat(&mut self, record_batch: &RecordBatch) {
461        let batch_rows = record_batch.num_rows();
462        if batch_rows == 0 {
463            return;
464        }
465
466        let time_index_pos = time_index_column_index(record_batch.num_columns());
467        let timestamps = record_batch.column(time_index_pos);
468        let Some((ts_values, _unit)) = timestamp_array_to_primitive(timestamps) else {
469            return;
470        };
471        let values = ts_values.values();
472
473        // Checks if there's a boundary between the last batch and this batch
474        if let Some(last_ts) = self.last_timestamp {
475            if values[0] <= last_ts {
476                self.series_count += 1;
477            }
478        } else {
479            // First batch, counts as first series
480            self.series_count = 1;
481        }
482
483        // Counts series boundaries within this batch.
484        for i in 0..batch_rows - 1 {
485            // We assumes the same timestamp as a new series, which is different from
486            // how we split batches.
487            if values[i] >= values[i + 1] {
488                self.series_count += 1;
489            }
490        }
491
492        // Updates the last timestamp
493        self.last_timestamp = Some(values[batch_rows - 1]);
494    }
495
496    /// Returns the estimated number of series.
497    pub(crate) fn finish(&mut self) -> u64 {
498        self.last_timestamp = None;
499        let count = self.series_count;
500        self.series_count = 0;
501
502        count
503    }
504}
505
506#[cfg(test)]
507mod tests {
508    use std::sync::Arc;
509
510    use ::parquet::arrow::AsyncArrowWriter;
511    use ::parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
512    use ::parquet::basic::LogicalType;
513    use ::parquet::variant::{VariantArray, VariantType, json_to_variant};
514    use common_query::prelude::greptime_native_histogram;
515    use datatypes::arrow::array::{
516        ArrayRef, BinaryArray, DictionaryArray, Int64Array, StringArray, StructArray,
517        TimestampMillisecondArray, UInt8Array, UInt32Array, UInt64Array,
518    };
519    use datatypes::arrow::datatypes::{DataType as ArrowDataType, Field, Schema, TimeUnit};
520    use datatypes::arrow::record_batch::RecordBatch;
521    use datatypes::extension::json::{Json2ExtensionType, Json2PhysicalLayout};
522    use datatypes::vectors::json::array::JsonArray;
523    use serde_json::json;
524
525    use super::*;
526
527    fn new_flat_record_batch(timestamps: &[i64]) -> RecordBatch {
528        // Flat format has: [fields..., time_index, __primary_key, __sequence, __op_type]
529        let num_cols = 4; // time_index + 3 internal columns
530        let time_index_pos = time_index_column_index(num_cols);
531        assert_eq!(time_index_pos, 0); // For 4 columns, time index should be at position 0
532
533        let time_array = Arc::new(TimestampMillisecondArray::from(timestamps.to_vec()));
534        let pk_array = Arc::new(DictionaryArray::new(
535            UInt32Array::from(vec![0; timestamps.len()]),
536            Arc::new(BinaryArray::from(vec![b"test".as_slice()])),
537        ));
538        let seq_array = Arc::new(UInt64Array::from(vec![1; timestamps.len()]));
539        let op_array = Arc::new(UInt8Array::from(vec![1; timestamps.len()]));
540
541        let schema = Arc::new(Schema::new(vec![
542            Field::new(
543                "time",
544                ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
545                false,
546            ),
547            Field::new_dictionary(
548                "__primary_key",
549                ArrowDataType::UInt32,
550                ArrowDataType::Binary,
551                false,
552            ),
553            Field::new("__sequence", ArrowDataType::UInt64, false),
554            Field::new("__op_type", ArrowDataType::UInt8, false),
555        ]));
556
557        RecordBatch::try_new(schema, vec![time_array, pk_array, seq_array, op_array]).unwrap()
558    }
559
560    #[test]
561    fn test_series_estimator_flat_empty_batch() {
562        let mut estimator = SeriesEstimator::default();
563        let record_batch = new_flat_record_batch(&[]);
564        estimator.update_flat(&record_batch);
565        assert_eq!(0, estimator.finish());
566    }
567
568    #[test]
569    fn test_series_estimator_flat_single_batch() {
570        let mut estimator = SeriesEstimator::default();
571        let record_batch = new_flat_record_batch(&[1, 2, 3]);
572        estimator.update_flat(&record_batch);
573        assert_eq!(1, estimator.finish());
574    }
575
576    #[test]
577    fn test_series_estimator_flat_series_boundary_within_batch() {
578        let mut estimator = SeriesEstimator::default();
579        // Timestamps decrease from 3 to 2, indicating a series boundary
580        let record_batch = new_flat_record_batch(&[1, 2, 3, 2, 4, 5]);
581        estimator.update_flat(&record_batch);
582        // Should detect boundary at position 3 (3 >= 2)
583        assert_eq!(2, estimator.finish());
584    }
585
586    #[test]
587    fn test_series_estimator_flat_multiple_boundaries_within_batch() {
588        let mut estimator = SeriesEstimator::default();
589        // Multiple series boundaries: 5>=4, 6>=3
590        let record_batch = new_flat_record_batch(&[1, 2, 5, 4, 6, 3, 7]);
591        estimator.update_flat(&record_batch);
592        assert_eq!(3, estimator.finish());
593    }
594
595    #[test]
596    fn test_series_estimator_flat_equal_timestamps() {
597        let mut estimator = SeriesEstimator::default();
598        // Equal timestamps are considered as new series
599        let record_batch = new_flat_record_batch(&[1, 2, 2, 3, 3, 3, 4]);
600        estimator.update_flat(&record_batch);
601        // Boundaries at: 2>=2, 3>=3, 3>=3
602        assert_eq!(4, estimator.finish());
603    }
604
605    #[test]
606    fn test_series_estimator_flat_multiple_batches_continuation() {
607        let mut estimator = SeriesEstimator::default();
608
609        // First batch: timestamps 1, 2, 3
610        let batch1 = new_flat_record_batch(&[1, 2, 3]);
611        estimator.update_flat(&batch1);
612
613        // Second batch: timestamps 4, 5, 6 (continuation)
614        let batch2 = new_flat_record_batch(&[4, 5, 6]);
615        estimator.update_flat(&batch2);
616
617        assert_eq!(1, estimator.finish());
618    }
619
620    #[test]
621    fn test_series_estimator_flat_multiple_batches_new_series() {
622        let mut estimator = SeriesEstimator::default();
623
624        // First batch: timestamps 1, 2, 3
625        let batch1 = new_flat_record_batch(&[1, 2, 3]);
626        estimator.update_flat(&batch1);
627
628        // Second batch: timestamps 2, 3, 4 (goes back to 2, new series)
629        let batch2 = new_flat_record_batch(&[2, 3, 4]);
630        estimator.update_flat(&batch2);
631
632        assert_eq!(2, estimator.finish());
633    }
634
635    #[test]
636    fn test_series_estimator_flat_boundary_at_batch_edge_equal() {
637        let mut estimator = SeriesEstimator::default();
638
639        // First batch ending at 5
640        let batch1 = new_flat_record_batch(&[1, 2, 5]);
641        estimator.update_flat(&batch1);
642
643        // Second batch starting at 5 (equal timestamp, new series)
644        let batch2 = new_flat_record_batch(&[5, 6, 7]);
645        estimator.update_flat(&batch2);
646
647        assert_eq!(2, estimator.finish());
648    }
649
650    #[test]
651    fn test_series_estimator_flat_mixed_batches() {
652        let mut estimator = SeriesEstimator::default();
653
654        // Batch 1: single series [10, 20, 30]
655        let batch1 = new_flat_record_batch(&[10, 20, 30]);
656        estimator.update_flat(&batch1);
657
658        // Batch 2: starts new series [5, 15], boundary within batch [15, 10, 25]
659        let batch2 = new_flat_record_batch(&[5, 15, 10, 25]);
660        estimator.update_flat(&batch2);
661
662        // Batch 3: continues from 25 to [30, 35]
663        let batch3 = new_flat_record_batch(&[30, 35]);
664        estimator.update_flat(&batch3);
665
666        // Expected: 1 (batch1) + 1 (batch2 start) + 1 (within batch2) = 3
667        assert_eq!(3, estimator.finish());
668    }
669
670    #[test]
671    fn test_series_estimator_flat_descending_timestamps() {
672        let mut estimator = SeriesEstimator::default();
673        // Strictly descending timestamps - each pair creates a boundary
674        let record_batch = new_flat_record_batch(&[10, 9, 8, 7, 6]);
675        estimator.update_flat(&record_batch);
676        // Boundaries: 10>=9, 9>=8, 8>=7, 7>=6 = 4 boundaries + 1 initial = 5 series
677        assert_eq!(5, estimator.finish());
678    }
679
680    #[test]
681    fn test_series_estimator_flat_finish_resets_state() {
682        let mut estimator = SeriesEstimator::default();
683
684        let batch1 = new_flat_record_batch(&[1, 2, 3]);
685        estimator.update_flat(&batch1);
686
687        assert_eq!(1, estimator.finish());
688
689        // After finish, state should be reset
690        let batch2 = new_flat_record_batch(&[4, 5, 6]);
691        estimator.update_flat(&batch2);
692
693        assert_eq!(1, estimator.finish());
694    }
695
696    /// Build a native-histogram struct field whose top-level `PARQUET:field_id`
697    /// is `column_id` (as `with_field_id` does on the real write path).
698    fn histogram_field(name: &str, column_id: u32) -> Field {
699        use common_query::native_histogram::native_histogram_value_type;
700        use datatypes::data_type::DataType;
701        with_field_id(
702            Field::new(name, native_histogram_value_type().as_arrow_type(), true),
703            column_id,
704        )
705    }
706
707    /// Asserts `field` is a stamped native-histogram struct: it carries the
708    /// `greptime.histogram` extension and every sub-field (and list element)
709    /// carries its reserved `PARQUET:field_id` namespaced by `column_id`.
710    fn assert_histogram_stamped(field: &Field, column_id: i32) {
711        use arrow_schema::extension::ExtensionType;
712        use common_query::native_histogram::{
713            native_histogram_list_element_id, native_histogram_subfield_id,
714        };
715        use datatypes::extension::histogram::HistogramExtensionType;
716
717        assert_eq!(
718            field
719                .metadata()
720                .get(arrow_schema::extension::EXTENSION_TYPE_NAME_KEY)
721                .map(|s| s.as_str()),
722            Some(HistogramExtensionType::NAME),
723            "histogram field must carry the greptime.histogram extension"
724        );
725        let ArrowDataType::Struct(children) = field.data_type() else {
726            panic!("expected a struct, got {:?}", field.data_type());
727        };
728        for child in children {
729            let expected = native_histogram_subfield_id(column_id, child.name())
730                .unwrap_or_else(|| panic!("no id for sub-field {}", child.name()));
731            let got: i32 = child
732                .metadata()
733                .get(PARQUET_FIELD_ID_KEY)
734                .unwrap_or_else(|| panic!("sub-field {} missing field id", child.name()))
735                .parse()
736                .unwrap();
737            assert_eq!(got, expected, "sub-field {} id", child.name());
738            if let ArrowDataType::List(elem) = child.data_type() {
739                let elem_expected =
740                    native_histogram_list_element_id(column_id, child.name()).unwrap();
741                let elem_got: i32 = elem
742                    .metadata()
743                    .get(PARQUET_FIELD_ID_KEY)
744                    .unwrap_or_else(|| panic!("list element of {} missing id", child.name()))
745                    .parse()
746                    .unwrap();
747                assert_eq!(
748                    elem_got,
749                    elem_expected,
750                    "list element id of {}",
751                    child.name()
752                );
753            }
754        }
755    }
756
757    #[test]
758    fn test_maybe_wrap_schema_native_histogram() {
759        let schema = Arc::new(Schema::new(vec![
760            Field::new(
761                "greptime_timestamp",
762                ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
763                false,
764            ),
765            histogram_field(greptime_native_histogram(), 1),
766        ]));
767
768        let wrapped = maybe_wrap_schema(&schema).unwrap();
769        let hist = wrapped
770            .field_with_name(greptime_native_histogram())
771            .expect("histogram field present");
772        // The struct has 18 sub-fields.
773        let ArrowDataType::Struct(children) = hist.data_type() else {
774            unreachable!()
775        };
776        assert_eq!(children.len(), 18);
777        assert_histogram_stamped(hist, 1);
778    }
779
780    #[test]
781    fn test_maybe_wrap_schema_multiple_histograms_disjoint_ids() {
782        // Two histogram columns with distinct parent column ids get disjoint
783        // sub-field ids (defensive: the metric engine yields at most one
784        // histogram column, but the scheme must stay correct if more appear).
785        use common_query::native_histogram::native_histogram_subfield_id;
786
787        let schema = Arc::new(Schema::new(vec![
788            histogram_field(greptime_native_histogram(), 1),
789            histogram_field(greptime_native_histogram(), 7),
790        ]));
791        let wrapped = maybe_wrap_schema(&schema).unwrap();
792        let h1 = &wrapped.fields()[0];
793        let h2 = &wrapped.fields()[1];
794        assert_histogram_stamped(h1, 1);
795        assert_histogram_stamped(h2, 7);
796        // The same sub-field name resolves to different ids across columns.
797        assert_ne!(
798            native_histogram_subfield_id(1, "sum"),
799            native_histogram_subfield_id(7, "sum")
800        );
801    }
802
803    #[test]
804    fn test_maybe_wrap_schema_recognizes_histogram_by_type() {
805        let schema = Arc::new(Schema::new(vec![histogram_field("custom_histogram", 5)]));
806
807        let wrapped = maybe_wrap_schema(&schema).unwrap();
808        let hist = wrapped.field_with_name("custom_histogram").unwrap();
809        assert_histogram_stamped(hist, 5);
810    }
811
812    #[test]
813    fn test_maybe_wrap_schema_plain_struct_not_stamped() {
814        use arrow_schema::extension::EXTENSION_TYPE_NAME_KEY;
815
816        let plain = ArrowDataType::Struct(
817            vec![
818                Arc::new(Field::new("a", ArrowDataType::Int32, true)),
819                Arc::new(Field::new("b", ArrowDataType::Utf8, true)),
820            ]
821            .into(),
822        );
823        let schema = Arc::new(Schema::new(vec![
824            Field::new(
825                "ts",
826                ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
827                false,
828            ),
829            Field::new("data", plain, true),
830        ]));
831
832        let wrapped = maybe_wrap_schema(&schema).unwrap();
833        let data = wrapped.field_with_name("data").unwrap();
834        assert!(
835            data.metadata().get(EXTENSION_TYPE_NAME_KEY).is_none(),
836            "non-histogram struct must not get the extension"
837        );
838        if let ArrowDataType::Struct(children) = data.data_type() {
839            for child in children {
840                assert!(
841                    child.metadata().get(PARQUET_FIELD_ID_KEY).is_none(),
842                    "non-histogram sub-field {} must not get a field id",
843                    child.name()
844                );
845            }
846        }
847    }
848
849    #[test]
850    fn test_maybe_wrap_schema_no_struct_unchanged() {
851        let schema: Arc<Schema> = Arc::new(Schema::new(vec![
852            Field::new(
853                "ts",
854                ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
855                false,
856            ),
857            Field::new("v", ArrowDataType::Float64, true),
858        ]));
859        let wrapped = maybe_wrap_schema(&schema).unwrap();
860        assert!(
861            Arc::ptr_eq(&wrapped, &schema),
862            "a schema without any struct column must be returned unchanged"
863        );
864    }
865
866    /// Writes `schema` through `maybe_wrap_schema` and a real parquet
867    /// [`ArrowWriter`], then returns the arrow schema read back from the file
868    /// footer. This proves the `greptime.histogram` extension and the nested
869    /// `PARQUET:field_id`s actually land on disk, not just in memory.
870    ///
871    /// `maybe_wrap_schema` is exactly what the SST parquet writer hands to
872    /// `AsyncArrowWriter` (see `writer.rs`); the sync [`ArrowWriter`] shares
873    /// the same arrow-to-parquet schema conversion, so the footer it emits is
874    /// the on-disk contract this change introduces. An empty batch suffices
875    /// because the parquet footer always carries the schema.
876    fn parquet_footer_arrow_schema(schema: &SchemaRef) -> SchemaRef {
877        use ::parquet::arrow::ArrowWriter;
878        use ::parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
879        use ::parquet::file::properties::WriterProperties;
880        use bytes::Bytes;
881
882        let wrapped = maybe_wrap_schema(schema).unwrap();
883        let mut bytes = Vec::new();
884        let props = WriterProperties::builder().build();
885        let mut writer = ArrowWriter::try_new(&mut bytes, wrapped.clone(), Some(props)).unwrap();
886        writer
887            .write(&RecordBatch::new_empty(wrapped.clone()))
888            .unwrap();
889        writer.close().unwrap();
890
891        ParquetRecordBatchReaderBuilder::try_new(Bytes::from(bytes))
892            .unwrap()
893            .schema()
894            .clone()
895    }
896
897    #[test]
898    fn test_maybe_wrap_schema_survives_parquet_roundtrip() {
899        // On-disk contract: after writing through the parquet writer path, the
900        // footer still carries the greptime.histogram extension and every
901        // nested (sub-field + list-element) PARQUET:field_id.
902        let schema = Arc::new(Schema::new(vec![
903            Field::new(
904                "greptime_timestamp",
905                ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
906                false,
907            ),
908            histogram_field(greptime_native_histogram(), 3),
909        ]));
910
911        let on_disk = parquet_footer_arrow_schema(&schema);
912        let hist = on_disk
913            .field_with_name(greptime_native_histogram())
914            .expect("histogram field present");
915        assert_histogram_stamped(hist, 3);
916    }
917
918    #[test]
919    fn test_parquet_roundtrip_recognizes_histogram_by_type() {
920        // The persisted type, rather than a process-local configured name,
921        // identifies native histograms across upgrades and prefix changes.
922        let schema = Arc::new(Schema::new(vec![
923            Field::new(
924                "ts",
925                ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
926                false,
927            ),
928            histogram_field("custom_histogram", 9),
929        ]));
930
931        let on_disk = parquet_footer_arrow_schema(&schema);
932        let hist = on_disk.field_with_name("custom_histogram").unwrap();
933        assert_histogram_stamped(hist, 9);
934    }
935
936    #[test]
937    fn test_maybe_wrap_schema_overflows_return_error() {
938        // A column id of 12_582_912 makes the derived sub-field id overflow
939        // i32 (BASE + column_id*64 == i32::MAX + 1). The write path must
940        // surface this as an error rather than silently dropping the field
941        // id, wrapping, or panicking.
942        let schema = Arc::new(Schema::new(vec![
943            Field::new(
944                "greptime_timestamp",
945                ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
946                false,
947            ),
948            histogram_field(greptime_native_histogram(), 12_582_912),
949        ]));
950        let err = maybe_wrap_schema(&schema).unwrap_err();
951        assert!(
952            matches!(
953                err,
954                crate::error::Error::InvalidNativeHistogramSubfield { .. }
955            ),
956            "expected InvalidNativeHistogramSubfield, got {:?}",
957            err
958        );
959    }
960
961    #[test]
962    fn test_maybe_wrap_schema_missing_field_id_returns_error() {
963        // The parent column's PARQUET:field_id namespaces every sub-field id.
964        // If it is absent (e.g. a histogram struct handed to the writer
965        // without the write path's stamping), the writer must fail loudly
966        // rather than silently namespace under column 0, which would collide
967        // with that column's nested ids.
968        use common_query::native_histogram::native_histogram_value_type;
969        use datatypes::data_type::DataType;
970
971        let field = Field::new(
972            greptime_native_histogram(),
973            native_histogram_value_type().as_arrow_type(),
974            true,
975        );
976        assert!(
977            field.metadata().get(PARQUET_FIELD_ID_KEY).is_none(),
978            "fixture must not carry a field id"
979        );
980        let schema = Arc::new(Schema::new(vec![field]));
981        let err = maybe_wrap_schema(&schema).unwrap_err();
982        assert!(
983            matches!(
984                err,
985                crate::error::Error::InvalidNativeHistogramFieldId { .. }
986            ),
987            "expected InvalidNativeHistogramFieldId, got {:?}",
988            err
989        );
990    }
991
992    #[test]
993    fn test_maybe_wrap_schema_field_id_above_i32_max_returns_error() {
994        // `with_field_id` serializes the column id from a u32, so a valid id
995        // above i32::MAX (e.g. u32::MAX) must not be silently parsed as a
996        // failed i32 and collapsed onto column 0's nested ids. It must
997        // surface a checked-conversion error instead.
998        let schema = Arc::new(Schema::new(vec![
999            Field::new(
1000                "greptime_timestamp",
1001                ArrowDataType::Timestamp(TimeUnit::Millisecond, None),
1002                false,
1003            ),
1004            histogram_field(greptime_native_histogram(), u32::MAX),
1005        ]));
1006        let err = maybe_wrap_schema(&schema).unwrap_err();
1007        assert!(
1008            matches!(
1009                err,
1010                crate::error::Error::InvalidNativeHistogramFieldId { .. }
1011            ),
1012            "expected InvalidNativeHistogramFieldId, got {:?}",
1013            err
1014        );
1015    }
1016
1017    fn json2_v2_test_type() -> ArrowDataType {
1018        ArrowDataType::Struct(
1019            vec![
1020                Arc::new(Field::new("active", ArrowDataType::Boolean, true)),
1021                Arc::new(Field::new("hot", ArrowDataType::Int64, true)),
1022                Arc::new(Field::new("name", ArrowDataType::Utf8, true)),
1023            ]
1024            .into(),
1025        )
1026    }
1027
1028    /// Validates the persisted-format foundation for the JSON2 v2 remainder.
1029    ///
1030    /// JSON2 will store `!__remainder__!` as a nested Variant child of its root
1031    /// Struct. Before enabling that layout in production, this test ensures the
1032    /// SST schema wrapper and Arrow writer preserve the Variant extension,
1033    /// encode the Parquet Variant logical type, and round-trip the values
1034    /// without changing the surrounding Struct.
1035    #[tokio::test]
1036    async fn test_nested_variant_survives_sst_writer_schema_roundtrip()
1037    -> Result<(), Box<dyn std::error::Error>> {
1038        let json: ArrayRef = Arc::new(StringArray::from(vec![
1039            Some(r#"{}"#),
1040            Some(r#"{"name":"Alice","active":true}"#),
1041            Some(r#"{"nested":{"count":42},"items":[1,"two",null]}"#),
1042            Some(r#"{"\u5b57\u6bb5":"\u503c"}"#),
1043            None,
1044        ]));
1045        let remainder = json_to_variant(&json)?;
1046        let remainder_field = remainder.field("!__remainder__!");
1047        let remainder_array = ArrayRef::from(remainder);
1048        let hot_field = Field::new("hot", ArrowDataType::Int64, true);
1049        let data_array = Arc::new(StructArray::new(
1050            vec![remainder_field.clone(), hot_field.clone()].into(),
1051            vec![
1052                remainder_array,
1053                Arc::new(Int64Array::from(vec![
1054                    Some(1),
1055                    Some(2),
1056                    Some(3),
1057                    Some(4),
1058                    None,
1059                ])),
1060            ],
1061            None,
1062        ));
1063        let data_field = Field::new(
1064            "data",
1065            ArrowDataType::Struct(vec![remainder_field, hot_field].into()),
1066            true,
1067        )
1068        .with_extension_type(Json2ExtensionType::default());
1069        let schema = Arc::new(Schema::new(vec![data_field]));
1070        let source = RecordBatch::try_new(schema.clone(), vec![data_array])?;
1071
1072        let wrapped = maybe_wrap_schema(&schema)?;
1073        let mut buffer = Vec::new();
1074        let mut writer = AsyncArrowWriter::try_new(&mut buffer, wrapped, None)?;
1075        writer.write(&source).await?;
1076        writer.close().await?;
1077
1078        let builder = ParquetRecordBatchReaderBuilder::try_new(bytes::Bytes::from(buffer))?;
1079        let parquet_remainder =
1080            &builder.parquet_schema().root_schema().get_fields()[0].get_fields()[0];
1081        assert_eq!(
1082            parquet_remainder.get_basic_info().logical_type_ref(),
1083            Some(&LogicalType::variant(None))
1084        );
1085
1086        let ArrowDataType::Struct(children) = builder.schema().field_with_name("data")?.data_type()
1087        else {
1088            unreachable!();
1089        };
1090        assert!(children[0].has_valid_extension_type::<VariantType>());
1091
1092        let mut reader = builder.build()?;
1093        let result = reader.next().unwrap()?;
1094        assert_eq!(source, result);
1095        let result_field = result.schema().field(0).clone();
1096        let result = result
1097            .column(0)
1098            .as_any()
1099            .downcast_ref::<StructArray>()
1100            .unwrap();
1101        VariantArray::try_new(result.column(0))?;
1102        let result: ArrayRef = Arc::new(result.clone());
1103        let result =
1104            JsonArray::from(&result).project_to_v2(&result_field, &json2_v2_test_type())?;
1105        assert_eq!(
1106            json!({"active": true, "hot": 2, "name": "Alice"}),
1107            JsonArray::from(&result).try_get_value(1)?
1108        );
1109        Ok(())
1110    }
1111
1112    /// Ensures future readers retain compatibility with the first JSON2 v2 layout.
1113    #[test]
1114    fn test_read_json2_v2_fixture() -> Result<(), Box<dyn std::error::Error>> {
1115        let bytes = bytes::Bytes::from_static(include_bytes!("../test-data/json2-v2.parquet"));
1116        let builder = ParquetRecordBatchReaderBuilder::try_new(bytes)?;
1117        let field = builder.schema().field(0).clone();
1118        assert!(Json2PhysicalLayout::try_from_root(&field)?.is_version_2());
1119
1120        let batch = builder.build()?.next().unwrap()?;
1121        let data = batch
1122            .column(0)
1123            .as_any()
1124            .downcast_ref::<StructArray>()
1125            .unwrap();
1126        VariantArray::try_new(data.column(0))?;
1127        let data: ArrayRef = Arc::new(data.clone());
1128        let data = JsonArray::from(&data).project_to_v2(&field, &json2_v2_test_type())?;
1129        assert_eq!(
1130            json!({"active": true, "hot": 2, "name": "Alice"}),
1131            JsonArray::from(&data).try_get_value(1)?
1132        );
1133        Ok(())
1134    }
1135}